Ollama und Persistenz-Patch

This commit is contained in:
2026-08-04 05:27:55 +02:00
parent 4355f6bfd6
commit 23339a8be3
31 changed files with 2510 additions and 263 deletions
+34 -1
View File
@@ -12,11 +12,44 @@ 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
# Local models
# Batched disk persistence. AI-THINK drafts, the GLPI-KB cache and graph-state
# are written sequentially at this interval and once more during shutdown.
BRAIN_PERSIST_INTERVAL=5m
# Ollama pool. OLLAMA_URLS takes precedence over legacy OLLAMA_URL.
OLLAMA_URL=http://localhost:11434
OLLAMA_URLS=http://localhost:11434
OLLAMA_NODE_NAMES=local-gpu
OLLAMA_NODE_WEIGHTS=1
OLLAMA_ROUTING_MODE=least_inflight
OLLAMA_NODE_MAX_INFLIGHT=1
OLLAMA_NODE_HEALTH_INTERVAL=15s
OLLAMA_NODE_FAILURE_COOLDOWN=30s
OLLAMA_NODE_REQUEST_TIMEOUT=8m
OLLAMA_FAILOVER_ENABLED=true
# 0 = try all configured nodes.
OLLAMA_FAILOVER_ATTEMPTS=0
OLLAMA_REQUIRE_SAME_MODEL_DIGEST=true
OLLAMA_REQUIRE_EMBEDDING_MODEL=true
OLLAMA_CHAT_MODEL=qwen3:8b
OLLAMA_EMBEDDING_MODEL=embeddinggemma
# Optional read-only GLPI Knowledge Base ingest.
GLPI_KB_ENABLED=false
GLPI_URL=https://glpi.example.invalid
GLPI_API_VERSION=v2.3
GLPI_CLIENT_ID=
GLPI_CLIENT_SECRET=
GLPI_USERNAME=
GLPI_PASSWORD=
GLPI_ALLOW_INSECURE_HTTP=false
GLPI_TIMEOUT=20s
GLPI_KB_PATH=auto
GLPI_KB_FILTER=
GLPI_KB_LIMIT=500
GLPI_KB_SYNC_INTERVAL=10m
GLPI_KB_SOURCE=GLPI Knowledge Base
# Sequential enrichment
BRAIN_AUTO_ENRICH=true
BRAIN_SCAN_INTERVAL=20s
+22 -18
View File
@@ -1,39 +1,43 @@
# Architektur und Vertrauensgrenzen
```text
Agent knowledge/ ──ro──┐
Agent runs.jsonl ──ro──┼──► Ingest ─► Graph + Vectors ─► Retrieval/Inference ─► SSE/Web UI
KB staging/ ──────rw───┘ │
├──► AI Edge (staging)
├──► optional SearXNG evidence
└──► AI-THINK JSON in staging
Lokale Knowledge-JSONs ──ro──┐
GLPI Knowledge Base ─────ro──┼──► Ingest ─► In-Memory Graph + Vectors
Agent runs.jsonl ─────────ro──┘ │
├──► Ollama Pool
├──► AI-THINK / Research
├──► SSE / Hierarchical LOD UI
└──► Batched Persistence
├── staging drafts
├── GLPI cache
└── graph-state.json
```
## Edge-Klassen
- `categorized_as`, `mentions`, `derived_from`: deterministisch aus JSON.
- `research_evidence`: aus einer explizit gestarteten Recherche, weiterhin `staging`.
- `categorized_as`, `mentions`, `derived_from`: deterministisch aus lokalen oder GLPI-Daten.
- `research_evidence`: aus explizit kontrollierter Recherche, weiterhin `staging`.
- `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; verhindert Endlosschleifen.
- `rejected`: intern gespeicherte Prüfung ohne sichtbare Beziehung.
## Ereignismodell
## Schreibmodell
`Activity` ist von der Graph-Persistenz getrennt. Eine Anfrage verändert deshalb nicht automatisch die Wissensbasis. Erst der sequenzielle Enrichment-Schritt darf eine neue Edge oder einen AI-THINK-Entwurf anlegen.
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.
## Datenhoheit
## Ollama-Pool
Produktive JSON-Dateien werden niemals verändert. Der einzige schreibende Fremdpfad ist das explizit konfigurierte Staging-Verzeichnis. Alle Dateien werden atomar über Temporärdatei und Rename geschrieben.
Chat- und Embeddingaufrufe werden als vollständige Requests geroutet. Healthchecks prüfen Erreichbarkeit und Modelle. Ein Per-Node-Inflight-Limit verhindert lokale Überlastung. Failover verändert keine produktiven Daten, da Schreiboperationen erst nach erfolgreicher Inferenz deterministisch ausgeführt werden.
## Hierarchische Render-Schicht
Der persistierte Wissensgraph bleibt unverändert vollständig. Das Frontend leitet daraus eine temporäre, hierarchische Render-Schicht ab:
Der persistierte Wissensgraph bleibt vollständig. Das Frontend leitet eine temporäre LOD-Struktur ab:
```text
Full Graph
└── Cortex-Region
└── Themenwolke (LOD 2)
└── lokale Gruppe (LOD 1)
└── einzelner Wissensknoten (LOD 0)
└── Themenwolke
└── lokale Gruppe
└── einzelner Wissensknoten
```
Die sichtbare Ebene wird durch Zoom, Fokus und Aktivitätszeiten bestimmt. Außenkanten werden nach sichtbaren Endpunkten, Relationstyp, Richtung, Herkunft und Status aggregiert. Eine Aktivierung öffnet ausschließlich den betroffenen Pfad der Hierarchie; nach Ablauf der Offenhaltezeit wird dieser wieder verdichtet. Die KI-, Retrieval- und Persistenzlogik arbeitet weiterhin ausschließlich mit den echten Node- und Edge-IDs.
Außenkanten werden nach sichtbaren Endpunkten, Relationstyp, Richtung, Herkunft und Status aggregiert. Aktivität öffnet nur den betroffenen Hierarchiepfad.
+30
View File
@@ -0,0 +1,30 @@
# Erweiterung: GLPI-KB, Ollama-Pool und gebündelte Persistenz
## GLPI-Knowledgebase
- read-only OAuth-Client für GLPI;
- automatische KnowbaseItem-Pfaderkennung über OpenAPI;
- periodischer und manueller Sync;
- Kategorien-, Quellen- und Wissens-Nodes mit verifizierten Edges;
- In-Memory-Nutzung sofort, Cache-Persistenz gebündelt;
- unveränderte Sync-Ergebnisse erzeugen keine neue Graphversion.
## Ollama-Pool
- bis zu 64 Nodes;
- `least_inflight`, `round_robin`, `weighted`, `fastest_recent`;
- Per-Node-Inflight-Limit;
- Healthchecks über `/api/tags`;
- optionale strikte Digest-Konsistenz;
- Cooldown und Failover bei retryfähigen Fehlern;
- Poolstatus über `/api/status`.
## Persistenz
- zentrale Queue für AI-THINK- und Cache-Dateien;
- Deduplizierung nach Zielpfad;
- serielles, atomares Schreiben;
- Graph-Snapshot erst nach den Wissensdateien;
- Standardintervall fünf Minuten;
- finaler Flush beim geregelten Shutdown;
- manueller Flush über `POST /api/flush`.
+23
View File
@@ -0,0 +1,23 @@
# GLPI-Knowledgebase-Ingest
Die GLPI-Integration ist read-only. Sie verwendet OAuth-Passwortfluss und die konfigurierte GLPI-API-Version. Sichtbarkeit und Entitätsrechte werden vom verwendeten GLPI-Service-Account durchgesetzt.
## Datenmodell
Jeder GLPI-Beitrag wird zu einem `knowledge`-Node:
```text
GLPI-KB-Beitrag
├── categorized_as ─► GLPI/ITIL-Themenbereich
└── derived_from ───► GLPI Knowledge Base
```
Metadaten enthalten GLPI-ID, KB-Kategorie-IDs, Sprache und Änderungszeit. HTML-Inhalt wird für Embedding und Analyse in Text überführt; der produktive GLPI-Beitrag wird nicht verändert.
## OpenAPI-Pfaderkennung
Mit `GLPI_KB_PATH=auto` wird `/api.php/doc.json` gelesen und die am besten passende lesbare `KnowbaseItem`-Collection gewählt. Für Installationen mit abweichender Route kann ein absoluter API-Pfad gesetzt werden.
## Cache
Der normalisierte GLPI-Teilgraph wird unter `BRAIN_DATA_DIR/glpi-kb-cache.json` gepuffert. Beim Start kann er sofort geladen werden. Aktualisierungen werden über denselben gebündelten Persistenzmechanismus wie der Graph geschrieben.
+33
View File
@@ -0,0 +1,33 @@
# Ollama-Pool
Der Pool verteilt vollständige Inferenzrequests auf bis zu 64 unabhängige Ollama-Instanzen. Er teilt kein einzelnes Modell über mehrere Hosts auf. Jeder verwendete Node lädt das jeweilige Modell lokal.
## Ablauf
```text
Embedding / AI-THINK Request
│
▼
Pool-Router
├─ Health/Modelprüfung
├─ Routing
├─ Per-Node Inflight-Limit
├─ Cooldown bei retryfähigem Fehler
└─ optionaler Failover
│
├── Ollama A
├── Ollama B
└── Ollama C
```
Retryfähig sind Netzwerkfehler, Timeouts, HTTP 408, HTTP 429, HTTP 5xx und ungültiges äußeres Ollama-JSON. Andere 4xx-Fehler werden nicht auf weitere Nodes gespiegelt.
`OLLAMA_FAILOVER_ATTEMPTS=0` bedeutet maximal alle konfigurierten Nodes versuchen.
## Modellkonsistenz
Mit `OLLAMA_REQUIRE_SAME_MODEL_DIGEST=true` werden die Digests aus `/api/tags` verglichen. Unterschiedliche Chat- oder Embedding-Digests machen den Pool inkompatibel, statt still unterschiedliche Ergebnisse zu erzeugen.
## Sicherheit
Ollama sollte nicht öffentlich erreichbar sein. Nutze ein separates KI-Netz, Firewall-Regeln, VPN oder einen abgesicherten Reverse Proxy. Das Brain implementiert keine zusätzliche Authentifizierung gegenüber Ollama.
+15
View File
@@ -0,0 +1,15 @@
# Gebündelte 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.
## Reihenfolge
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.
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.
## Trade-off
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.
+129 -112
View File
@@ -1,165 +1,182 @@
# Neural Knowledge Brain
Eigenständiger Go-Dienst für deine beiden GLPI-Projekte. Er liest die produktive Wissensbasis und Agent-Audits **read-only**, erzeugt einen dynamischen Wissensgraphen und rendert dessen Aktivität als fullscreen „Gehirn“. Neue KI-Synthesen werden ausschließlich als **AI-THINK** in das vorhandene Knowledgebase-Staging geschrieben.
Eigenständiger Go-Dienst für Agent, lokale Knowledgebase, GLPI-Knowledgebase und autonome AI-THINK-Anreicherung. Der Dienst hält einen vollständigen Wissensgraphen im Arbeitsspeicher und rendert eine hierarchisch verdichtete Fullscreen-Ansicht als lebende Hirnaktivität.
## Was bereits implementiert ist
## Kernfunktionen
- Fullscreen Canvas-Rendering in Gehirnform, ohne Frontend-Framework oder externe CDN-Abhängigkeit.
- Echtzeitaktivierung über Server-Sent Events: Nodes glühen, Edges leuchten, Partikel laufen entlang verwendeter Verbindungen.
- Direkter Ingest der vorhandenen Knowledge-JSONs einschließlich Kategorien, Keywords, Quellen und Staging-Status.
- Read-only Tailing von `runs.jsonl` des Agents; neue Agent-Läufe erscheinen als Hirnaktivität.
- Optionale, nicht blockierende Telemetrie-Patches für Agent- und Knowledgebase-Suchanfragen.
- Embeddings über Ollama `embeddinggemma`; lokaler Feature-Hash-Fallback, falls Ollama gerade nicht erreichbar ist.
- Suchanfragen über Agent, Knowledgebase oder API: Retrieval, aktivierte Nodes/Edges und strukturierte Verarbeitung durch `qwen3:8b`. Die Webansicht bleibt bewusst eine reine Visualisierung ohne Eingabefeld.
- Sequenzielle, automatische Verknüpfungsanalyse mit sichtbarem Worker-Status. Pro Zyklus werden standardmäßig bis zu drei Kandidaten nacheinander geprüft; niemals parallel.
- KI-Edges mit Herkunft, Confidence, Erklärung und Evidenz. Abgelehnte Paare werden intern markiert, damit sie nicht endlos erneut geprüft werden.
- Optional kontrollierte Recherche über eine eigene SearXNG-Instanz.
- Automatische AI-THINK-Beiträge im bestehenden Staging-JSON-Format, stets mit `auto_reply: false`.
- Fullscreen-Canvas mit 3D-Hirnform, semantischen Cortex-Regionen und hierarchischem Level-of-Detail.
- Echtzeitaktivierung über Server-Sent Events: Nodes glühen, aggregierte Edges leuchten und Partikel folgen tatsächlichen Wissenspfaden.
- Ingest lokaler produktiver Knowledge-JSONs sowie separater AI-THINK-Staging-Dateien.
- Optionaler read-only Ingest sichtbarer Beiträge aus der GLPI-Knowledgebase.
- Read-only Tailing der Agent-`runs.jsonl` und optionale Suchtelemetrie aus Agent und KB.
- Ollama-Pool mit mehreren unabhängigen Instanzen, Routing, Healthchecks, Cooldown und Failover.
- Embeddings über `embeddinggemma`, Beziehungsanalyse über `qwen3:8b`.
- Sequenzieller autonomer AI-THINK-Worker mit optionaler SearXNG-Recherche.
- Gebündelte Festplattenpersistenz: Graph, GLPI-KB-Cache und AI-THINK-Dateien werden standardmäßig nur alle fünf Minuten sequenziell geschrieben.
## Vertrauens- und Schreibgrenzen
## Automatische Visualzustände
Das Brain liest produktive Quellen und GLPI-KB ausschließlich. Es schreibt niemals in produktive Knowledge-Dateien oder nach GLPI zurück.
Die Fullscreen-Ansicht wechselt selbstständig zwischen drei Darstellungsstufen:
Schreibbar sind nur:
- **LIVING:** ruhige Eigenaktivität mit langsamer Atmung semantischer Cortex-Areale, vereinzelten internen Impulsen und sanfter Kamerabewegung. Diese Mikroaktivität wird nicht als wichtiges Feed-Ereignis protokolliert.
- **ACTIVATION:** Agent-Suchen, Knowledgebase-Suchen und Graph-Updates fokussieren automatisch den betroffenen Wissensbereich. Aktive Regionen dehnen sich leicht aus; relevante Edges transportieren Partikel.
- **AI-THINK / RESEARCH:** Beziehungsanalyse und Recherche erhalten einen stärkeren visuellen Modus mit fokussierter Kamera, konzentrischen Wellen, Synapsen-Bursts und statusabhängigen Farben.
- das eigene `BRAIN_DATA_DIR` für Graphzustand und GLPI-KB-Cache;
- der erste Pfad aus `BRAIN_STAGING_DIRS` für AI-THINK-Entwürfe.
Die Themenstruktur wird nicht nur aus dem ersten Kategorie-Feld abgeleitet. Kategorien bilden feste Anker; Konzepte, Quellen und externe Recherche-Nodes übernehmen über gewichtete Nachbarschafts-Propagation das stärkste verbundene Themengebiet. Verwandte Bereiche ziehen sich an, nicht verwandte Bereiche stoßen sich ab. Die zwölf stärksten Cortex-Areale erhalten bewusst deutlich getrennte Farben.
AI-THINK bleibt `auto_reply: false`, trägt die Kategorien `AI-THINK` und `AI-Staging` und wird erst nach einer expliziten Freigabe im bestehenden Editor produktiv.
Der Aktivitätsfeed zeigt nur relevante Ereignisse und ergänzt – sofern vorhanden – Cortex-Bereich, Trefferzahl, verwendete Quellen, Laufzeit, semantische Nähe, Relationstyp, Konfidenz, Recherchequellen, Ticket-ID und Staging-Pfad. Der Statusblock zeigt außerdem, ob der autonome Worker läuft, wartet, keine Kandidaten findet oder Ollama nicht erreicht. Über den wieder eingeblendeten **AI-THINK**-Button kann jederzeit ein manueller Zyklus angestoßen werden.
### Hierarchisches Semantic-LOD
Die Visualisierung rendert nicht mehr permanent jeden Knoten und jede Kante. Räumlich nahe, inaktive Elemente derselben Cortex-Region werden in zwei Stufen verdichtet:
- **Themenwolke (Level 2):** größere ruhende Bereiche;
- **lokale Wissensgruppe (Level 1):** kleinere Gruppen innerhalb einer Themenwolke;
- **Einzel-Node (Level 0):** konkrete Wissenseinträge bei Aktivität oder starkem Zoom.
Alle Originaldaten bleiben im Browser erhalten. Nur der abgeleitete `renderGraph` wird reduziert. Außenkanten einer Gruppe werden nach Quellgruppe, Zielgruppe, Relationstyp, Richtung, Herkunft und Status dedupliziert. `edgeCount` und `weightSum` bleiben an der aggregierten Kante erhalten; interne Kanten werden am Supernode gezählt.
Wird ein enthaltenes Wissenselement durch Agent, Knowledgebase, AI-THINK oder Recherche aktiviert, öffnet sich zunächst seine Themenwolke in engere Gruppen und anschließend die betroffene lokale Gruppe in Einzel-Nodes. Nach 30 bis 42 Sekunden ohne erneute Aktivität fällt der Bereich wieder zusammen. Zoom-Hysterese verhindert Flackern. Der **LOD**-Schalter kann die Verdichtung zu Diagnosezwecken deaktivieren; die Kennzahl **Render** im Header zeigt die tatsächlich pro Frame gezeichneten Nodes.
Die Grenzwerte stehen am Anfang von `internal/web/static/app.js` in `LOD_CONFIG` (`localDistance`, `coarseDistance`, Gruppengrößen und Offenhaltezeiten).
## Schutz der Basisprojekte
Das Brain bekommt nur:
- `knowledge/` **read-only**
- `data/runs.jsonl` **read-only**
- `staging/` **read-write**
Der Agent erhält weiterhin keinen Zugriff auf das Staging. Im integrierten Compose-Stack erhält auch `kb-search` nur ein leeres, flüchtiges Staging; ausschließlich der Prüf-Editor und das Brain sehen die echten Entwürfe. Die bestehende Knowledgebase nimmt AI-THINK erst nach deiner Freigabe in den produktiven Bestand. Damit kann das Brain die Entwürfe bereits darstellen und beim Denken berücksichtigen, während Agent und produktive Suche sie noch nicht sehen.
## Schnellstart nativ
## Schnellstart
```bash
cp .env.example .env
# Pfade in .env anpassen
ollama pull qwen3:8b
ollama pull embeddinggemma
# Pfade, Ollama-Pool und optional GLPI-Zugangsdaten anpassen.
set -a; . ./.env; set +a
go run ./cmd/brain
```
Windows PowerShell:
```powershell
Copy-Item .env.example .env
# Variablen aus .env setzen oder direkt in der Sitzung definieren
go run ./cmd/brain
```
Oberfläche: `http://localhost:8090`
Ohne Ollama startet die Visualisierung trotzdem. Retrieval verwendet dann einen deterministischen lokalen Fallback; Qwen-Synthesen und belastbare AI-Inferenz benötigen Ollama.
## Docker
Passe in `docker-compose.yml` die drei Host-Pfade an und starte:
Docker:
```bash
docker compose up -d --build
```
Auf Linux ist `host.docker.internal` über `extra_hosts` eingebunden. Alternativ kann das Brain in dasselbe Docker-Netz wie Ollama aufgenommen und `OLLAMA_URL=http://ollama:11434` gesetzt werden.
## Mehrere Ollama-Instanzen
## Ablauf einer sichtbaren Anfrage
`OLLAMA_URLS` hat Vorrang vor dem weiterhin unterstützten `OLLAMA_URL`.
1. Die Anfrage erzeugt eine Wahrnehmungswelle.
2. EmbeddingGemma bewertet passende Knowledge- und AI-THINK-Nodes.
3. Treffer leuchten nacheinander auf.
4. Vorhandene Verbindungen werden durchlaufen und mit Partikeln dargestellt.
5. Qwen3:8b erhält ausschließlich den ausgewählten Kontext.
6. Die final verwendeten Nodes und Edges pulsieren bei der Antwortsynthese.
```env
OLLAMA_URLS=http://10.20.30.21:11434,http://10.20.30.22:11434,http://10.20.30.23:11434
OLLAMA_NODE_NAMES=gpu-01,gpu-02,gpu-03
OLLAMA_NODE_WEIGHTS=1,1,4
OLLAMA_ROUTING_MODE=least_inflight
OLLAMA_NODE_MAX_INFLIGHT=1
OLLAMA_FAILOVER_ENABLED=true
OLLAMA_FAILOVER_ATTEMPTS=0
OLLAMA_REQUIRE_SAME_MODEL_DIGEST=true
OLLAMA_REQUIRE_EMBEDDING_MODEL=true
```
## Automatische Anreicherung
Unterstützte Routing-Modi:
Der Enrichment-Loop arbeitet bewusst seriell. Die Kandidatensuche verwendet einen rotierenden, begrenzten Anchor-Satz statt eines vollständigen O(n²)-Vergleichs über den gesamten Graphen. Dadurch bleibt sie auch bei zehntausenden Nodes reaktionsfähig und besucht langfristig trotzdem den gesamten Wissensbestand.
- `least_inflight` – bevorzugt den aktuell am wenigsten belasteten Node;
- `round_robin` – zyklische Verteilung;
- `weighted` – Verteilung entsprechend `OLLAMA_NODE_WEIGHTS`;
- `fastest_recent` – bevorzugt die zuletzt schnellsten Nodes.
1. einen rotierenden Bereich des Graphen nach dem stärksten noch ungeprüften Wissenspaar durchsuchen;
2. Qwen-Beziehungsanalyse mit festem JSON-Schema;
3. Edge als `staging` oder intern als `rejected` speichern;
4. bei Unklarheit optional SearXNG-Recherche durchführen;
5. externe Quellen als eigene Nodes mit Evidence-Edges anlegen;
6. AI-THINK-JSON atomar in `BRAIN_STAGING_DIRS` schreiben;
7. beim nächsten Scan den neuen Beitrag als sichtbaren und durchsuchbaren Staging-Node aufnehmen.
Jeder Node muss das Chatmodell besitzen. Mit `OLLAMA_REQUIRE_EMBEDDING_MODEL=true` muss außerdem jeder Node das Embeddingmodell besitzen. Bei aktivierter Digest-Prüfung wird ein Pool mit unterschiedlichen Modellständen fail-closed behandelt.
Ein erzeugter Entwurf enthält zusätzlich ein `ai_think`-Objekt mit Quell-Nodes, Relation, Confidence, Forschungsstatus und Evidenz. Die vorhandene Editor-Raw-JSON-Ansicht kann diese Daten bereits anzeigen.
Mehr Details: [`OLLAMA-POOL.md`](OLLAMA-POOL.md).
## GLPI-Knowledgebase einbeziehen
### AI-THINK-Taktung
Die Integration verwendet denselben read-only OAuth-/OpenAPI-Ansatz wie der Agent. GLPI entscheidet anhand des Service-Accounts, welche Beiträge sichtbar sind.
```env
GLPI_KB_ENABLED=true
GLPI_URL=https://glpi.example.org
GLPI_API_VERSION=v2.3
GLPI_CLIENT_ID=...
GLPI_CLIENT_SECRET=...
GLPI_USERNAME=...
GLPI_PASSWORD=...
GLPI_KB_PATH=auto
GLPI_KB_LIMIT=500
GLPI_KB_SYNC_INTERVAL=10m
GLPI_KB_SOURCE=GLPI Knowledge Base
```
Bei `GLPI_KB_PATH=auto` sucht das Brain im GLPI-OpenAPI-Dokument die lesbare `KnowbaseItem`-Collection. Die Beiträge werden als produktive Wissens-Nodes mit `glpi://KnowbaseItem/<id>`-URI aufgenommen. Kategorien und Quelle erzeugen verifizierte Edges. Es gibt keine GLPI-Schreiboperation.
Manueller Sync:
```bash
curl -X POST http://localhost:8090/api/glpi-kb/sync
```
Mehr Details: [`GLPI-KB.md`](GLPI-KB.md).
## Gebündelte Persistenz
```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:
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.
Manueller Flush:
```bash
curl -X POST http://localhost:8090/api/flush
```
Mehr Details: [`PERSISTENCE.md`](PERSISTENCE.md).
## Autonome Anreicherung
Der Worker arbeitet bewusst sequenziell:
1. Kandidatenpaar aus bestehenden Vektoren auswählen;
2. Beziehung durch Qwen mit festem JSON-Schema prüfen;
3. Edge als `staging` oder `rejected` im In-Memory-Graph ablegen;
4. optional kontrolliert recherchieren;
5. AI-THINK-Entwurf in die Schreibwarteschlange legen;
6. beim nächsten Persistenz-Flush atomar in Staging schreiben.
```env
BRAIN_AUTO_ENRICH=true
BRAIN_ENRICH_INTERVAL=90s
BRAIN_ENRICH_BATCH_SIZE=3
BRAIN_ENRICH_STEP_DELAY=3s
BRAIN_ENRICH_ANCHORS=48
```
`BRAIN_ENRICH_BATCH_SIZE` bestimmt, wie viele Beziehungen pro autonomem Zyklus nacheinander geprüft werden. `BRAIN_ENRICH_ANCHORS` begrenzt die CPU-seitige Kandidatensuche je Schritt. Die GPU wird nur während echter Qwen- oder Embedding-Aufrufe belastet; zwischen den Zyklen ist eine geringe oder null GPU-Auslastung normal.
## Minimale optionale Integrationen
Die Patches unter `integrations/` senden echte Suchanfragen und Treffer an `POST /api/events`:
Der manuelle AI-THINK-Button bleibt vorhanden. Alternativ:
```bash
# im jeweiligen Projekt-Root
git apply /pfad/glpi-neural-brain/integrations/agent/glpi-ai-agent-neural-brain.patch
git apply /pfad/glpi-neural-brain/integrations/knowledgebase/glpi-ai-knowledgebase-neural-brain.patch
curl -X POST 'http://localhost:8090/api/enrich?async=1'
```
Danach optional setzen:
```env
BRAIN_ACTIVITY_URL=http://brain:8090/api/events
BRAIN_ACTIVITY_API_KEY=
```
Ist `BRAIN_ACTIVITY_URL` leer, ist die Integration vollständig deaktiviert. Das Senden ist asynchron, fail-open, auf drei Sekunden begrenzt und kann weder Ticketverarbeitung noch KB-Suche blockieren. Der Agent funktioniert zusätzlich auch ohne Patch: Das Brain beobachtet weiterhin sein `runs.jsonl`.
## HTTP-Endpunkte
| Methode | Pfad | Zweck |
|---|---|---|
| `GET` | `/api/status` | Zustand, Modelle und Zähler |
| `GET` | `/api/graph` | kompletter aktueller Graph |
| `GET` | `/api/analysis` | Komponenten, Hubs, AI-Edges, Widersprüche und unverknüpftes Wissen |
| `GET` | `/api/status` | Gesamtstatus inklusive Ollama-Pool, GLPI-KB und Persistenzqueue |
| `GET` | `/api/graph` | vollständiger aktueller In-Memory-Graph |
| `GET` | `/api/analysis` | strukturelle Graphanalyse |
| `GET` | `/api/stream` | SSE-Aktivitätsstrom |
| `POST` | `/api/query` | sichtbare Wissensanfrage |
| `POST` | `/api/query` | programmatische Wissensanfrage; in der Fullscreen-UI verborgen |
| `POST` | `/api/events` | optionale Agent-/KB-Telemetrie |
| `POST` | `/api/reindex` | Scan und Embedding-Abgleich |
| `POST` | `/api/enrich` | einen AI-THINK-Schritt synchron ausführen |
| `POST` | `/api/enrich?async=1` | einen vollständigen sequenziellen AI-THINK-Zyklus einplanen |
| `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 |
Mit `BRAIN_API_KEY` werden alle POST-Endpunkte über `Authorization: Bearer …` oder `X-Brain-Key` geschützt. Die Webansicht ruft keine Query-POSTs mehr auf. Der API-Key schützt weiterhin Integrationen, Reindex, Enrichment und externe Query-Aufrufe; produktiv sollte der Dienst lokal oder hinter einem authentifizierenden Reverse Proxy betrieben werden.
Mit `BRAIN_API_KEY` werden POST-Endpunkte über `Authorization: Bearer …` oder `X-Brain-Key` geschützt.
## Grenzen des Prototyps
## Statusdiagnose
- Der Graphspeicher ist eine atomar geschriebene JSON-Datei und für einen einzelnen Brain-Prozess ausgelegt. Für sehr große Bestände wäre eine spätere Migration auf einen spezialisierten Graph-/Vektorspeicher sinnvoll.
- Webrecherche ist absichtlich nur über eine explizit konfigurierte SearXNG-Instanz aktiv.
- KI-Edges bleiben Hypothesen. Erst eine Freigabe des AI-THINK-Beitrags macht daraus produktives Knowledge; die Edge selbst trägt weiterhin ihre KI-Herkunft.
- Das System löst Widersprüche nicht stillschweigend auf. `contradicts` ist ein eigener Edge-Typ und bleibt sichtbar.
```bash
curl http://localhost:8090/api/status | jq
```
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;
- `enrich_*`: Zustand des autonomen AI-THINK-Workers.
## Validierung
```bash
go test ./...
go vet ./...
node --check internal/web/static/app.js
```
+29 -18
View File
@@ -1,34 +1,45 @@
b3704d36a5da4bca7faab55ad945a026e93843d169be7c942a3ebe66cb241e2e ./.env.example
cd2a385faa5059dc1c756533bce56666a124272682b738629831f042c7a0e0fb ./ARCHITECTURE.md
532a6031326fcbb3b9862042ef8a17dbec46dee87a873b595984d5092c807392 ./.env.example
97f83ffd2da3c3a89184d8bb460493f67eca5ab22d45dc7ea846fc6ea25ea5b9 ./ARCHITECTURE.md
91f75b87bd95480fd40f191b551fe2d5f9f64e590e7c5442c19ee6166a362e6e ./CHANGELOG-GLPI-POOL-PERSISTENCE.md
99171b262752a66278da5a4a8b2b48e83fbdb1eb9cdca3efa944d260d6e6bc01 ./Dockerfile
5534536965bf0479455f97324c242160202650ca1256f1ba0420b4ad67125e49 ./GLPI-KB.md
adcd3c9ba3bdf366afcc4e15a25423e068dd761e5d5d2d6f8cb20a3686302045 ./Makefile
221c2e42a6234fa3bae31eb71745182489f466fcedbc8f3d3e9d918ce6d35adf ./README.md
cf53db59fb1b3057019d54166bfbb0c95fbf84cb5032cad3ebea30af3aab69b3 ./cmd/brain/main.go
418310e2660f34890c4be97d6e4a0873cce57fe5b8eb6bfdff25bcd2d5faa7bf ./data/graph-state.json
2404066ac6c893852a3f2bba00b33883b7e0eeb542448caf733fff4862715ed0 ./deployment/README.md
2f3d69501091987783f8fd5e4d677d1191fd95fc435b12e907b0bcbf53d540aa ./deployment/docker-compose.full.yml
63b043e874d4cb8130d7bc5a8eca58cfb4cdd1069ff18c20a568ada0d373f9cd ./docker-compose.yml
381d7d6ac9e3c2e63c9ecdaa42ed4c73058f5d78e57c7532bb75a9663c919530 ./OLLAMA-POOL.md
10c14f4c08b9b84699cc4cf0dac3e6faf6c0ffd7e74b428463f14c628d0b533d ./PERSISTENCE.md
9de8804000cc41c7b5683b196e433caa4fb3e2e23c8c1cd1e5950c127d95dcb8 ./README.md
cfc2393e8300b792cec30c57b9dc7670b19dcf5abb190dddaa75f32b77cb1fc7 ./cmd/brain/main.go
21b51d0e1b7ed07c20f7f3a5da76dedab8df44a94a51724b67b0c3411599fe15 ./deployment/README.md
652c8188f9b591a085afdf95d0b56744c8fc618be82d8e6d03341c813f693ea7 ./deployment/docker-compose.full.yml
4b0a58e828cd3d808b6385be5d17eb06a6851526f2914aebf75759b492cec158 ./docker-compose.yml
9126f8bca1144abfc77747063b2c8f31acf12839a75e060b98369e10a00061ca ./go.mod
bf239391e61040b00d2f49d805b0d49a44bd3679f3c53849dfbc5e3b5a3fcc7a ./integrations/agent/README.md
a0105475dc054977223fac36618b8cd8137c55be1d11fddcd24e9a4d3074c170 ./integrations/agent/glpi-ai-agent-neural-brain.patch
73ab7c49600cdaa4795e76c50e66ce919dec171b3a07f2d16ff6f45ce5f73365 ./integrations/knowledgebase/README.md
36666043ebf4139e13610e5fcd0b0c6f45e4e6a53fed33ec47d03a66878e7b08 ./integrations/knowledgebase/glpi-ai-knowledgebase-neural-brain.patch
50d05fa2a183f5f3eaab0545cb48d3abb64be62eb5d84dc2c6d99c7125f7344b ./internal/activity/broker.go
094c75a5335170e300d9c88ad183db40fa54cc0070dad3f0fdbb62c719adc1e8 ./internal/config/config.go
686a2d9b0eb741e7ec163883e5e407e2b629e13b7b0f420158b5a71608345c0a ./internal/engine/engine.go
55af6324acec36d5a4ce9e6c48ade1dcba9a4acbd2eda00a572ab71ca576bfea ./internal/engine/engine_test.go
64f896f3421ffc41836e0dce7fdcf8beb639258960fee7df07552539c0db2038 ./internal/graph/store.go
c7f4c9702c5d4c0ab872001757a13564a065afa41ca2767475be42976bfb327f ./internal/config/config.go
6b6b9e92e1feee682a6663b1b617add06d53137de9305bdfc920c4536faeb866 ./internal/config/config_test.go
af76f57a74225b663bea7fbeb5ae4cf9d9bd8a71b292135e8a21c325923df76a ./internal/engine/engine.go
27c8ffa16cca14ab8a39298ecb11bc775cee6968d9d137a91ecc6a845064fd7e ./internal/engine/engine_test.go
b82980a646a92751bdd27a866ba1ffc6d34a3ba81d537f7b6e5a78e1432ee6fa ./internal/glpi/client.go
525102be56bc51ce8a08655b1b2bb53b67f4a1828903585fd664ed6a5133f617 ./internal/glpi/client_test.go
c41c845b3e342bf20a7453f16a34a8d1aa6c0f43424c2c86d6cf297ba4101381 ./internal/graph/store.go
5c32a4e4fca939165a2fa7c0f380fdf79cfab7cc65924e21fe7b9ac5e85a2b6d ./internal/graph/store_test.go
4476351d388d11c6becd78b4c918fc8d47dfd70500d8b3e2ddf81f7ff61a2cea ./internal/ingest/agent.go
fe0a813efe140fdc8961015cb2a493928c97e456a80c2d8cf2ddbe7c6e335018 ./internal/ingest/knowledge.go
5ed96aec4c3d2599cb9af3f6e951533271dc2b9a4826667482da814f9b450da1 ./internal/ingest/knowledge_test.go
a64fde7d9be8841b363bc5cf5e41c81ade8e068335f7b04f41272e5bc3628d93 ./internal/ingest/glpikb.go
774e9155eb7d17a208d483625f9fb6b60b9f9d44a70843c9cc51e29cf147a2b7 ./internal/ingest/glpikb_test.go
c0470d74bd3c3369cfefa6ba2343abcbdaa584a8898011704cba7d1cee620141 ./internal/ingest/knowledge.go
ebd47d134e61badf5e1eb35bdf1f45a3e4c9f99dcbf0a022850409bebc587e56 ./internal/ingest/knowledge_test.go
fb1af93afeb4ebc89cad95badeaee9a0cc3a236b9761e15a41a7b4b2c7192d00 ./internal/model/model.go
2ea6284f6aa2c23c6f6583ba66647ae8f48432527746b99898bc567c1be64a52 ./internal/ollama/client.go
db7a99b832bc41c61584cd726d5fc7f4f250197d716d58bc27e715923f7f5426 ./internal/ollama/client.go
61089700aa5912b58bff526b71c2a1bda74077b802f9ff69a0d7dc89c1e3e3a1 ./internal/ollama/client_test.go
e61a9a426b96851b409ac357bdc53853324b5c5da7568ff6f89fe9cc09f83a73 ./internal/persist/coordinator.go
d33318b43388f134358cf40f5b0f130eb10edcca06011955ea0544897b4e4ad1 ./internal/persist/coordinator_test.go
2eb686cfc9016b9b0cf15e5c342f693ec9b60c454bc4814c1c9583f6d5b5b806 ./internal/research/searxng.go
6346f7b213aa3fb36bc9f43134bba75e0b368c9a005cc1aef8d265f553ff8cef ./internal/research/searxng_test.go
dda84763f509405c39cb533ee1c1f2bbeed48718cd7de9cbb66dd1adc5582e2c ./internal/web/server.go
1a8db923dcfa3416edd015b08c71cc48145f6bf7b68429239134d58895ba1af3 ./internal/web/server.go
644105a4921f0394bc0ab3cc4c8658053e456e593acf7573255bd867d6cde8dd ./internal/web/static/app.css
c310356d4bfe2c02a7bf078140c9c49febf8a7c4f9ab2dfacd1016e918953b0d ./internal/web/static/app.js
df87659f906533144dd407a488d2d2abbf5ea57d95a11fcec70aa4c51729a5d5 ./internal/web/static/app.js
02007a3143603cda56794e36613f9a3cfc3cdbaf66279af0b3e4b9b5e8f4cbc2 ./internal/web/static/index.html
fcc7348cf65c04d023b0d0c1fa391e7628ae326d394e44c6fde9bb8eb6f665b5 ./neural-brain
39656b3c5e7f013545aa6a6b01c213f8f746a947f7158c3042b81a07429d5d2a ./neural-brain
83aded814b6225395935e61fe957963c3c470f368fc9089f505b6de23e959115 ./preview.png
+1 -1
View File
@@ -49,6 +49,6 @@ func main() {
<-ctx.Done()
shutdown, c := context.WithTimeout(context.Background(), 10*time.Second)
defer c()
_ = g.Persist()
_ = eng.Flush(shutdown)
_ = srv.Shutdown(shutdown)
}
File diff suppressed because one or more lines are too long
+19
View File
@@ -22,3 +22,22 @@ Ports: Agent `8080`, Editor `8081`, Suche `8082`, Brain `8090`.
## Staging-Isolation
Der Editor und das Brain teilen sich das echte Knowledgebase-Staging. `kb-search` erhält absichtlich nur ein leeres, flüchtiges Verzeichnis unter `/tmp/empty-staging`; der Agent mountet das Staging überhaupt nicht. AI-THINK-Drafts sind dadurch im Brain und im Prüf-Editor sichtbar, aber weder für Agent-Antworten noch für die produktive Knowledgebase-Suche verfügbar. Erst die bestehende Promotion im Editor kopiert einen freigegebenen Entwurf in die produktive Knowledge-Basis.
## Externer Ollama-Pool
Für mehrere Instanzen `OLLAMA_URLS`, `OLLAMA_NODE_NAMES` und optional `OLLAMA_NODE_WEIGHTS` in der Umgebung des Compose-Aufrufs setzen. Die lokale `ollama`-Service-Instanz ist keine harte Startabhängigkeit mehr; sie wird nur verwendet, wenn ihre URL im effektiven Pool steht.
```env
OLLAMA_URLS=http://gpu-01:11434,http://gpu-02:11434
OLLAMA_NODE_NAMES=gpu-01,gpu-02
OLLAMA_NODE_WEIGHTS=1,2
OLLAMA_ROUTING_MODE=least_inflight
```
## GLPI-KB für das Brain
`GLPI_KB_ENABLED=true` aktiviert einen zusätzlichen read-only Sync aus GLPI. Die gleichen GLPI-Zugangsdaten wie beim Agenten können verwendet werden, sofern der Service-Account Leserechte auf die Knowledgebase besitzt. Das Brain schreibt nicht nach GLPI zurück.
## Persistenz
`BRAIN_PERSIST_INTERVAL=5m` bündelt Graph-, Cache- und Staging-Schreibvorgänge. Bei einem geplanten Container-Stopp wird zusätzlich ein finaler Flush ausgeführt.
+32 -8
View File
@@ -25,6 +25,11 @@ services:
DATA_DIR: /app/data
KNOWLEDGE_DIR: /app/knowledge
OLLAMA_URL: http://ollama:11434
OLLAMA_URLS: ${OLLAMA_URLS:-http://ollama:11434}
OLLAMA_NODE_NAMES: ${OLLAMA_NODE_NAMES:-local-gpu}
OLLAMA_NODE_WEIGHTS: ${OLLAMA_NODE_WEIGHTS:-1}
OLLAMA_ROUTING_MODE: ${OLLAMA_ROUTING_MODE:-least_inflight}
OLLAMA_NODE_MAX_INFLIGHT: ${OLLAMA_NODE_MAX_INFLIGHT:-1}
BRAIN_ACTIVITY_URL: http://brain:8090/api/events
BRAIN_ACTIVITY_API_KEY: ${BRAIN_API_KEY:-}
ports:
@@ -35,8 +40,6 @@ services:
depends_on:
agent-data-init:
condition: service_completed_successfully
ollama:
condition: service_started
security_opt: [no-new-privileges:true]
cap_drop: [ALL]
read_only: true
@@ -110,24 +113,45 @@ services:
BRAIN_AGENT_RUNS_FILES: /sources/agent-data/runs.jsonl
BRAIN_API_KEY: ${BRAIN_API_KEY:-}
OLLAMA_URL: http://ollama:11434
OLLAMA_URLS: ${OLLAMA_URLS:-http://ollama:11434}
OLLAMA_NODE_NAMES: ${OLLAMA_NODE_NAMES:-local-gpu}
OLLAMA_NODE_WEIGHTS: ${OLLAMA_NODE_WEIGHTS:-1}
OLLAMA_ROUTING_MODE: ${OLLAMA_ROUTING_MODE:-least_inflight}
OLLAMA_NODE_MAX_INFLIGHT: ${OLLAMA_NODE_MAX_INFLIGHT:-1}
OLLAMA_NODE_HEALTH_INTERVAL: ${OLLAMA_NODE_HEALTH_INTERVAL:-15s}
OLLAMA_NODE_FAILURE_COOLDOWN: ${OLLAMA_NODE_FAILURE_COOLDOWN:-30s}
OLLAMA_NODE_REQUEST_TIMEOUT: ${OLLAMA_NODE_REQUEST_TIMEOUT:-8m}
OLLAMA_FAILOVER_ENABLED: ${OLLAMA_FAILOVER_ENABLED:-true}
OLLAMA_FAILOVER_ATTEMPTS: ${OLLAMA_FAILOVER_ATTEMPTS:-0}
OLLAMA_REQUIRE_SAME_MODEL_DIGEST: ${OLLAMA_REQUIRE_SAME_MODEL_DIGEST:-true}
OLLAMA_REQUIRE_EMBEDDING_MODEL: ${OLLAMA_REQUIRE_EMBEDDING_MODEL:-true}
OLLAMA_CHAT_MODEL: ${OLLAMA_CHAT_MODEL:-qwen3:8b}
OLLAMA_EMBEDDING_MODEL: ${OLLAMA_EMBEDDING_MODEL:-embeddinggemma}
BRAIN_AUTO_ENRICH: ${BRAIN_AUTO_ENRICH:-true}
BRAIN_SCAN_INTERVAL: ${BRAIN_SCAN_INTERVAL:-20s}
BRAIN_PERSIST_INTERVAL: ${BRAIN_PERSIST_INTERVAL:-5m}
BRAIN_ENRICH_INTERVAL: ${BRAIN_ENRICH_INTERVAL:-90s}
BRAIN_ENRICH_BATCH_SIZE: ${BRAIN_ENRICH_BATCH_SIZE:-3}
BRAIN_ENRICH_STEP_DELAY: ${BRAIN_ENRICH_STEP_DELAY:-3s}
BRAIN_ENRICH_ANCHORS: ${BRAIN_ENRICH_ANCHORS:-48}
BRAIN_RESEARCH_ENABLED: ${BRAIN_RESEARCH_ENABLED:-false}
SEARXNG_URL: ${SEARXNG_URL:-}
GLPI_KB_ENABLED: ${GLPI_KB_ENABLED:-false}
GLPI_URL: ${GLPI_URL:-}
GLPI_API_VERSION: ${GLPI_API_VERSION:-v2.3}
GLPI_CLIENT_ID: ${GLPI_CLIENT_ID:-}
GLPI_CLIENT_SECRET: ${GLPI_CLIENT_SECRET:-}
GLPI_USERNAME: ${GLPI_USERNAME:-}
GLPI_PASSWORD: ${GLPI_PASSWORD:-}
GLPI_ALLOW_INSECURE_HTTP: ${GLPI_ALLOW_INSECURE_HTTP:-false}
GLPI_TIMEOUT: ${GLPI_TIMEOUT:-20s}
GLPI_KB_PATH: ${GLPI_KB_PATH:-auto}
GLPI_KB_FILTER: ${GLPI_KB_FILTER:-}
GLPI_KB_LIMIT: ${GLPI_KB_LIMIT:-500}
GLPI_KB_SYNC_INTERVAL: ${GLPI_KB_SYNC_INTERVAL:-10m}
GLPI_KB_SOURCE: ${GLPI_KB_SOURCE:-GLPI Knowledge Base}
volumes:
- brain-data:/app/data
- ${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
depends_on:
ollama:
condition: service_started
ollama:
image: ollama/ollama:latest
+27
View File
@@ -12,10 +12,23 @@ services:
BRAIN_STAGING_DIRS: /sources/staging
BRAIN_AGENT_RUNS_FILES: /sources/agent-data/runs.jsonl
OLLAMA_URL: ${OLLAMA_URL:-http://host.docker.internal:11434}
OLLAMA_URLS: ${OLLAMA_URLS:-http://host.docker.internal:11434}
OLLAMA_NODE_NAMES: ${OLLAMA_NODE_NAMES:-local-gpu}
OLLAMA_NODE_WEIGHTS: ${OLLAMA_NODE_WEIGHTS:-1}
OLLAMA_ROUTING_MODE: ${OLLAMA_ROUTING_MODE:-least_inflight}
OLLAMA_NODE_MAX_INFLIGHT: ${OLLAMA_NODE_MAX_INFLIGHT:-1}
OLLAMA_NODE_HEALTH_INTERVAL: ${OLLAMA_NODE_HEALTH_INTERVAL:-15s}
OLLAMA_NODE_FAILURE_COOLDOWN: ${OLLAMA_NODE_FAILURE_COOLDOWN:-30s}
OLLAMA_NODE_REQUEST_TIMEOUT: ${OLLAMA_NODE_REQUEST_TIMEOUT:-8m}
OLLAMA_FAILOVER_ENABLED: ${OLLAMA_FAILOVER_ENABLED:-true}
OLLAMA_FAILOVER_ATTEMPTS: ${OLLAMA_FAILOVER_ATTEMPTS:-0}
OLLAMA_REQUIRE_SAME_MODEL_DIGEST: ${OLLAMA_REQUIRE_SAME_MODEL_DIGEST:-true}
OLLAMA_REQUIRE_EMBEDDING_MODEL: ${OLLAMA_REQUIRE_EMBEDDING_MODEL:-true}
OLLAMA_CHAT_MODEL: ${OLLAMA_CHAT_MODEL:-qwen3:8b}
OLLAMA_EMBEDDING_MODEL: ${OLLAMA_EMBEDDING_MODEL:-embeddinggemma}
BRAIN_AUTO_ENRICH: ${BRAIN_AUTO_ENRICH:-true}
BRAIN_SCAN_INTERVAL: ${BRAIN_SCAN_INTERVAL:-20s}
BRAIN_PERSIST_INTERVAL: ${BRAIN_PERSIST_INTERVAL:-5m}
BRAIN_ENRICH_INTERVAL: ${BRAIN_ENRICH_INTERVAL:-90s}
BRAIN_ENRICH_BATCH_SIZE: ${BRAIN_ENRICH_BATCH_SIZE:-3}
BRAIN_ENRICH_STEP_DELAY: ${BRAIN_ENRICH_STEP_DELAY:-3s}
@@ -25,6 +38,20 @@ services:
BRAIN_RESEARCH_ENABLED: ${BRAIN_RESEARCH_ENABLED:-false}
SEARXNG_URL: ${SEARXNG_URL:-}
BRAIN_API_KEY: ${BRAIN_API_KEY:-}
GLPI_KB_ENABLED: ${GLPI_KB_ENABLED:-false}
GLPI_URL: ${GLPI_URL:-}
GLPI_API_VERSION: ${GLPI_API_VERSION:-v2.3}
GLPI_CLIENT_ID: ${GLPI_CLIENT_ID:-}
GLPI_CLIENT_SECRET: ${GLPI_CLIENT_SECRET:-}
GLPI_USERNAME: ${GLPI_USERNAME:-}
GLPI_PASSWORD: ${GLPI_PASSWORD:-}
GLPI_ALLOW_INSECURE_HTTP: ${GLPI_ALLOW_INSECURE_HTTP:-false}
GLPI_TIMEOUT: ${GLPI_TIMEOUT:-20s}
GLPI_KB_PATH: ${GLPI_KB_PATH:-auto}
GLPI_KB_FILTER: ${GLPI_KB_FILTER:-}
GLPI_KB_LIMIT: ${GLPI_KB_LIMIT:-500}
GLPI_KB_SYNC_INTERVAL: ${GLPI_KB_SYNC_INTERVAL:-10m}
GLPI_KB_SOURCE: ${GLPI_KB_SOURCE:-GLPI Knowledge Base}
volumes:
- brain-data:/app/data
# Adjust the three host paths to your actual project directories.
+185 -42
View File
@@ -2,6 +2,7 @@ package config
import (
"fmt"
"net/url"
"os"
"path/filepath"
"strconv"
@@ -10,27 +11,55 @@ import (
)
type Config struct {
ListenAddr string
DataDir string
KnowledgeDirs []string
StagingDirs []string
AgentRunsFiles []string
OllamaURL string
ChatModel string
EmbeddingModel string
SearXNGURL string
ScanInterval time.Duration
EnrichInterval time.Duration
EnrichStepDelay time.Duration
EnrichBatchSize int
EnrichAnchors int
SimilarityThreshold float64
RelationThreshold float64
TopK int
MaxContextChars int
AutoEnrich bool
ResearchEnabled bool
APIKey string
ListenAddr string
DataDir string
KnowledgeDirs []string
StagingDirs []string
AgentRunsFiles []string
OllamaURL string
OllamaURLs []string
OllamaNodeNames []string
OllamaNodeWeights []int
OllamaRoutingMode string
OllamaNodeMaxInflight int
OllamaHealthInterval time.Duration
OllamaFailureCooldown time.Duration
OllamaRequestTimeout time.Duration
OllamaFailoverEnabled bool
OllamaFailoverAttempts int
OllamaRequireSameDigest bool
OllamaRequireEmbeddingModel bool
ChatModel string
EmbeddingModel string
SearXNGURL string
ScanInterval time.Duration
PersistInterval time.Duration
EnrichInterval time.Duration
EnrichStepDelay time.Duration
EnrichBatchSize int
EnrichAnchors int
SimilarityThreshold float64
RelationThreshold float64
TopK int
MaxContextChars int
AutoEnrich bool
ResearchEnabled bool
APIKey string
GLPIKBEnabled bool
GLPIURL string
GLPIAPIVersion string
GLPIClientID string
GLPIClientSecret string
GLPIUsername string
GLPIPassword string
GLPIAllowInsecure bool
GLPITimeout time.Duration
GLPIKBPath string
GLPIKBFilter string
GLPIKBLimit int
GLPIKBSyncInterval time.Duration
GLPIKBSource string
}
func Load() (Config, error) {
@@ -39,32 +68,70 @@ func Load() (Config, error) {
if err != nil {
return Config{}, err
}
legacyOllamaURL := strings.TrimRight(env("OLLAMA_URL", "http://localhost:11434"), "/")
ollamaURLs := stringList("OLLAMA_URLS")
if len(ollamaURLs) == 0 {
ollamaURLs = []string{legacyOllamaURL}
}
for i := range ollamaURLs {
ollamaURLs[i] = strings.TrimRight(ollamaURLs[i], "/")
}
cfg := Config{
ListenAddr: env("BRAIN_LISTEN_ADDR", ":8090"),
DataDir: abs,
KnowledgeDirs: paths("BRAIN_KNOWLEDGE_DIRS"),
StagingDirs: paths("BRAIN_STAGING_DIRS"),
AgentRunsFiles: paths("BRAIN_AGENT_RUNS_FILES"),
OllamaURL: strings.TrimRight(env("OLLAMA_URL", "http://localhost:11434"), "/"),
ChatModel: env("OLLAMA_CHAT_MODEL", "qwen3:8b"),
EmbeddingModel: env("OLLAMA_EMBEDDING_MODEL", "embeddinggemma"),
SearXNGURL: strings.TrimRight(strings.TrimSpace(os.Getenv("SEARXNG_URL")), "/"),
ScanInterval: duration("BRAIN_SCAN_INTERVAL", 20*time.Second),
EnrichInterval: duration("BRAIN_ENRICH_INTERVAL", 90*time.Second),
EnrichStepDelay: duration("BRAIN_ENRICH_STEP_DELAY", 3*time.Second),
EnrichBatchSize: integer("BRAIN_ENRICH_BATCH_SIZE", 3),
EnrichAnchors: integer("BRAIN_ENRICH_ANCHORS", 48),
SimilarityThreshold: number("BRAIN_SIMILARITY_THRESHOLD", 0.68),
RelationThreshold: number("BRAIN_RELATION_THRESHOLD", 0.72),
TopK: integer("BRAIN_TOP_K", 8),
MaxContextChars: integer("BRAIN_MAX_CONTEXT_CHARS", 16000),
AutoEnrich: boolean("BRAIN_AUTO_ENRICH", true),
ResearchEnabled: boolean("BRAIN_RESEARCH_ENABLED", false),
APIKey: strings.TrimSpace(os.Getenv("BRAIN_API_KEY")),
ListenAddr: env("BRAIN_LISTEN_ADDR", ":8090"),
DataDir: abs,
KnowledgeDirs: paths("BRAIN_KNOWLEDGE_DIRS"),
StagingDirs: paths("BRAIN_STAGING_DIRS"),
AgentRunsFiles: paths("BRAIN_AGENT_RUNS_FILES"),
OllamaURL: legacyOllamaURL,
OllamaURLs: ollamaURLs,
OllamaNodeNames: stringList("OLLAMA_NODE_NAMES"),
OllamaNodeWeights: intList("OLLAMA_NODE_WEIGHTS"),
OllamaRoutingMode: strings.ToLower(env("OLLAMA_ROUTING_MODE", "least_inflight")),
OllamaNodeMaxInflight: integer("OLLAMA_NODE_MAX_INFLIGHT", 1),
OllamaHealthInterval: duration("OLLAMA_NODE_HEALTH_INTERVAL", 15*time.Second),
OllamaFailureCooldown: duration("OLLAMA_NODE_FAILURE_COOLDOWN", 30*time.Second),
OllamaRequestTimeout: duration("OLLAMA_NODE_REQUEST_TIMEOUT", 8*time.Minute),
OllamaFailoverEnabled: boolean("OLLAMA_FAILOVER_ENABLED", true),
OllamaFailoverAttempts: integer("OLLAMA_FAILOVER_ATTEMPTS", 0),
OllamaRequireSameDigest: boolean("OLLAMA_REQUIRE_SAME_MODEL_DIGEST", true),
OllamaRequireEmbeddingModel: boolean("OLLAMA_REQUIRE_EMBEDDING_MODEL", true),
ChatModel: env("OLLAMA_CHAT_MODEL", "qwen3:8b"),
EmbeddingModel: env("OLLAMA_EMBEDDING_MODEL", "embeddinggemma"),
SearXNGURL: strings.TrimRight(strings.TrimSpace(os.Getenv("SEARXNG_URL")), "/"),
ScanInterval: duration("BRAIN_SCAN_INTERVAL", 20*time.Second),
PersistInterval: duration("BRAIN_PERSIST_INTERVAL", 5*time.Minute),
EnrichInterval: duration("BRAIN_ENRICH_INTERVAL", 90*time.Second),
EnrichStepDelay: duration("BRAIN_ENRICH_STEP_DELAY", 3*time.Second),
EnrichBatchSize: integer("BRAIN_ENRICH_BATCH_SIZE", 3),
EnrichAnchors: integer("BRAIN_ENRICH_ANCHORS", 48),
SimilarityThreshold: number("BRAIN_SIMILARITY_THRESHOLD", 0.68),
RelationThreshold: number("BRAIN_RELATION_THRESHOLD", 0.72),
TopK: integer("BRAIN_TOP_K", 8),
MaxContextChars: integer("BRAIN_MAX_CONTEXT_CHARS", 16000),
AutoEnrich: boolean("BRAIN_AUTO_ENRICH", true),
ResearchEnabled: boolean("BRAIN_RESEARCH_ENABLED", false),
APIKey: strings.TrimSpace(os.Getenv("BRAIN_API_KEY")),
GLPIKBEnabled: boolean("GLPI_KB_ENABLED", false),
GLPIURL: strings.TrimRight(strings.TrimSpace(os.Getenv("GLPI_URL")), "/"),
GLPIAPIVersion: env("GLPI_API_VERSION", "v2.3"),
GLPIClientID: strings.TrimSpace(os.Getenv("GLPI_CLIENT_ID")),
GLPIClientSecret: strings.TrimSpace(os.Getenv("GLPI_CLIENT_SECRET")),
GLPIUsername: strings.TrimSpace(os.Getenv("GLPI_USERNAME")),
GLPIPassword: strings.TrimSpace(os.Getenv("GLPI_PASSWORD")),
GLPIAllowInsecure: boolean("GLPI_ALLOW_INSECURE_HTTP", false),
GLPITimeout: duration("GLPI_TIMEOUT", 20*time.Second),
GLPIKBPath: env("GLPI_KB_PATH", "auto"),
GLPIKBFilter: strings.TrimSpace(os.Getenv("GLPI_KB_FILTER")),
GLPIKBLimit: integer("GLPI_KB_LIMIT", 500),
GLPIKBSyncInterval: duration("GLPI_KB_SYNC_INTERVAL", 10*time.Minute),
GLPIKBSource: env("GLPI_KB_SOURCE", "GLPI Knowledge Base"),
}
if cfg.ScanInterval < 2*time.Second {
return Config{}, fmt.Errorf("BRAIN_SCAN_INTERVAL must be at least 2s")
}
if cfg.PersistInterval < 10*time.Second {
return Config{}, fmt.Errorf("BRAIN_PERSIST_INTERVAL must be at least 10s")
}
if cfg.EnrichInterval < 10*time.Second {
return Config{}, fmt.Errorf("BRAIN_ENRICH_INTERVAL must be at least 10s")
}
@@ -86,6 +153,60 @@ func Load() (Config, error) {
if cfg.ResearchEnabled && cfg.SearXNGURL == "" {
return Config{}, fmt.Errorf("BRAIN_RESEARCH_ENABLED requires SEARXNG_URL")
}
if len(cfg.OllamaURLs) < 1 || len(cfg.OllamaURLs) > 64 {
return Config{}, fmt.Errorf("OLLAMA_URLS must contain between 1 and 64 nodes")
}
for _, raw := range cfg.OllamaURLs {
u, err := url.Parse(raw)
if err != nil || u.Scheme == "" || u.Host == "" || (u.Scheme != "http" && u.Scheme != "https") {
return Config{}, fmt.Errorf("invalid Ollama URL %q", raw)
}
}
if len(cfg.OllamaNodeNames) > 0 && len(cfg.OllamaNodeNames) != len(cfg.OllamaURLs) {
return Config{}, fmt.Errorf("OLLAMA_NODE_NAMES must contain one name per OLLAMA_URLS entry")
}
if len(cfg.OllamaNodeWeights) > 0 && len(cfg.OllamaNodeWeights) != len(cfg.OllamaURLs) {
return Config{}, fmt.Errorf("OLLAMA_NODE_WEIGHTS must contain one weight per OLLAMA_URLS entry")
}
for _, weight := range cfg.OllamaNodeWeights {
if weight < 1 || weight > 1000 {
return Config{}, fmt.Errorf("OLLAMA_NODE_WEIGHTS values must be between 1 and 1000")
}
}
switch cfg.OllamaRoutingMode {
case "least_inflight", "round_robin", "weighted", "fastest_recent":
default:
return Config{}, fmt.Errorf("invalid OLLAMA_ROUTING_MODE %q", cfg.OllamaRoutingMode)
}
if cfg.OllamaNodeMaxInflight < 1 || cfg.OllamaNodeMaxInflight > 64 {
return Config{}, fmt.Errorf("OLLAMA_NODE_MAX_INFLIGHT must be between 1 and 64")
}
if cfg.OllamaHealthInterval < time.Second || cfg.OllamaFailureCooldown < time.Second || cfg.OllamaRequestTimeout < time.Second {
return Config{}, fmt.Errorf("Ollama health, cooldown and request timeout must be at least 1s")
}
if cfg.GLPIKBEnabled {
u, err := url.Parse(cfg.GLPIURL)
if err != nil || u.Scheme == "" || u.Host == "" {
return Config{}, fmt.Errorf("GLPI_URL must be a valid absolute URL when GLPI_KB_ENABLED=true")
}
if u.Scheme != "https" && !(u.Scheme == "http" && cfg.GLPIAllowInsecure) {
return Config{}, fmt.Errorf("GLPI_URL must use https unless GLPI_ALLOW_INSECURE_HTTP=true")
}
for name, value := range map[string]string{"GLPI_CLIENT_ID": cfg.GLPIClientID, "GLPI_CLIENT_SECRET": cfg.GLPIClientSecret, "GLPI_USERNAME": cfg.GLPIUsername, "GLPI_PASSWORD": cfg.GLPIPassword} {
if value == "" {
return Config{}, fmt.Errorf("%s is required when GLPI_KB_ENABLED=true", name)
}
}
if cfg.GLPIKBLimit < 1 || cfg.GLPIKBLimit > 5000 {
return Config{}, fmt.Errorf("GLPI_KB_LIMIT must be between 1 and 5000")
}
if cfg.GLPIKBSyncInterval < time.Minute {
return Config{}, fmt.Errorf("GLPI_KB_SYNC_INTERVAL must be at least 1m")
}
if cfg.GLPIKBPath != "auto" && !strings.HasPrefix(cfg.GLPIKBPath, "/") {
return Config{}, fmt.Errorf("GLPI_KB_PATH must be auto or an absolute API path")
}
}
if err := os.MkdirAll(cfg.DataDir, 0o750); err != nil {
return Config{}, err
}
@@ -109,6 +230,28 @@ func paths(k string) []string {
}
return out
}
func stringList(k string) []string {
var out []string
for _, value := range strings.Split(os.Getenv(k), ",") {
if value = strings.TrimSpace(value); value != "" {
out = append(out, value)
}
}
return out
}
func intList(k string) []int {
var out []int
for _, value := range stringList(k) {
n, err := strconv.Atoi(value)
if err != nil {
out = append(out, 0)
continue
}
out = append(out, n)
}
return out
}
func duration(k string, d time.Duration) time.Duration {
v := strings.TrimSpace(os.Getenv(k))
if v == "" {
+39
View File
@@ -0,0 +1,39 @@
package config
import (
"testing"
)
func TestLoadOllamaPoolAndGLPIKB(t *testing.T) {
t.Setenv("BRAIN_DATA_DIR", t.TempDir())
t.Setenv("OLLAMA_URLS", "http://gpu-a:11434,http://gpu-b:11434")
t.Setenv("OLLAMA_NODE_NAMES", "a,b")
t.Setenv("OLLAMA_NODE_WEIGHTS", "1,3")
t.Setenv("OLLAMA_ROUTING_MODE", "weighted")
t.Setenv("GLPI_KB_ENABLED", "true")
t.Setenv("GLPI_URL", "http://glpi.test")
t.Setenv("GLPI_ALLOW_INSECURE_HTTP", "true")
t.Setenv("GLPI_CLIENT_ID", "id")
t.Setenv("GLPI_CLIENT_SECRET", "secret")
t.Setenv("GLPI_USERNAME", "user")
t.Setenv("GLPI_PASSWORD", "password")
cfg, err := Load()
if err != nil {
t.Fatal(err)
}
if len(cfg.OllamaURLs) != 2 || cfg.OllamaNodeNames[1] != "b" || cfg.OllamaNodeWeights[1] != 3 || cfg.OllamaRoutingMode != "weighted" {
t.Fatalf("unexpected pool config: %+v", cfg)
}
if !cfg.GLPIKBEnabled || cfg.GLPIURL != "http://glpi.test" {
t.Fatalf("unexpected GLPI config: %+v", cfg)
}
}
func TestLoadRejectsMismatchedPoolMetadata(t *testing.T) {
t.Setenv("BRAIN_DATA_DIR", t.TempDir())
t.Setenv("OLLAMA_URLS", "http://gpu-a:11434,http://gpu-b:11434")
t.Setenv("OLLAMA_NODE_NAMES", "only-one")
if _, err := Load(); err == nil {
t.Fatal("expected pool metadata validation error")
}
}
+75 -28
View File
@@ -19,10 +19,12 @@ import (
"github.com/local/glpi-neural-brain/internal/activity"
"github.com/local/glpi-neural-brain/internal/config"
"github.com/local/glpi-neural-brain/internal/glpi"
"github.com/local/glpi-neural-brain/internal/graph"
"github.com/local/glpi-neural-brain/internal/ingest"
"github.com/local/glpi-neural-brain/internal/model"
"github.com/local/glpi-neural-brain/internal/ollama"
"github.com/local/glpi-neural-brain/internal/persist"
"github.com/local/glpi-neural-brain/internal/research"
)
@@ -37,12 +39,14 @@ type EnrichOutcome struct {
}
type Engine struct {
Cfg config.Config
Graph *graph.Store
Broker *activity.Broker
Ollama *ollama.Client
Research *research.Client
Scanner *ingest.KnowledgeScanner
Cfg config.Config
Graph *graph.Store
Broker *activity.Broker
Ollama *ollama.Client
Research *research.Client
Scanner *ingest.KnowledgeScanner
GLPIKB *ingest.GLPIKBSyncer
Persistence *persist.Coordinator
mu sync.Mutex
stateMu sync.RWMutex
@@ -68,13 +72,46 @@ func New(cfg config.Config, g *graph.Store, b *activity.Broker) *Engine {
if cfg.EnrichAnchors < 1 {
cfg.EnrichAnchors = 48
}
e := &Engine{Cfg: cfg, Graph: g, Broker: b, Ollama: ollama.New(cfg.OllamaURL, cfg.ChatModel, cfg.EmbeddingModel), Scanner: &ingest.KnowledgeScanner{Graph: g, ProductionDirs: cfg.KnowledgeDirs, StagingDirs: cfg.StagingDirs}, enrichRequests: make(chan string, 1)}
ollamaURLs := append([]string(nil), cfg.OllamaURLs...)
if len(ollamaURLs) == 0 && strings.TrimSpace(cfg.OllamaURL) != "" {
ollamaURLs = []string{cfg.OllamaURL}
}
nodes := make([]ollama.NodeConfig, 0, len(ollamaURLs))
for i, rawURL := range ollamaURLs {
name := fmt.Sprintf("ollama-%d", i+1)
if i < len(cfg.OllamaNodeNames) && strings.TrimSpace(cfg.OllamaNodeNames[i]) != "" {
name = strings.TrimSpace(cfg.OllamaNodeNames[i])
}
weight := 1
if i < len(cfg.OllamaNodeWeights) && cfg.OllamaNodeWeights[i] > 0 {
weight = cfg.OllamaNodeWeights[i]
}
nodes = append(nodes, ollama.NodeConfig{Name: name, URL: rawURL, Weight: weight})
}
pool := ollama.NewPool(ollama.PoolConfig{
Nodes: nodes, RoutingMode: cfg.OllamaRoutingMode, NodeMaxInflight: cfg.OllamaNodeMaxInflight,
HealthInterval: cfg.OllamaHealthInterval, FailureCooldown: cfg.OllamaFailureCooldown,
RequestTimeout: cfg.OllamaRequestTimeout, FailoverEnabled: cfg.OllamaFailoverEnabled,
FailoverAttempts: cfg.OllamaFailoverAttempts, RequireSameModelDigest: cfg.OllamaRequireSameDigest,
RequireEmbeddingModel: cfg.OllamaRequireEmbeddingModel,
}, cfg.ChatModel, cfg.EmbeddingModel)
persistence := persist.New(g, b, cfg.PersistInterval)
e := &Engine{Cfg: cfg, Graph: g, Broker: b, Ollama: pool, Persistence: persistence, Scanner: &ingest.KnowledgeScanner{Graph: g, ProductionDirs: cfg.KnowledgeDirs, StagingDirs: cfg.StagingDirs}, enrichRequests: make(chan string, 1)}
if cfg.SearXNGURL != "" {
e.Research = research.New(cfg.SearXNGURL)
}
if cfg.GLPIKBEnabled {
client := glpi.New(cfg.GLPIURL, cfg.GLPIAPIVersion, cfg.GLPIClientID, cfg.GLPIClientSecret, cfg.GLPIUsername, cfg.GLPIPassword, cfg.GLPITimeout)
e.GLPIKB = ingest.NewGLPIKBSyncer(ingest.GLPIKBConfig{Enabled: true, Path: cfg.GLPIKBPath, Filter: cfg.GLPIKBFilter, Limit: cfg.GLPIKBLimit, SyncInterval: cfg.GLPIKBSyncInterval, Source: cfg.GLPIKBSource, CachePath: filepath.Join(cfg.DataDir, "glpi-kb-cache.json")}, client, g, b, persistence)
}
return e
}
func (e *Engine) Start(ctx context.Context) {
e.Persistence.Start(ctx)
e.Ollama.Start(ctx)
if e.GLPIKB != nil {
e.GLPIKB.Start(ctx)
}
go func() {
if err := e.Scan(ctx); err != nil {
slog.Error("initial brain scan failed", "error", err)
@@ -271,11 +308,15 @@ func (e *Engine) idle(ctx context.Context) {
func (e *Engine) Scan(ctx context.Context) error {
e.mu.Lock()
defer e.mu.Unlock()
e.Broker.Publish(model.Activity{Type: "scan.started", Source: "brain", Phase: "ingest", Message: "Wissensräume werden synchronisiert", Strength: .45})
beforeVersion := e.Graph.Version()
count, err := e.Scanner.Scan()
if err != nil {
return err
}
pendingEmbeddings := len(e.Graph.NodesForEmbedding())
if e.Graph.Version() != beforeVersion || pendingEmbeddings > 0 {
e.Broker.Publish(model.Activity{Type: "scan.started", Source: "brain", Phase: "ingest", Message: "Neue oder geänderte Wissenselemente werden verarbeitet", Strength: .45, Metadata: map[string]any{"pending_embeddings": pendingEmbeddings}})
}
pingCtx, pingCancel := context.WithTimeout(ctx, 3*time.Second)
pingErr := e.Ollama.Ping(pingCtx)
pingCancel()
@@ -295,14 +336,13 @@ func (e *Engine) Scan(ctx context.Context) error {
e.setOllamaOK(true)
}
}
if err := e.Graph.Persist(); err != nil {
return err
}
e.stateMu.Lock()
e.lastScan = time.Now().UTC()
e.stateMu.Unlock()
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}})
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}})
}
return nil
}
func (e *Engine) ensureEmbeddings(ctx context.Context) error {
@@ -530,14 +570,11 @@ func (e *Engine) enrichOne(ctx context.Context, trigger string) (EnrichOutcome,
return outcome, err
}
outcome.Created = true
e.Broker.Publish(model.Activity{Type: "think.created", Source: "brain", Phase: "staging", NodeIDs: []string{a.ID, b.ID}, EdgeIDs: []string{edge.ID}, Message: "Neuer AI-THINK-Beitrag wurde im Staging erzeugt", Strength: 1, Metadata: map[string]any{"trigger": trigger, "path": path, "relation_type": safeRelation(decision.RelationType), "confidence": decision.Confidence, "semantic_similarity": sim, "research_result_count": len(researchResults), "title": decision.Title}})
e.Broker.Publish(model.Activity{Type: "think.created", Source: "brain", Phase: "staging", NodeIDs: []string{a.ID, b.ID}, EdgeIDs: []string{edge.ID}, Message: "Neuer AI-THINK-Beitrag wurde erzeugt und für den gebündelten Staging-Schreibvorgang vorgemerkt", Strength: 1, Metadata: map[string]any{"trigger": trigger, "path": path, "relation_type": safeRelation(decision.RelationType), "confidence": decision.Confidence, "semantic_similarity": sim, "research_result_count": len(researchResults), "title": decision.Title, "write_pending": true}})
} else {
outcome.Rejected = true
e.Broker.Publish(model.Activity{Type: "think.rejected", Source: "brain", Phase: "validation", NodeIDs: []string{a.ID, b.ID}, Message: "Ähnlichkeit geprüft, aber nicht als belastbare Edge übernommen", Strength: .42, Metadata: map[string]any{"trigger": trigger, "relation_type": safeRelation(decision.RelationType), "confidence": decision.Confidence, "semantic_similarity": sim, "explanation": decision.Explanation}})
}
if err := e.Graph.Persist(); err != nil {
return outcome, err
}
return outcome, nil
}
@@ -555,9 +592,6 @@ func (e *Engine) writeAIThink(a, b model.Node, d model.RelationDecision, results
return "", fmt.Errorf("no BRAIN_STAGING_DIRS configured")
}
dir := e.Cfg.StagingDirs[0]
if err := os.MkdirAll(dir, 0o750); err != nil {
return "", err
}
pair := a.ID + "\x00" + b.ID
sum := sha256.Sum256([]byte(pair))
short := strings.ToUpper(hex.EncodeToString(sum[:5]))
@@ -576,17 +610,13 @@ func (e *Engine) writeAIThink(a, b model.Node, d model.RelationDecision, results
return "", err
}
path := filepath.Join(dir, strings.ToLower(id)+".json")
if e.Persistence.Pending(path) {
return path, nil
}
if _, err := os.Stat(path); err == nil {
return path, nil
}
tmp := path + ".tmp"
if err := os.WriteFile(tmp, append(bts, '\n'), 0o640); err != nil {
return "", err
}
if err := os.Rename(tmp, path); err != nil {
return "", err
}
return path, nil
return e.Persistence.QueueFile(path, append(bts, '\n'), 0o640)
}
func (e *Engine) Status() map[string]any {
s := e.Graph.Snapshot()
@@ -600,11 +630,28 @@ func (e *Engine) Status() map[string]any {
"enrich_rejected": e.enrichRejected, "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(),
}
if e.GLPIKB != nil {
status["glpi_kb"] = e.GLPIKB.Status()
} else {
status["glpi_kb"] = ingest.GLPIKBStatus{Enabled: false}
}
e.stateMu.RUnlock()
return status
}
func (e *Engine) Flush(ctx context.Context) error {
return e.Persistence.Flush(ctx, "manual")
}
func (e *Engine) SyncGLPIKB(ctx context.Context) error {
if e.GLPIKB == nil {
return fmt.Errorf("GLPI knowledge-base integration is disabled")
}
return e.GLPIKB.Sync(ctx, "manual")
}
func relationContextWithResearch(a, b model.Node, sim float64, results []model.ResearchResult) string {
var out strings.Builder
out.WriteString(relationContext(a, b, sim))
+4 -1
View File
@@ -20,7 +20,7 @@ func TestEnrichWritesAIThinkOnlyAfterStructuredQwenDecision(t *testing.T) {
w.Header().Set("Content-Type", "application/json")
switch r.URL.Path {
case "/api/tags":
_, _ = w.Write([]byte(`{"models":[]}`))
_, _ = w.Write([]byte(`{"models":[{"name":"qwen3:8b","digest":"chat-digest"},{"name":"embeddinggemma:latest","digest":"embed-digest"}]}`))
case "/api/embed":
var req struct {
Input []string `json:"input"`
@@ -72,6 +72,9 @@ func TestEnrichWritesAIThinkOnlyAfterStructuredQwenDecision(t *testing.T) {
if err := e.EnrichOne(context.Background()); err != nil {
t.Fatal(err)
}
if err := e.Flush(context.Background()); err != nil {
t.Fatal(err)
}
files, err := filepath.Glob(filepath.Join(staging, "*.json"))
if err != nil || len(files) != 1 {
t.Fatalf("expected one AI-THINK file, files=%v err=%v", files, err)
+386
View File
@@ -0,0 +1,386 @@
package glpi
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"sort"
"strconv"
"strings"
"sync"
"time"
)
type KnowledgeItem struct {
ID int64
Title string
Content string
CategoryIDs []int64
Language string
ModifiedAt string
}
type ITILCategory struct {
ID int64
Name string
CompleteName string
KnowbaseCategoryID int64
}
type Client struct {
baseURL, version, clientID, clientSecret, username, password string
http *http.Client
mu sync.Mutex
token string
tokenExpiry time.Time
}
func New(baseURL, version, clientID, clientSecret, username, password string, timeout time.Duration) *Client {
return &Client{baseURL: strings.TrimRight(baseURL, "/"), version: strings.Trim(version, "/"), clientID: clientID, clientSecret: clientSecret, username: username, password: password, http: &http.Client{Timeout: timeout}}
}
func (c *Client) APIBase() string { return c.baseURL + "/api.php/" + c.version }
func (c *Client) authenticate(ctx context.Context, force bool) (string, error) {
c.mu.Lock()
defer c.mu.Unlock()
if !force && c.token != "" && time.Until(c.tokenExpiry) > 60*time.Second {
return c.token, nil
}
form := url.Values{"grant_type": {"password"}, "client_id": {c.clientID}, "client_secret": {c.clientSecret}, "username": {c.username}, "password": {c.password}, "scope": {"api"}}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/api.php/token", strings.NewReader(form.Encode()))
if err != nil {
return "", err
}
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
req.Header.Set("Accept", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return "", err
}
defer resp.Body.Close()
body, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
if resp.StatusCode/100 != 2 {
return "", fmt.Errorf("GLPI OAuth failed: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body)))
}
var token struct {
AccessToken string `json:"access_token"`
ExpiresIn int `json:"expires_in"`
}
if err := json.Unmarshal(body, &token); err != nil {
return "", err
}
if token.AccessToken == "" {
return "", errors.New("GLPI OAuth response contains no access_token")
}
if token.ExpiresIn <= 0 {
token.ExpiresIn = 3600
}
c.token = token.AccessToken
c.tokenExpiry = time.Now().Add(time.Duration(token.ExpiresIn) * time.Second)
return c.token, nil
}
func (c *Client) do(ctx context.Context, method, path string, query url.Values, body any) ([]byte, error) {
var payload []byte
var err error
if body != nil {
payload, err = json.Marshal(body)
if err != nil {
return nil, err
}
}
for attempt := 0; attempt < 2; attempt++ {
token, err := c.authenticate(ctx, attempt > 0)
if err != nil {
return nil, err
}
u := c.APIBase() + path
if len(query) > 0 {
u += "?" + query.Encode()
}
req, err := http.NewRequestWithContext(ctx, method, u, bytes.NewReader(payload))
if err != nil {
return nil, err
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Accept", "application/json")
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
resp, err := c.http.Do(req)
if err != nil {
return nil, err
}
data, _ := io.ReadAll(io.LimitReader(resp.Body, 8<<20))
resp.Body.Close()
if resp.StatusCode == http.StatusUnauthorized && attempt == 0 {
continue
}
if resp.StatusCode/100 != 2 {
return nil, fmt.Errorf("GLPI %s %s failed: HTTP %d: %s", method, path, resp.StatusCode, strings.TrimSpace(string(data)))
}
return data, nil
}
return nil, errors.New("GLPI request failed after token refresh")
}
func (c *Client) FetchOpenAPI(ctx context.Context) (map[string]any, error) {
token, err := c.authenticate(ctx, false)
if err != nil {
return nil, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+"/api.php/doc.json", nil)
if err != nil {
return nil, err
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Accept", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return nil, fmt.Errorf("GLPI OpenAPI HTTP %d", resp.StatusCode)
}
var doc map[string]any
if err := json.NewDecoder(io.LimitReader(resp.Body, 16<<20)).Decode(&doc); err != nil {
return nil, err
}
return doc, nil
}
func (c *Client) DiscoverKnowledgeBasePath(ctx context.Context, configured string) (string, error) {
configured = strings.TrimSpace(configured)
if configured != "" && !strings.EqualFold(configured, "auto") {
if !strings.HasPrefix(configured, "/") {
return "", errors.New("GLPI KB path must be absolute")
}
return configured, nil
}
doc, err := c.FetchOpenAPI(ctx)
if err != nil {
return "", err
}
paths, ok := doc["paths"].(map[string]any)
if !ok {
return "", errors.New("GLPI OpenAPI document has no paths map")
}
type candidate struct {
path string
score int
}
var candidates []candidate
for path, raw := range paths {
if strings.Contains(path, "{") {
continue
}
ops, ok := raw.(map[string]any)
if !ok || ops["get"] == nil {
continue
}
lower := strings.ToLower(path)
score := 0
if strings.Contains(lower, "knowbaseitem") {
score += 100
}
if strings.Contains(lower, "knowledge") {
score += 40
}
if strings.Contains(lower, "knowbase") {
score += 40
}
if strings.HasSuffix(lower, "/knowbaseitem") {
score += 30
}
if score > 0 {
candidates = append(candidates, candidate{path: path, score: score})
}
}
if len(candidates) == 0 {
return "", errors.New("GLPI OpenAPI exposes no readable KnowbaseItem collection route")
}
sort.Slice(candidates, func(i, j int) bool {
if candidates[i].score == candidates[j].score {
return candidates[i].path < candidates[j].path
}
return candidates[i].score > candidates[j].score
})
path := candidates[0].path
if i := strings.Index(path, "/api.php/"); i >= 0 {
rest := path[i+len("/api.php/"):]
if slash := strings.Index(rest, "/"); slash >= 0 {
path = rest[slash:]
}
}
versionPrefix := "/" + c.version
if strings.HasPrefix(path, versionPrefix+"/") {
path = strings.TrimPrefix(path, versionPrefix)
}
if !strings.HasPrefix(path, "/") {
path = "/" + path
}
return path, nil
}
func (c *Client) ListKnowledgeBaseItems(ctx context.Context, path string, limit int, filter string) ([]KnowledgeItem, error) {
query := url.Values{"limit": {strconv.Itoa(limit)}, "sort": {"date_mod"}, "order": {"DESC"}}
if strings.TrimSpace(filter) != "" {
query.Set("filter", filter)
}
data, err := c.do(ctx, http.MethodGet, path, query, nil)
if err != nil {
return nil, err
}
items, err := extractArray(data)
if err != nil {
return nil, err
}
out := make([]KnowledgeItem, 0, len(items))
for _, raw := range items {
id := int64Val(raw["id"])
if id <= 0 {
continue
}
if firstString(raw, "answer", "content", "text", "description") == "" {
if detail, detailErr := c.do(ctx, http.MethodGet, strings.TrimRight(path, "/")+"/"+strconv.FormatInt(id, 10), nil, nil); detailErr == nil {
var full map[string]any
if json.Unmarshal(detail, &full) == nil {
for k, v := range full {
raw[k] = v
}
}
}
}
title := firstString(raw, "name", "title", "subject")
content := firstString(raw, "answer", "content", "text", "description")
if title == "" || content == "" {
continue
}
out = append(out, KnowledgeItem{ID: id, Title: title, Content: content, CategoryIDs: knowledgeCategoryIDs(raw), Language: firstString(raw, "language", "locale"), ModifiedAt: firstString(raw, "date_mod", "modified_at", "date_creation")})
}
return out, nil
}
func (c *Client) GetCategories(ctx context.Context) ([]ITILCategory, error) {
data, err := c.do(ctx, http.MethodGet, "/Dropdowns/ITILCategory", url.Values{"limit": {"1000"}}, nil)
if err != nil {
return nil, err
}
items, err := extractArray(data)
if err != nil {
return nil, err
}
out := make([]ITILCategory, 0, len(items))
for _, raw := range items {
out = append(out, ITILCategory{ID: int64Val(raw["id"]), Name: firstString(raw, "name"), CompleteName: firstString(raw, "completename"), KnowbaseCategoryID: firstRefID(raw, "knowbase_category", "knowbasecategory")})
}
return out, nil
}
func extractArray(data []byte) ([]map[string]any, error) {
var arr []map[string]any
if json.Unmarshal(data, &arr) == nil {
return arr, nil
}
var object map[string]json.RawMessage
if err := json.Unmarshal(data, &object); err != nil {
return nil, err
}
for _, key := range []string{"data", "items", "results"} {
if raw, ok := object[key]; ok && json.Unmarshal(raw, &arr) == nil {
return arr, nil
}
}
return nil, fmt.Errorf("unexpected GLPI collection response: %.200s", string(data))
}
func knowledgeCategoryIDs(raw map[string]any) []int64 {
seen := map[int64]struct{}{}
var out []int64
var add func(any)
add = func(v any) {
switch value := v.(type) {
case []any:
for _, item := range value {
add(item)
}
case map[string]any:
if id := int64Val(value["id"]); id > 0 {
if _, ok := seen[id]; !ok {
seen[id] = struct{}{}
out = append(out, id)
}
return
}
for _, nested := range value {
add(nested)
}
default:
if id := int64Val(value); id > 0 {
if _, ok := seen[id]; !ok {
seen[id] = struct{}{}
out = append(out, id)
}
}
}
}
for _, key := range []string{"knowbase_category", "knowbase_categories", "knowbaseitemcategory", "knowbaseitemcategories", "categories", "category"} {
if value, ok := raw[key]; ok {
add(value)
}
}
sort.Slice(out, func(i, j int) bool { return out[i] < out[j] })
return out
}
func firstRefID(raw map[string]any, keys ...string) int64 {
for _, key := range keys {
if value, ok := raw[key]; ok {
if id := refID(value); id > 0 {
return id
}
}
}
return 0
}
func refID(v any) int64 {
if m, ok := v.(map[string]any); ok {
return int64Val(m["id"])
}
return int64Val(v)
}
func firstString(raw map[string]any, keys ...string) string {
for _, key := range keys {
if value, ok := raw[key]; ok {
if text := strings.TrimSpace(fmt.Sprint(value)); text != "" && text != "<nil>" {
return text
}
}
}
return ""
}
func int64Val(v any) int64 {
switch x := v.(type) {
case float64:
return int64(x)
case int:
return int64(x)
case int64:
return x
case json.Number:
n, _ := x.Int64()
return n
case string:
n, _ := strconv.ParseInt(x, 10, 64)
return n
default:
return 0
}
}
+52
View File
@@ -0,0 +1,52 @@
package glpi
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"time"
)
func TestClientDiscoversAndReadsKnowledgeBase(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch r.URL.Path {
case "/api.php/token":
_ = json.NewEncoder(w).Encode(map[string]any{"access_token": "token", "expires_in": 3600})
case "/api.php/doc.json":
_ = json.NewEncoder(w).Encode(map[string]any{"paths": map[string]any{"/api.php/v2.3/Knowledge/KnowbaseItem": map[string]any{"get": map[string]any{}}}})
case "/api.php/v2.3/Knowledge/KnowbaseItem":
_ = json.NewEncoder(w).Encode(map[string]any{"data": []map[string]any{{"id": 5, "name": "VPN", "answer": "<p>Gateway prüfen</p>", "knowbase_category": map[string]any{"id": 9}, "date_mod": "now"}}})
case "/api.php/v2.3/Dropdowns/ITILCategory":
_ = json.NewEncoder(w).Encode([]map[string]any{{"id": 7, "name": "Netzwerk", "completename": "Infra > Netzwerk", "knowbase_category": map[string]any{"id": 9}}})
default:
http.NotFound(w, r)
}
}))
defer server.Close()
client := New(server.URL, "v2.3", "id", "secret", "user", "password", time.Second)
path, err := client.DiscoverKnowledgeBasePath(context.Background(), "auto")
if err != nil {
t.Fatal(err)
}
if path != "/Knowledge/KnowbaseItem" {
t.Fatalf("unexpected path %q", path)
}
items, err := client.ListKnowledgeBaseItems(context.Background(), path, 100, "")
if err != nil {
t.Fatal(err)
}
if len(items) != 1 || items[0].ID != 5 || len(items[0].CategoryIDs) != 1 || items[0].CategoryIDs[0] != 9 {
t.Fatalf("unexpected items: %+v", items)
}
categories, err := client.GetCategories(context.Background())
if err != nil {
t.Fatal(err)
}
if len(categories) != 1 || categories[0].KnowbaseCategoryID != 9 {
t.Fatalf("unexpected categories: %+v", categories)
}
}
+35 -4
View File
@@ -130,7 +130,12 @@ func (s *Store) LookupExternal(id string) (model.Node, bool) {
func (s *Store) SetVector(id string, v []float64) {
s.mu.Lock()
defer s.mu.Unlock()
old, ok := s.vectors[id]
if ok && floatSlicesEqual(old, v) {
return
}
s.vectors[id] = append([]float64(nil), v...)
s.version++
}
func (s *Store) Vector(id string) ([]float64, bool) {
s.mu.RLock()
@@ -248,6 +253,12 @@ func (s *Store) ReplaceOrigins(origins []string, nodes []model.Node, edges []mod
}
s.version++
}
func (s *Store) Version() uint64 {
s.mu.RLock()
defer s.mu.RUnlock()
return s.version
}
func (s *Store) Snapshot() model.Snapshot {
s.mu.RLock()
defer s.mu.RUnlock()
@@ -267,6 +278,11 @@ func (s *Store) Snapshot() model.Snapshot {
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 {
@@ -276,18 +292,21 @@ func (s *Store) Persist() error {
d.Edges = append(d.Edges, e)
}
for k, v := range s.vectors {
d.Vectors[k] = v
d.Vectors[k] = append([]float64(nil), v...)
}
s.mu.RUnlock()
b, err := json.Marshal(d)
if err != nil {
return err
return 0, err
}
tmp := s.path + ".tmp"
if err = os.WriteFile(tmp, b, 0o640); err != nil {
return err
return 0, err
}
return os.Rename(tmp, s.path)
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()
@@ -428,6 +447,18 @@ func edgeBetweenLocked(edges map[string]model.Edge, a, b string) bool {
}
return false
}
func floatSlicesEqual(a, b []float64) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
func cosine(a, b []float64) float64 {
if len(a) == 0 || len(a) != len(b) {
return 0
+330
View File
@@ -0,0 +1,330 @@
package ingest
import (
"context"
"encoding/json"
"fmt"
htmlstd "html"
"os"
"path/filepath"
"regexp"
"sort"
"strconv"
"strings"
"sync"
"time"
"github.com/local/glpi-neural-brain/internal/activity"
"github.com/local/glpi-neural-brain/internal/glpi"
"github.com/local/glpi-neural-brain/internal/graph"
"github.com/local/glpi-neural-brain/internal/model"
"github.com/local/glpi-neural-brain/internal/persist"
)
type GLPIKBConfig struct {
Enabled bool
Path string
Filter string
Limit int
SyncInterval time.Duration
Source string
CachePath string
}
type GLPIKBStatus struct {
Enabled bool `json:"enabled"`
Path string `json:"path,omitempty"`
LastSync time.Time `json:"last_sync,omitempty"`
LastError string `json:"last_error,omitempty"`
Documents int `json:"documents"`
CachePath string `json:"cache_path,omitempty"`
NextSync time.Time `json:"next_sync,omitempty"`
}
type GLPIKBSource interface {
DiscoverKnowledgeBasePath(context.Context, string) (string, error)
ListKnowledgeBaseItems(context.Context, string, int, string) ([]glpi.KnowledgeItem, error)
GetCategories(context.Context) ([]glpi.ITILCategory, error)
}
type GLPIKBSyncer struct {
cfg GLPIKBConfig
client GLPIKBSource
graph *graph.Store
broker *activity.Broker
persistence *persist.Coordinator
syncMu sync.Mutex
mu sync.RWMutex
path string
lastSync time.Time
lastError string
documents int
nextSync time.Time
lastFingerprint string
}
type glpiKBCache struct {
SyncedAt time.Time `json:"synced_at"`
Path string `json:"path"`
Nodes []model.Node `json:"nodes"`
Edges []model.Edge `json:"edges"`
}
func NewGLPIKBSyncer(cfg GLPIKBConfig, client GLPIKBSource, g *graph.Store, b *activity.Broker, p *persist.Coordinator) *GLPIKBSyncer {
return &GLPIKBSyncer{cfg: cfg, client: client, graph: g, broker: b, persistence: p}
}
func (s *GLPIKBSyncer) Start(ctx context.Context) {
if !s.cfg.Enabled {
return
}
_ = s.LoadCache()
go func() {
s.runSync(ctx, "startup")
ticker := time.NewTicker(s.cfg.SyncInterval)
defer ticker.Stop()
for {
s.setNextSync(time.Now().Add(s.cfg.SyncInterval))
select {
case <-ctx.Done():
return
case <-ticker.C:
s.runSync(ctx, "interval")
}
}
}()
}
func (s *GLPIKBSyncer) runSync(ctx context.Context, trigger string) {
timeout := s.cfg.SyncInterval / 2
if timeout < 30*time.Second {
timeout = 30 * time.Second
}
if timeout > 5*time.Minute {
timeout = 5 * time.Minute
}
syncCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
if err := s.Sync(syncCtx, trigger); err != nil {
s.mu.Lock()
s.lastError = err.Error()
s.mu.Unlock()
if s.broker != nil {
s.broker.Publish(model.Activity{Type: "glpi.kb.failed", Source: "glpi", Phase: "knowledge-sync", Message: "GLPI-Wissensdatenbank konnte nicht synchronisiert werden", Strength: .35, Metadata: map[string]any{"error": err.Error(), "trigger": trigger}})
}
}
}
func (s *GLPIKBSyncer) LoadCache() error {
if strings.TrimSpace(s.cfg.CachePath) == "" {
return nil
}
data, err := os.ReadFile(s.cfg.CachePath)
if os.IsNotExist(err) {
return nil
}
if err != nil {
return err
}
var cache glpiKBCache
if err := json.Unmarshal(data, &cache); err != nil {
return err
}
s.graph.ReplaceOrigins([]string{"glpi-kb", "glpi-kb-taxonomy"}, cache.Nodes, cache.Edges)
fingerprint, _ := knowledgeFingerprint(cache.Nodes, cache.Edges)
s.mu.Lock()
s.path = cache.Path
s.lastSync = cache.SyncedAt
s.documents = countKnowledge(cache.Nodes)
s.lastError = ""
s.lastFingerprint = fingerprint
s.mu.Unlock()
return nil
}
func (s *GLPIKBSyncer) Sync(ctx context.Context, trigger string) error {
s.syncMu.Lock()
defer s.syncMu.Unlock()
path := s.currentPath()
if path == "" {
var err error
path, err = s.client.DiscoverKnowledgeBasePath(ctx, s.cfg.Path)
if err != nil {
return fmt.Errorf("discover GLPI KB path: %w", err)
}
}
items, err := s.client.ListKnowledgeBaseItems(ctx, path, s.cfg.Limit, s.cfg.Filter)
if err != nil {
return fmt.Errorf("list GLPI KB items: %w", err)
}
categories, err := s.client.GetCategories(ctx)
if err != nil {
return fmt.Errorf("load GLPI categories: %w", err)
}
nodes, edges := s.buildGraph(items, categories)
fingerprint, err := knowledgeFingerprint(nodes, edges)
if err != nil {
return err
}
now := time.Now().UTC()
s.mu.RLock()
unchanged := fingerprint != "" && fingerprint == s.lastFingerprint
s.mu.RUnlock()
if unchanged {
s.mu.Lock()
s.path = path
s.lastSync = now
s.lastError = ""
s.documents = len(items)
s.mu.Unlock()
if s.broker != nil {
s.broker.Publish(model.Activity{Type: "glpi.kb.unchanged", Source: "glpi", Phase: "knowledge-sync", Message: "GLPI-Wissensdatenbank ist unverändert", Strength: .18, Metadata: map[string]any{"documents": len(items), "path": path, "trigger": trigger}})
}
return nil
}
s.graph.ReplaceOrigins([]string{"glpi-kb", "glpi-kb-taxonomy"}, nodes, edges)
cache := glpiKBCache{SyncedAt: now, Path: path, Nodes: nodes, Edges: edges}
if s.persistence != nil && s.cfg.CachePath != "" {
data, err := json.Marshal(cache)
if err != nil {
return err
}
if _, err := s.persistence.QueueFile(s.cfg.CachePath, data, 0o640); err != nil {
return err
}
}
s.mu.Lock()
s.path = path
s.lastSync = now
s.lastError = ""
s.documents = len(items)
s.lastFingerprint = fingerprint
s.mu.Unlock()
if s.broker != nil {
s.broker.Publish(model.Activity{Type: "glpi.kb.synced", Source: "glpi", Phase: "knowledge-sync", Message: fmt.Sprintf("%d GLPI-KB-Beiträge wurden in den Wissensgraphen übernommen", len(items)), Strength: .64, Metadata: map[string]any{"documents": len(items), "nodes": len(nodes), "edges": len(edges), "path": path, "trigger": trigger}})
}
return nil
}
func (s *GLPIKBSyncer) Status() GLPIKBStatus {
s.mu.RLock()
defer s.mu.RUnlock()
return GLPIKBStatus{Enabled: s.cfg.Enabled, Path: s.path, LastSync: s.lastSync, LastError: s.lastError, Documents: s.documents, CachePath: s.cfg.CachePath, NextSync: s.nextSync}
}
func (s *GLPIKBSyncer) setNextSync(t time.Time) { s.mu.Lock(); s.nextSync = t; s.mu.Unlock() }
func (s *GLPIKBSyncer) currentPath() string { s.mu.RLock(); defer s.mu.RUnlock(); return s.path }
func (s *GLPIKBSyncer) buildGraph(items []glpi.KnowledgeItem, categories []glpi.ITILCategory) ([]model.Node, []model.Edge) {
kbCategoryNames := map[int64][]string{}
for _, category := range categories {
if category.KnowbaseCategoryID <= 0 {
continue
}
name := strings.TrimSpace(category.CompleteName)
if name == "" {
name = strings.TrimSpace(category.Name)
}
if name != "" {
kbCategoryNames[category.KnowbaseCategoryID] = append(kbCategoryNames[category.KnowbaseCategoryID], name)
}
}
now := time.Now().UTC()
nodesByID := map[string]model.Node{}
var edges []model.Edge
sourceLabel := s.cfg.Source
if sourceLabel == "" {
sourceLabel = "GLPI Knowledge Base"
}
sourceID := graph.ID("source", strings.ToLower(sourceLabel))
nodesByID[sourceID] = model.Node{ID: sourceID, Kind: "source", Label: sourceLabel, Origin: "glpi-kb-taxonomy", ExternalID: sourceLabel, Weight: .85, UpdatedAt: now}
for _, item := range items {
text := cleanGLPIHTML(item.Content)
if text == "" {
continue
}
externalID := "GLPI-KB-" + strconv.FormatInt(item.ID, 10)
nodeID := graph.ID("knowledge", externalID)
categoriesForNode := []string{"GLPI KB"}
for _, categoryID := range item.CategoryIDs {
categoriesForNode = append(categoriesForNode, kbCategoryNames[categoryID]...)
}
categoriesForNode = uniqueStrings(categoriesForNode)
keywords := extractKeywords(item.Title, categoriesForNode)
nodesByID[nodeID] = model.Node{ID: nodeID, Kind: "knowledge", Label: strings.TrimSpace(item.Title), Summary: clamp(text, 1100), Status: "production", Origin: "glpi-kb", ExternalID: externalID, URI: "glpi://KnowbaseItem/" + strconv.FormatInt(item.ID, 10), Categories: categoriesForNode, Keywords: keywords, Weight: 1.35, Metadata: map[string]any{"source": sourceLabel, "glpi_id": item.ID, "glpi_category_ids": item.CategoryIDs, "language": item.Language, "modified_at": item.ModifiedAt}, UpdatedAt: now}
edges = append(edges, model.Edge{Source: nodeID, Target: sourceID, Type: "derived_from", Origin: "glpi-kb", Status: "verified", Confidence: 1, Weight: .45})
for _, category := range categoriesForNode {
categoryID := graph.ID("category", strings.ToLower(category))
if _, ok := nodesByID[categoryID]; !ok {
nodesByID[categoryID] = model.Node{ID: categoryID, Kind: "category", Label: category, Origin: "glpi-kb-taxonomy", ExternalID: category, Weight: .8, UpdatedAt: now}
}
edges = append(edges, model.Edge{Source: nodeID, Target: categoryID, Type: "categorized_as", Origin: "glpi-kb", Status: "verified", Confidence: 1, Weight: .58})
}
}
nodes := make([]model.Node, 0, len(nodesByID))
for _, node := range nodesByID {
nodes = append(nodes, node)
}
sort.Slice(nodes, func(i, j int) bool { return nodes[i].ID < nodes[j].ID })
return nodes, edges
}
var glpiTagRE = regexp.MustCompile(`(?s)<[^>]*>`)
var glpiSpaceRE = regexp.MustCompile(`[\t\r\n ]+`)
var keywordSplitRE = regexp.MustCompile(`[^\pL\pN_-]+`)
func cleanGLPIHTML(value string) string {
replacer := strings.NewReplacer("<br>", "\n", "<br/>", "\n", "<br />", "\n", "</p>", "\n", "</li>", "\n")
value = replacer.Replace(value)
value = glpiTagRE.ReplaceAllString(value, " ")
value = htmlstd.UnescapeString(value)
return strings.TrimSpace(glpiSpaceRE.ReplaceAllString(value, " "))
}
func extractKeywords(title string, categories []string) []string {
var out []string
for _, part := range keywordSplitRE.Split(title, -1) {
part = strings.TrimSpace(part)
if len([]rune(part)) >= 4 {
out = append(out, part)
}
}
out = append(out, categories...)
return uniqueStrings(out)
}
func uniqueStrings(in []string) []string {
seen := map[string]struct{}{}
var out []string
for _, value := range in {
value = strings.TrimSpace(value)
key := strings.ToLower(value)
if value == "" {
continue
}
if _, ok := seen[key]; ok {
continue
}
seen[key] = struct{}{}
out = append(out, value)
}
return out
}
func countKnowledge(nodes []model.Node) int {
count := 0
for _, node := range nodes {
if node.Kind == "knowledge" {
count++
}
}
return count
}
func ensureCacheDir(path string) error {
if path == "" {
return nil
}
return os.MkdirAll(filepath.Dir(path), 0o750)
}
+72
View File
@@ -0,0 +1,72 @@
package ingest
import (
"context"
"os"
"path/filepath"
"testing"
"time"
"github.com/local/glpi-neural-brain/internal/glpi"
"github.com/local/glpi-neural-brain/internal/graph"
"github.com/local/glpi-neural-brain/internal/persist"
)
type fakeGLPIKB struct{}
func (fakeGLPIKB) DiscoverKnowledgeBasePath(context.Context, string) (string, error) {
return "/Knowledge/KnowbaseItem", nil
}
func (fakeGLPIKB) ListKnowledgeBaseItems(context.Context, string, int, string) ([]glpi.KnowledgeItem, error) {
return []glpi.KnowledgeItem{{ID: 42, Title: "VPN Zugriff", Content: "<p>Gateway prüfen</p>", CategoryIDs: []int64{9}, Language: "de-DE", ModifiedAt: "now"}}, nil
}
func (fakeGLPIKB) GetCategories(context.Context) ([]glpi.ITILCategory, error) {
return []glpi.ITILCategory{{ID: 7, Name: "Netzwerk", CompleteName: "Infrastruktur > Netzwerk", KnowbaseCategoryID: 9}}, nil
}
func TestGLPIKBSyncAddsGraphAndQueuesCache(t *testing.T) {
dir := t.TempDir()
g, err := graph.Open(dir)
if err != nil {
t.Fatal(err)
}
p := persist.New(g, nil, time.Minute)
cachePath := filepath.Join(dir, "cache", "glpi-kb-cache.json")
s := NewGLPIKBSyncer(GLPIKBConfig{Enabled: true, Path: "auto", Limit: 100, SyncInterval: time.Minute, Source: "glpi-kb", CachePath: cachePath}, fakeGLPIKB{}, g, nil, p)
if err := s.Sync(context.Background(), "test"); err != nil {
t.Fatal(err)
}
status := s.Status()
if status.Documents != 1 || status.Path != "/Knowledge/KnowbaseItem" {
t.Fatalf("unexpected status: %+v", status)
}
snap := g.Snapshot()
found := false
for _, n := range snap.Nodes {
if n.ExternalID == "GLPI-KB-42" {
found = true
if n.Status != "production" || n.Origin != "glpi-kb" {
t.Fatalf("unexpected node: %+v", n)
}
}
}
if !found {
t.Fatal("GLPI KB node missing")
}
if _, err := os.Stat(cachePath); !os.IsNotExist(err) {
t.Fatalf("cache must be buffered before flush: %v", err)
}
version := g.Version()
if err := s.Sync(context.Background(), "test-unchanged"); err != nil {
t.Fatal(err)
}
if g.Version() != version {
t.Fatalf("unchanged GLPI sync changed graph version: %d -> %d", version, g.Version())
}
if err := p.Flush(context.Background(), "test"); err != nil {
t.Fatal(err)
}
if _, err := os.Stat(cachePath); err != nil {
t.Fatal(err)
}
}
+43 -3
View File
@@ -1,6 +1,8 @@
package ingest
import (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io/fs"
@@ -15,9 +17,10 @@ import (
)
type KnowledgeScanner struct {
Graph *graph.Store
ProductionDirs []string
StagingDirs []string
Graph *graph.Store
ProductionDirs []string
StagingDirs []string
lastFingerprint string
}
func (s *KnowledgeScanner) Scan() (int, error) {
@@ -56,7 +59,15 @@ func (s *KnowledgeScanner) Scan() (int, error) {
}
edges = append(edges, e...)
}
fingerprint, err := knowledgeFingerprint(nodes, edges)
if err != nil {
return 0, err
}
if fingerprint == s.lastFingerprint {
return len(nodes), nil
}
s.Graph.ReplaceOrigins([]string{"knowledge-production", "knowledge-staging", "knowledge-taxonomy"}, nodes, edges)
s.lastFingerprint = fingerprint
return len(nodes), nil
}
@@ -168,6 +179,35 @@ func scanDir(root, origin, status string) ([]model.Node, []model.Edge, error) {
return docs, edges, nil
}
func knowledgeFingerprint(nodes []model.Node, edges []model.Edge) (string, error) {
nodeCopies := append([]model.Node(nil), nodes...)
edgeCopies := append([]model.Edge(nil), edges...)
for i := range nodeCopies {
nodeCopies[i].UpdatedAt = time.Time{}
nodeCopies[i].X, nodeCopies[i].Y, nodeCopies[i].Z = 0, 0, 0
}
for i := range edgeCopies {
edgeCopies[i].ID = ""
edgeCopies[i].CreatedAt = time.Time{}
edgeCopies[i].UpdatedAt = time.Time{}
}
sort.Slice(nodeCopies, func(i, j int) bool { return nodeCopies[i].ID < nodeCopies[j].ID })
sort.Slice(edgeCopies, func(i, j int) bool {
a := edgeCopies[i].Source + "\x00" + edgeCopies[i].Target + "\x00" + edgeCopies[i].Type + "\x00" + edgeCopies[i].Origin
b := edgeCopies[j].Source + "\x00" + edgeCopies[j].Target + "\x00" + edgeCopies[j].Type + "\x00" + edgeCopies[j].Origin
return a < b
})
data, err := json.Marshal(struct {
Nodes []model.Node `json:"nodes"`
Edges []model.Edge `json:"edges"`
}{Nodes: nodeCopies, Edges: edgeCopies})
if err != nil {
return "", err
}
sum := sha256.Sum256(data)
return hex.EncodeToString(sum[:]), nil
}
func firstString(m map[string]any, keys ...string) string {
for _, k := range keys {
if v, ok := m[k]; ok {
+26
View File
@@ -43,3 +43,29 @@ func TestKnowledgeScannerIncludesAIThinkStaging(t *testing.T) {
t.Fatalf("missing nodes: prod=%v think=%v", prodOK, thinkOK)
}
}
func TestKnowledgeScannerSkipsUnchangedGraphReplacement(t *testing.T) {
root := t.TempDir()
prod := filepath.Join(root, "prod")
if err := os.MkdirAll(prod, 0o755); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(filepath.Join(prod, "one.json"), []byte(`{"id":"KB-ONE","title":"One","text":"Same","categories":["Test"]}`), 0o644); err != nil {
t.Fatal(err)
}
g, err := graph.Open(filepath.Join(root, "data"))
if err != nil {
t.Fatal(err)
}
scanner := KnowledgeScanner{Graph: g, ProductionDirs: []string{prod}}
if _, err := scanner.Scan(); err != nil {
t.Fatal(err)
}
v1 := g.Version()
if _, err := scanner.Scan(); err != nil {
t.Fatal(err)
}
if v2 := g.Version(); v2 != v1 {
t.Fatalf("unchanged scan changed graph version: %d -> %d", v1, v2)
}
}
+467 -23
View File
@@ -4,27 +4,150 @@ import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"sort"
"strings"
"sync"
"time"
)
type Client struct {
BaseURL, ChatModel, EmbeddingModel string
HTTP *http.Client
type NodeConfig struct {
Name string
URL string
Weight int
}
func New(base, chat, embed string) *Client {
return &Client{BaseURL: strings.TrimRight(base, "/"), ChatModel: chat, EmbeddingModel: embed, HTTP: &http.Client{Timeout: 8 * time.Minute}}
type PoolConfig struct {
Nodes []NodeConfig
RoutingMode string
NodeMaxInflight int
HealthInterval time.Duration
FailureCooldown time.Duration
RequestTimeout time.Duration
FailoverEnabled bool
FailoverAttempts int
RequireSameModelDigest bool
RequireEmbeddingModel bool
}
type NodeStatus struct {
Name string `json:"name"`
URL string `json:"url"`
Weight int `json:"weight"`
Healthy bool `json:"healthy"`
Compatible bool `json:"compatible"`
ChatModel bool `json:"chat_model"`
EmbeddingModel bool `json:"embedding_model"`
ChatDigest string `json:"chat_digest,omitempty"`
EmbeddingDigest string `json:"embedding_digest,omitempty"`
Inflight int `json:"inflight"`
Requests uint64 `json:"requests"`
Failures uint64 `json:"failures"`
AverageDurationMS float64 `json:"average_duration_ms"`
CooldownUntil time.Time `json:"cooldown_until,omitempty"`
LastCheck time.Time `json:"last_check,omitempty"`
LastError string `json:"last_error,omitempty"`
}
type nodeState struct {
cfg NodeConfig
healthy bool
compatible bool
chatModel bool
embeddingModel bool
chatDigest string
embeddingDigest string
inflight int
requests uint64
failures uint64
totalDuration time.Duration
averageDuration time.Duration
cooldownUntil time.Time
lastCheck time.Time
lastError string
}
type Client struct {
ChatModel, EmbeddingModel string
HTTP *http.Client
cfg PoolConfig
mu sync.Mutex
nodes []*nodeState
roundRobin uint64
healthReady bool
healthMu sync.Mutex
}
type requestError struct {
err error
retryable bool
statusCode int
}
func (e *requestError) Error() string { return e.err.Error() }
func (e *requestError) Unwrap() error { return e.err }
func New(base, chat, embed string) *Client {
return NewPool(PoolConfig{Nodes: []NodeConfig{{Name: "ollama-1", URL: base, Weight: 1}}, RoutingMode: "least_inflight", NodeMaxInflight: 1, HealthInterval: 15 * time.Second, FailureCooldown: 30 * time.Second, RequestTimeout: 8 * time.Minute, FailoverEnabled: true, FailoverAttempts: 1, RequireSameModelDigest: true, RequireEmbeddingModel: true}, chat, embed)
}
func NewPool(cfg PoolConfig, chat, embed string) *Client {
if cfg.NodeMaxInflight < 1 {
cfg.NodeMaxInflight = 1
}
if cfg.HealthInterval <= 0 {
cfg.HealthInterval = 15 * time.Second
}
if cfg.FailureCooldown <= 0 {
cfg.FailureCooldown = 30 * time.Second
}
if cfg.RequestTimeout <= 0 {
cfg.RequestTimeout = 8 * time.Minute
}
if cfg.RoutingMode == "" {
cfg.RoutingMode = "least_inflight"
}
c := &Client{ChatModel: chat, EmbeddingModel: embed, cfg: cfg, HTTP: &http.Client{Timeout: cfg.RequestTimeout}}
for i, raw := range cfg.Nodes {
raw.URL = strings.TrimRight(strings.TrimSpace(raw.URL), "/")
if raw.Name == "" {
raw.Name = fmt.Sprintf("ollama-%d", i+1)
}
if raw.Weight < 1 {
raw.Weight = 1
}
c.nodes = append(c.nodes, &nodeState{cfg: raw})
}
return c
}
func (c *Client) Start(ctx context.Context) {
go func() {
_ = c.refreshHealth(ctx)
ticker := time.NewTicker(c.cfg.HealthInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
_ = c.refreshHealth(ctx)
}
}
}()
}
func (c *Client) Embed(ctx context.Context, texts []string) ([][]float64, error) {
body := map[string]any{"model": c.EmbeddingModel, "input": texts, "truncate": true}
var out struct {
Embeddings [][]float64 `json:"embeddings"`
}
if err := c.post(ctx, "/api/embed", body, &out); err != nil {
if err := c.doJSON(ctx, "embedding", "/api/embed", body, &out); err != nil {
return nil, err
}
if len(out.Embeddings) != len(texts) {
@@ -32,6 +155,7 @@ func (c *Client) Embed(ctx context.Context, texts []string) ([][]float64, error)
}
return out.Embeddings, nil
}
func (c *Client) ChatJSON(ctx context.Context, system, user string, schema any, target any) error {
body := map[string]any{"model": c.ChatModel, "messages": []map[string]string{{"role": "system", "content": system}, {"role": "user", "content": user}}, "stream": false, "think": false, "format": schema, "options": map[string]any{"temperature": 0}}
var env struct {
@@ -39,7 +163,7 @@ func (c *Client) ChatJSON(ctx context.Context, system, user string, schema any,
Content string `json:"content"`
} `json:"message"`
}
if err := c.post(ctx, "/api/chat", body, &env); err != nil {
if err := c.doJSON(ctx, "chat", "/api/chat", body, &env); err != nil {
return err
}
if strings.TrimSpace(env.Message.Content) == "" {
@@ -50,45 +174,365 @@ func (c *Client) ChatJSON(ctx context.Context, system, user string, schema any,
}
return nil
}
func (c *Client) Ping(ctx context.Context) error {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.BaseURL+"/api/tags", nil)
if err != nil {
if err := c.ensureHealth(ctx); err != nil {
return err
}
resp, err := c.HTTP.Do(req)
if err != nil {
return err
c.mu.Lock()
defer c.mu.Unlock()
for _, n := range c.nodes {
if n.healthy && n.compatible && n.chatModel && (!c.cfg.RequireEmbeddingModel || n.embeddingModel) {
return nil
}
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return fmt.Errorf("ollama HTTP %d", resp.StatusCode)
}
return nil
return errors.New("no healthy compatible Ollama node")
}
func (c *Client) post(ctx context.Context, path string, in, out any) error {
func (c *Client) RoutingMode() string { return c.cfg.RoutingMode }
func (c *Client) NodeStatuses() []NodeStatus {
c.mu.Lock()
defer c.mu.Unlock()
out := make([]NodeStatus, 0, len(c.nodes))
for _, n := range c.nodes {
avg := float64(n.averageDuration.Milliseconds())
out = append(out, NodeStatus{Name: n.cfg.Name, URL: n.cfg.URL, Weight: n.cfg.Weight, Healthy: n.healthy, Compatible: n.compatible, ChatModel: n.chatModel, EmbeddingModel: n.embeddingModel, ChatDigest: n.chatDigest, EmbeddingDigest: n.embeddingDigest, Inflight: n.inflight, Requests: n.requests, Failures: n.failures, AverageDurationMS: avg, CooldownUntil: n.cooldownUntil, LastCheck: n.lastCheck, LastError: n.lastError})
}
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out
}
func (c *Client) PoolStatus() map[string]any {
statuses := c.NodeStatuses()
healthy, available := 0, 0
now := time.Now()
for _, s := range statuses {
if s.Healthy && s.Compatible {
healthy++
if now.After(s.CooldownUntil) && s.Inflight < c.cfg.NodeMaxInflight {
available++
}
}
}
return map[string]any{"routing_mode": c.cfg.RoutingMode, "node_count": len(statuses), "healthy_nodes": healthy, "available_nodes": available, "node_max_inflight": c.cfg.NodeMaxInflight, "failover_enabled": c.cfg.FailoverEnabled, "nodes": statuses}
}
func (c *Client) doJSON(ctx context.Context, capability, path string, in, out any) error {
if err := c.ensureHealth(ctx); err != nil {
return err
}
attemptLimit := c.cfg.FailoverAttempts
if attemptLimit <= 0 || attemptLimit > len(c.nodes) {
attemptLimit = len(c.nodes)
}
if !c.cfg.FailoverEnabled && attemptLimit > 1 {
attemptLimit = 1
}
tried := make(map[*nodeState]bool, attemptLimit)
var errs []string
for attempt := 0; attempt < attemptLimit; attempt++ {
n, err := c.acquireNode(ctx, capability, tried)
if err != nil {
if len(errs) > 0 {
return fmt.Errorf("ollama pool failed: %s; %w", strings.Join(errs, "; "), err)
}
return err
}
tried[n] = true
started := time.Now()
err = c.postNode(ctx, n, path, in, out)
c.releaseNode(n, time.Since(started), err)
if err == nil {
return nil
}
errs = append(errs, n.cfg.Name+": "+err.Error())
var re *requestError
if !errors.As(err, &re) || !re.retryable || !c.cfg.FailoverEnabled {
break
}
}
return fmt.Errorf("ollama pool request failed: %s", strings.Join(errs, "; "))
}
func (c *Client) ensureHealth(ctx context.Context) error {
c.mu.Lock()
ready := c.healthReady
c.mu.Unlock()
if ready {
return nil
}
return c.refreshHealth(ctx)
}
func (c *Client) acquireNode(ctx context.Context, capability string, tried map[*nodeState]bool) (*nodeState, error) {
for {
c.mu.Lock()
candidates := make([]*nodeState, 0, len(c.nodes))
now := time.Now()
viable := 0
busy := false
for _, n := range c.nodes {
if tried[n] || !n.healthy || !n.compatible {
continue
}
if capability == "embedding" && !n.embeddingModel {
continue
}
if capability == "chat" && !n.chatModel {
continue
}
viable++
if now.Before(n.cooldownUntil) {
continue
}
if n.inflight >= c.cfg.NodeMaxInflight {
busy = true
continue
}
candidates = append(candidates, n)
}
if len(candidates) > 0 {
n := c.chooseNode(candidates)
n.inflight++
c.mu.Unlock()
return n, nil
}
c.mu.Unlock()
if viable == 0 || !busy {
return nil, errors.New("no additional compatible Ollama node available")
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(20 * time.Millisecond):
}
}
}
func (c *Client) chooseNode(nodes []*nodeState) *nodeState {
switch c.cfg.RoutingMode {
case "round_robin":
sort.Slice(nodes, func(i, j int) bool { return nodes[i].cfg.Name < nodes[j].cfg.Name })
n := nodes[c.roundRobin%uint64(len(nodes))]
c.roundRobin++
return n
case "weighted":
sort.Slice(nodes, func(i, j int) bool {
a := float64(nodes[i].requests+uint64(nodes[i].inflight)) / float64(nodes[i].cfg.Weight)
b := float64(nodes[j].requests+uint64(nodes[j].inflight)) / float64(nodes[j].cfg.Weight)
if a == b {
return nodes[i].cfg.Name < nodes[j].cfg.Name
}
return a < b
})
return nodes[0]
case "fastest_recent":
sort.Slice(nodes, func(i, j int) bool {
if nodes[i].requests == 0 || nodes[j].requests == 0 {
if nodes[i].requests == nodes[j].requests {
return nodes[i].cfg.Name < nodes[j].cfg.Name
}
return nodes[i].requests == 0
}
a := nodes[i].averageDuration
b := nodes[j].averageDuration
if a == b {
return nodes[i].cfg.Name < nodes[j].cfg.Name
}
return a < b
})
return nodes[0]
default: // least_inflight
sort.Slice(nodes, func(i, j int) bool {
if nodes[i].inflight != nodes[j].inflight {
return nodes[i].inflight < nodes[j].inflight
}
if nodes[i].requests != nodes[j].requests {
return nodes[i].requests < nodes[j].requests
}
return nodes[i].cfg.Name < nodes[j].cfg.Name
})
return nodes[0]
}
}
func (c *Client) releaseNode(n *nodeState, duration time.Duration, err error) {
c.mu.Lock()
defer c.mu.Unlock()
if n.inflight > 0 {
n.inflight--
}
n.requests++
n.totalDuration += duration
if n.averageDuration == 0 {
n.averageDuration = duration
} else {
n.averageDuration = time.Duration(float64(n.averageDuration)*0.8 + float64(duration)*0.2)
}
if err == nil {
n.lastError = ""
return
}
n.failures++
n.lastError = err.Error()
var re *requestError
if errors.As(err, &re) && re.retryable {
n.cooldownUntil = time.Now().Add(c.cfg.FailureCooldown)
}
}
func (c *Client) postNode(ctx context.Context, n *nodeState, path string, in, out any) error {
b, err := json.Marshal(in)
if err != nil {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.BaseURL+path, bytes.NewReader(b))
req, err := http.NewRequestWithContext(ctx, http.MethodPost, n.cfg.URL+path, bytes.NewReader(b))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.HTTP.Do(req)
if err != nil {
return err
return &requestError{err: err, retryable: true}
}
defer resp.Body.Close()
data, err := io.ReadAll(io.LimitReader(resp.Body, 16<<20))
if err != nil {
return err
return &requestError{err: err, retryable: true}
}
if resp.StatusCode/100 != 2 {
return fmt.Errorf("ollama %s HTTP %d: %s", path, resp.StatusCode, strings.TrimSpace(string(data)))
retryable := resp.StatusCode == http.StatusRequestTimeout || resp.StatusCode == http.StatusTooManyRequests || resp.StatusCode >= 500
return &requestError{err: fmt.Errorf("ollama %s HTTP %d: %s", path, resp.StatusCode, strings.TrimSpace(string(data))), retryable: retryable, statusCode: resp.StatusCode}
}
if err := json.Unmarshal(data, out); err != nil {
return err
return &requestError{err: fmt.Errorf("decode Ollama response: %w", err), retryable: true}
}
return nil
}
func (c *Client) refreshHealth(ctx context.Context) error {
c.healthMu.Lock()
defer c.healthMu.Unlock()
type result struct {
n *nodeState
models map[string]string
err error
}
results := make(chan result, len(c.nodes))
for _, n := range c.nodes {
go func(n *nodeState) {
models, err := c.fetchTags(ctx, n.cfg.URL)
results <- result{n: n, models: models, err: err}
}(n)
}
collected := make([]result, 0, len(c.nodes))
for range c.nodes {
collected = append(collected, <-results)
}
c.mu.Lock()
defer c.mu.Unlock()
chatDigests := map[string]struct{}{}
embedDigests := map[string]struct{}{}
now := time.Now().UTC()
for _, r := range collected {
n := r.n
n.lastCheck = now
n.compatible = false
if r.err != nil {
n.healthy = false
n.chatModel = false
n.embeddingModel = false
n.lastError = r.err.Error()
continue
}
n.healthy = true
n.lastError = ""
n.chatDigest, n.chatModel = findModelDigest(r.models, c.ChatModel)
n.embeddingDigest, n.embeddingModel = findModelDigest(r.models, c.EmbeddingModel)
if n.chatModel && n.chatDigest != "" {
chatDigests[n.chatDigest] = struct{}{}
}
if n.embeddingModel && n.embeddingDigest != "" {
embedDigests[n.embeddingDigest] = struct{}{}
}
}
divergentChat := c.cfg.RequireSameModelDigest && len(chatDigests) > 1
divergentEmbed := c.cfg.RequireSameModelDigest && len(embedDigests) > 1
for _, n := range c.nodes {
if !n.healthy || !n.chatModel || (c.cfg.RequireEmbeddingModel && !n.embeddingModel) {
continue
}
if c.cfg.RequireSameModelDigest && (n.chatDigest == "" || (c.cfg.RequireEmbeddingModel && n.embeddingDigest == "")) {
n.lastError = "model digest missing while strict digest validation is enabled"
continue
}
if divergentChat || divergentEmbed {
n.lastError = "model digest mismatch inside Ollama pool"
continue
}
n.compatible = true
}
c.healthReady = true
for _, n := range c.nodes {
if n.healthy && n.compatible {
return nil
}
}
return errors.New("no healthy compatible Ollama node")
}
func (c *Client) fetchTags(ctx context.Context, base string) (map[string]string, error) {
reqCtx, cancel := context.WithTimeout(ctx, minDuration(c.cfg.RequestTimeout, 20*time.Second))
defer cancel()
req, err := http.NewRequestWithContext(reqCtx, http.MethodGet, base+"/api/tags", nil)
if err != nil {
return nil, err
}
resp, err := c.HTTP.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return nil, fmt.Errorf("ollama tags HTTP %d", resp.StatusCode)
}
var env struct {
Models []struct {
Name string `json:"name"`
Model string `json:"model"`
Digest string `json:"digest"`
} `json:"models"`
}
if err := json.NewDecoder(io.LimitReader(resp.Body, 4<<20)).Decode(&env); err != nil {
return nil, err
}
out := make(map[string]string)
for _, m := range env.Models {
name := m.Name
if name == "" {
name = m.Model
}
out[name] = m.Digest
}
return out, nil
}
func findModelDigest(models map[string]string, wanted string) (string, bool) {
wanted = strings.TrimSpace(wanted)
for name, digest := range models {
if name == wanted || strings.TrimSuffix(name, ":latest") == strings.TrimSuffix(wanted, ":latest") {
return digest, true
}
}
return "", false
}
func minDuration(a, b time.Duration) time.Duration {
if a < b {
return a
}
return b
}
+78
View File
@@ -0,0 +1,78 @@
package ollama
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
)
func poolTestServer(t *testing.T, calls *atomic.Int64, fail bool) *httptest.Server {
t.Helper()
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch r.URL.Path {
case "/api/tags":
_ = json.NewEncoder(w).Encode(map[string]any{"models": []map[string]any{{"name": "qwen3:8b", "digest": "chat"}, {"name": "embeddinggemma:latest", "digest": "embed"}}})
case "/api/chat":
calls.Add(1)
if fail {
http.Error(w, "busy", http.StatusServiceUnavailable)
return
}
_ = json.NewEncoder(w).Encode(map[string]any{"message": map[string]any{"content": `{"ok":true}`}})
case "/api/embed":
calls.Add(1)
_ = json.NewEncoder(w).Encode(map[string]any{"embeddings": [][]float64{{0.1, 0.2}}})
default:
http.NotFound(w, r)
}
}))
}
func TestPoolLeastInflightBalancesSerialRequests(t *testing.T) {
var aCalls, bCalls atomic.Int64
a := poolTestServer(t, &aCalls, false)
defer a.Close()
b := poolTestServer(t, &bCalls, false)
defer b.Close()
c := NewPool(PoolConfig{Nodes: []NodeConfig{{Name: "a", URL: a.URL}, {Name: "b", URL: b.URL}}, RoutingMode: "least_inflight", NodeMaxInflight: 1, HealthInterval: time.Minute, FailureCooldown: time.Second, RequestTimeout: time.Second, FailoverEnabled: true, RequireSameModelDigest: true, RequireEmbeddingModel: true}, "qwen3:8b", "embeddinggemma")
if err := c.Ping(context.Background()); err != nil {
t.Fatal(err)
}
for i := 0; i < 4; i++ {
var out struct {
OK bool `json:"ok"`
}
if err := c.ChatJSON(context.Background(), "system", "user", map[string]any{"type": "object"}, &out); err != nil {
t.Fatal(err)
}
}
if aCalls.Load() != 2 || bCalls.Load() != 2 {
t.Fatalf("distribution a=%d b=%d", aCalls.Load(), bCalls.Load())
}
}
func TestPoolFailsOverOnRetryableError(t *testing.T) {
var badCalls, goodCalls atomic.Int64
bad := poolTestServer(t, &badCalls, true)
defer bad.Close()
good := poolTestServer(t, &goodCalls, false)
defer good.Close()
c := NewPool(PoolConfig{Nodes: []NodeConfig{{Name: "a-bad", URL: bad.URL}, {Name: "b-good", URL: good.URL}}, RoutingMode: "least_inflight", NodeMaxInflight: 1, HealthInterval: time.Minute, FailureCooldown: time.Second, RequestTimeout: time.Second, FailoverEnabled: true, FailoverAttempts: 2, RequireSameModelDigest: true, RequireEmbeddingModel: true}, "qwen3:8b", "embeddinggemma")
if err := c.Ping(context.Background()); err != nil {
t.Fatal(err)
}
var out struct {
OK bool `json:"ok"`
}
if err := c.ChatJSON(context.Background(), "system", "user", map[string]any{"type": "object"}, &out); err != nil {
t.Fatal(err)
}
if badCalls.Load() != 1 || goodCalls.Load() != 1 {
t.Fatalf("failover bad=%d good=%d", badCalls.Load(), goodCalls.Load())
}
}
+225
View File
@@ -0,0 +1,225 @@
package persist
import (
"context"
"fmt"
"os"
"path/filepath"
"sort"
"sync"
"time"
"github.com/local/glpi-neural-brain/internal/activity"
"github.com/local/glpi-neural-brain/internal/graph"
"github.com/local/glpi-neural-brain/internal/model"
)
type pendingFile struct {
path string
data []byte
perm os.FileMode
}
type Status struct {
Interval string `json:"interval"`
PendingFiles int `json:"pending_files"`
GraphDirty bool `json:"graph_dirty"`
LastFlush time.Time `json:"last_flush,omitempty"`
LastError string `json:"last_error,omitempty"`
LastPersistedVersion uint64 `json:"last_persisted_version"`
CurrentGraphVersion uint64 `json:"current_graph_version"`
SuccessfulFlushes uint64 `json:"successful_flushes"`
FailedFlushes uint64 `json:"failed_flushes"`
}
type Coordinator struct {
graph *graph.Store
broker *activity.Broker
interval time.Duration
mu sync.RWMutex
files map[string]pendingFile
lastFlush time.Time
lastError string
lastPersistedVersion uint64
successfulFlushes uint64
failedFlushes uint64
flushMu sync.Mutex
}
func New(g *graph.Store, b *activity.Broker, interval time.Duration) *Coordinator {
if interval <= 0 {
interval = 5 * time.Minute
}
version := uint64(0)
if g != nil {
version = g.Version()
}
return &Coordinator{graph: g, broker: b, interval: interval, files: make(map[string]pendingFile), lastPersistedVersion: version}
}
func (c *Coordinator) Start(ctx context.Context) {
go func() {
ticker := time.NewTicker(c.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := c.Flush(ctx, "interval"); err != nil && c.broker != nil {
c.broker.Publish(model.Activity{Type: "persistence.failed", Source: "brain", Phase: "storage", Message: "Gebündeltes Schreiben ist fehlgeschlagen", Strength: .35, Metadata: map[string]any{"error": err.Error()}})
}
}
}
}()
}
func (c *Coordinator) QueueFile(path string, data []byte, perm os.FileMode) (string, error) {
if path == "" {
return "", fmt.Errorf("persistence path is empty")
}
if perm == 0 {
perm = 0o640
}
abs, err := filepath.Abs(path)
if err != nil {
return "", err
}
c.mu.Lock()
c.files[abs] = pendingFile{path: abs, data: append([]byte(nil), data...), perm: perm}
c.mu.Unlock()
return abs, nil
}
func (c *Coordinator) Pending(path string) bool {
abs, err := filepath.Abs(path)
if err != nil {
return false
}
c.mu.RLock()
_, ok := c.files[abs]
c.mu.RUnlock()
return ok
}
func (c *Coordinator) Flush(ctx context.Context, trigger string) error {
c.flushMu.Lock()
defer c.flushMu.Unlock()
select {
case <-ctx.Done():
return ctx.Err()
default:
}
c.mu.RLock()
files := make([]pendingFile, 0, len(c.files))
for _, f := range c.files {
files = append(files, pendingFile{path: f.path, data: append([]byte(nil), f.data...), perm: f.perm})
}
lastVersion := c.lastPersistedVersion
c.mu.RUnlock()
sort.Slice(files, func(i, j int) bool { return files[i].path < files[j].path })
currentVersion := lastVersion
if c.graph != nil {
currentVersion = c.graph.Version()
}
graphDirty := currentVersion != lastVersion
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
// not yet reached disk.
for _, f := range files {
if err := writeAtomic(f.path, f.data, f.perm); err != nil {
c.recordFailure(err)
return err
}
c.mu.Lock()
if pending, ok := c.files[f.path]; ok && string(pending.data) == string(f.data) {
delete(c.files, f.path)
}
c.mu.Unlock()
}
if graphDirty && c.graph != nil {
persistedVersion, err := c.graph.PersistVersion()
if err != nil {
c.recordFailure(err)
return err
}
currentVersion = persistedVersion
}
now := time.Now().UTC()
c.mu.Lock()
c.lastFlush = now
c.lastError = ""
c.lastPersistedVersion = currentVersion
c.successfulFlushes++
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}})
}
return nil
}
func (c *Coordinator) recordFailure(err error) {
c.mu.Lock()
c.lastError = err.Error()
c.failedFlushes++
c.mu.Unlock()
}
func (c *Coordinator) Status() Status {
current := uint64(0)
if c.graph != nil {
current = c.graph.Version()
}
c.mu.RLock()
defer c.mu.RUnlock()
return Status{
Interval: c.interval.String(),
PendingFiles: len(c.files),
GraphDirty: current != c.lastPersistedVersion,
LastFlush: c.lastFlush,
LastError: c.lastError,
LastPersistedVersion: c.lastPersistedVersion,
CurrentGraphVersion: current,
SuccessfulFlushes: c.successfulFlushes,
FailedFlushes: c.failedFlushes,
}
}
func writeAtomic(path string, data []byte, perm os.FileMode) error {
if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
return err
}
tmp, err := os.CreateTemp(filepath.Dir(path), ".brain-write-*")
if err != nil {
return err
}
tmpName := tmp.Name()
defer os.Remove(tmpName)
if err := tmp.Chmod(perm); err != nil {
tmp.Close()
return err
}
if _, err := tmp.Write(data); err != nil {
tmp.Close()
return err
}
if err := tmp.Sync(); err != nil {
tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
return os.Rename(tmpName, path)
}
+52
View File
@@ -0,0 +1,52 @@
package persist
import (
"context"
"os"
"path/filepath"
"testing"
"time"
"github.com/local/glpi-neural-brain/internal/graph"
"github.com/local/glpi-neural-brain/internal/model"
)
func TestCoordinatorBatchesFilesAndGraph(t *testing.T) {
dir := t.TempDir()
g, err := graph.Open(dir)
if err != nil {
t.Fatal(err)
}
c := New(g, nil, time.Minute)
g.UpsertNode(model.Node{ID: "n1", Kind: "knowledge", Label: "One", Origin: "test"})
path := filepath.Join(dir, "staging", "draft.json")
if _, err := c.QueueFile(path, []byte("first"), 0o640); err != nil {
t.Fatal(err)
}
if _, err := c.QueueFile(path, []byte("second"), 0o640); err != nil {
t.Fatal(err)
}
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-state.json")); !os.IsNotExist(err) {
t.Fatalf("graph should remain buffered before flush: %v", err)
}
if err := c.Flush(context.Background(), "test"); err != nil {
t.Fatal(err)
}
data, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
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)
}
}
+30
View File
@@ -36,6 +36,8 @@ func (s *Server) Handler() http.Handler {
mux.HandleFunc("POST /api/events", s.handleEvent)
mux.HandleFunc("POST /api/reindex", s.handleReindex)
mux.HandleFunc("POST /api/enrich", s.handleEnrich)
mux.HandleFunc("POST /api/glpi-kb/sync", s.handleGLPIKBSync)
mux.HandleFunc("POST /api/flush", s.handleFlush)
sub, _ := fs.Sub(assets, "static")
mux.Handle("GET /", http.FileServer(http.FS(sub)))
return s.headers(mux)
@@ -109,6 +111,34 @@ func (s *Server) handleReindex(w http.ResponseWriter, r *http.Request) {
}
writeJSON(w, 200, s.Engine.Status())
}
func (s *Server) handleGLPIKBSync(w http.ResponseWriter, r *http.Request) {
if !s.authorized(r) {
writeJSON(w, 401, map[string]string{"error": "unauthorized"})
return
}
ctx, cancel := contextTimeout(r, 10*time.Minute)
defer cancel()
if err := s.Engine.SyncGLPIKB(ctx); err != nil {
writeJSON(w, 500, map[string]string{"error": err.Error()})
return
}
writeJSON(w, 200, s.Engine.Status())
}
func (s *Server) handleFlush(w http.ResponseWriter, r *http.Request) {
if !s.authorized(r) {
writeJSON(w, 401, map[string]string{"error": "unauthorized"})
return
}
ctx, cancel := contextTimeout(r, 2*time.Minute)
defer cancel()
if err := s.Engine.Flush(ctx); err != nil {
writeJSON(w, 500, map[string]string{"error": err.Error()})
return
}
writeJSON(w, 200, s.Engine.Status())
}
func (s *Server) handleEnrich(w http.ResponseWriter, r *http.Request) {
if !s.authorized(r) {
writeJSON(w, 401, map[string]string{"error": "unauthorized"})
+16 -3
View File
@@ -132,6 +132,12 @@
cls = 'waiting';
}
panel.className = `autonomy-status ${cls}`;
const pool = status.ollama_pool || {};
const persistence = status.persistence || {};
const glpi = status.glpi_kb || {};
const diagnostics = [`Ollama ${pool.healthy_nodes ?? 0}/${pool.node_count ?? 0}`, `Diskqueue ${persistence.pending_files ?? 0}`];
if (glpi.enabled) diagnostics.push(`GLPI-KB ${glpi.documents ?? 0}`);
panel.title = diagnostics.join(' · ');
const span = panel.querySelector('span');
if (span) span.textContent = text;
button.disabled = running;
@@ -1269,7 +1275,7 @@
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.created', 'think.rejected', 'think.failed', 'think.paused', 'research.started', 'agent.run']);
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.created', 'think.rejected', 'think.failed', 'think.paused', 'research.started', '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;
@@ -1302,12 +1308,19 @@
'think.failed': 'AI-THINK Fehler',
'think.paused': 'AI-THINK pausiert',
'research.started': 'Recherche gestartet',
'agent.run': 'Agent-Lauf'
'agent.run': 'Agent-Lauf',
'glpi.kb.synced': 'GLPI-KB synchronisiert',
'glpi.kb.failed': 'GLPI-KB Fehler',
'persistence.flushed': 'Gebündelt gespeichert',
'persistence.failed': 'Speicherfehler'
};
const title = titleMap[evt.type] || evt.phase || evt.source || 'Aktivität';
const meta = [];
if (evt.metadata?.nodes) meta.push(`${Number(evt.metadata.nodes).toLocaleString('de-DE')} Nodes`);
if (evt.metadata?.edges) meta.push(`${Number(evt.metadata.edges).toLocaleString('de-DE')} Edges`);
if (evt.metadata?.documents) meta.push(`${Number(evt.metadata.documents).toLocaleString('de-DE')} GLPI-Beiträge`);
if (evt.metadata?.files !== undefined) meta.push(`${Number(evt.metadata.files).toLocaleString('de-DE')} Dateien`);
if (evt.metadata?.graph_version !== undefined) meta.push(`Graph v${Number(evt.metadata.graph_version).toLocaleString('de-DE')}`);
if (evt.metadata?.duration_ms) meta.push(`${Number(evt.metadata.duration_ms).toLocaleString('de-DE')} ms`);
if (evt.metadata?.hit_count) meta.push(`${Number(evt.metadata.hit_count).toLocaleString('de-DE')} Treffer`);
if (evt.metadata?.used_nodes) meta.push(`${Number(evt.metadata.used_nodes).toLocaleString('de-DE')} Quellen`);
@@ -1337,7 +1350,7 @@
} else if (evt.type === 'query.started' && evt.query) {
message = `${evt.source === 'agent' ? 'Agent' : evt.source === 'knowledgebase' ? 'Knowledgebase' : 'Brain'} verarbeitet eine Anfrage.`;
}
return {time, title, message, meta, query: eventQuery, regions, cls: evt.type?.includes('think') ? 'think' : evt.type?.includes('research') ? 'research' : evt.type === 'graph.updated' || evt.type === 'scan.started' ? 'graph' : evt.source === 'agent' ? 'agent' : ''};
return {time, title, message, meta, query: eventQuery, regions, cls: evt.type?.includes('think') ? '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) {
BIN
View File
Binary file not shown.