Persistenz und Performance optimiert
All checks were successful
release-tag / release-image (push) Successful in 1m40s

This commit is contained in:
2026-07-29 08:56:04 +02:00
parent 720b2aecdd
commit 827c348e69
9 changed files with 1002 additions and 84 deletions

View File

@@ -231,7 +231,7 @@ Ein Auto-Reply ist nur erlaubt, wenn `language` und `communication_style` des fr
`categories` begrenzt Auto-Reply auf die angegebenen Zielkategorien. Eine leere Liste bedeutet keine zusätzliche Kategorie-Einschränkung. `min_score` kann die globale Schwelle je Artikel verschärfen.
Bei aktiviertem RAG erzeugt Ollama Embeddings über `/api/embed`; der Cache landet in `data/embeddings.json`. Für Ticket und Knowledge wird dasselbe Embedding-Modell verwendet.
Bei aktiviertem RAG erzeugt Ollama Embeddings über `/api/embed`. Der aktive lokale Index wird persistent unter `DATA_DIR/knowledge-index/snapshot.gob` gespeichert. Ein vorhandenes altes `data/embeddings.json` wird nur noch als einmalige Migrationsquelle verwendet. Für Ticket und Knowledge wird dasselbe Embedding-Modell verwendet.
### Realistisches Hybrid-Scoring und dynamische Kandidatenauswahl
@@ -515,6 +515,26 @@ Die interne Knowledge Base kann bei `KNOWLEDGE_WEB_EDIT_ENABLED=true` direkt im
### GLPI-KB Rich Text
Rich-Text-Formatierungen aus synchronisierten GLPI-KB-Artikeln bleiben in Ticketantworten erhalten; RAG und LLM sehen weiterhin nur bereinigten Plaintext.
### Große Knowledge-Verzeichnisse
### Große Knowledge-Verzeichnisse und persistenter Index
Bei großen lokalen Korpora startet das Dashboard sofort und zeigt den Hintergrundaufbau des Knowledge-Index an. Der Agent verarbeitet noch keine Tickets, solange die lokale KB nicht `ready` ist. Unter `/api/status` stehen unter anderem `knowledge_init_phase`, `knowledge_init_processed_files`, `knowledge_init_loaded_docs`, `knowledge_init_indexed_docs`, `knowledge_init_cache_hits` und `knowledge_init_pending_embeddings` zur Verfügung.
Für große lokale Korpora wird der vollständige lokale Index nach dem ersten erfolgreichen Aufbau persistent unter `DATA_DIR/knowledge-index/snapshot.gob` gespeichert. Beim normalen Neustart mit `KNOWLEDGE_INDEX_MODE=incremental` wird dieser Snapshot zuerst geladen; die Ticketverarbeitung kann anschließend mit dem letzten konsistenten Index starten. Die Quelldateien werden danach im Hintergrund inkrementell geprüft.
Empfohlene Werte:
```env
KNOWLEDGE_INDEX_MODE=incremental
KNOWLEDGE_EMBED_BATCH_SIZE=64
KNOWLEDGE_INDEX_SCAN_INTERVAL=5m
```
Beim Delta-Scan werden zunächst nur Dateiname, Größe und `mtime` geprüft. Unveränderte Dateien werden weder geöffnet noch geparst. Erst bei geänderten Metadaten wird der Dateiinhalt gelesen und gehasht; nur tatsächlich geänderte Retrieval-Inhalte werden erneut an Ollama `/api/embed` geschickt. Gelöschte Dateien werden aus dem Index entfernt.
Index-Modi:
- `incremental`: vorhandenen Snapshot sofort laden und Änderungen im Hintergrund nachziehen. Empfohlen.
- `rebuild`: Quelldateien beim Start vollständig neu einlesen; gültige Embeddings aus dem alten Cache können bei der Migration weiterhin wiederverwendet werden.
- `readonly`: ausschließlich einen vorhandenen persistenten Snapshot verwenden; ohne kompatiblen Snapshot schlägt der Start fehl.
`/api/status` und das Dashboard zeigen unter anderem Snapshot-Zeitpunkt, letzten Delta-Scan, geänderte/gelöschte Dateien, wiederverwendete Vektoren, Embedding-Batchgröße und Scanintervall. Das komplette Vektorindex-Map wird bei einer Suche nicht mehr pro Ticket kopiert; Suchläufe lesen den warmen Index direkt unter einem Read-Lock.
Beim allerersten Aufbau ohne Snapshot startet das Dashboard weiterhin sofort und zeigt Scan-/Embedding-Fortschritt. Die Ticketverarbeitung wartet in diesem Fall, bis der erste konsistente Index fertig ist.

View File

@@ -155,4 +155,30 @@ Währenddessen gilt:
- `/api/status` und das Dashboard zeigen Phase, Datei-/Dokumentfortschritt, Cache-Treffer, offene Embeddings und Fehler.
- Bei einem fehlerhaften KB-Dokument bleibt das WebUI erreichbar und zeigt den Initialisierungsfehler an.
Die Embedding-Erzeugung verarbeitet große Korpora dokumentweise in Batches und schreibt periodische Cache-Checkpoints. Dadurch wird ein Verzeichnis mit vielen tausend Dateien nicht mehr als ein riesiger Embedding-Request im Speicher aufgebaut.
Die Embedding-Erzeugung verarbeitet große Korpora dokumentweise in Batches. Nach dem ersten vollständigen Aufbau wird ein atomarer persistenter Snapshot geschrieben; spätere Starts verwenden diesen Snapshot und führen nur Delta-Scans aus.
## Persistenter inkrementeller Knowledge-Index
Für große lokale KB-Bestände sollte die bestehende `.env` ergänzt werden:
```env
KNOWLEDGE_INDEX_MODE=incremental
KNOWLEDGE_EMBED_BATCH_SIZE=64
KNOWLEDGE_INDEX_SCAN_INTERVAL=5m
```
Der neue Snapshot liegt unter `DATA_DIR/knowledge-index/snapshot.gob`. Bei Docker muss `DATA_DIR` deshalb dauerhaft gemountet und für UID/GID `65532:65532` beschreibbar bleiben. `docker compose down -v` bzw. das Löschen des Host-Verzeichnisses entfernt auch den persistenten Index.
Beim ersten Start dieser Version existiert noch kein Snapshot. Der Agent kann vorhandene gültige Vektoren aus dem bisherigen `DATA_DIR/embeddings.json` übernehmen und schreibt nach erfolgreichem Aufbau den neuen Snapshot. Danach wird `embeddings.json` für die lokale KB nicht mehr als primärer Index benötigt.
Normaler Neustart in `incremental`:
1. Snapshot laden.
2. Knowledge sofort als `ready` markieren.
3. Ticketverarbeitung starten.
4. Quelldateien im Hintergrund per Größe/`mtime` vergleichen.
5. Nur geänderte Dateien lesen/hashen/parsen und nur geänderte Retrieval-Texte neu embedden.
6. Geänderten Snapshot atomar ersetzen.
`KNOWLEDGE_INDEX_MODE=rebuild` erzwingt einen vollständigen Quellen-Scan. `readonly` verwendet ausschließlich den vorhandenen Snapshot und führt keine lokalen Delta-Scans aus.

View File

@@ -72,6 +72,7 @@ func main() {
k, err := knowledge.NewStore(cfg.KnowledgeDir, cfg.DataDir, o, cfg.RAGEnabled, cfg.KnowledgeAllowedSources, knowledge.ScoringConfig{
SemanticWeight: cfg.KnowledgeSemanticWeight, TitleWeight: cfg.KnowledgeTitleWeight, LexicalWeight: cfg.KnowledgeLexicalWeight, KeywordWeight: cfg.KnowledgeKeywordWeight, CategoryWeight: cfg.KnowledgeCategoryWeight,
EmbeddingProfile: embeddingProfile, EmbeddingIdentity: cfg.OllamaEmbeddingModel, ChunkWords: cfg.KnowledgeChunkWords, ChunkOverlap: cfg.KnowledgeChunkOverlapWords, MaxChunksPerDoc: cfg.KnowledgeMaxChunksPerDoc, MaxQueryChunks: cfg.KnowledgeMaxQueryChunks,
IndexMode: cfg.KnowledgeIndexMode, EmbedBatchSize: cfg.KnowledgeEmbedBatchSize, IndexScanInterval: cfg.KnowledgeIndexScanInterval,
CategoryMode: cfg.KnowledgeCategoryMode, CategoryMapFile: cfg.KnowledgeCategoryMapFile, IgnoreGlobs: cfg.KnowledgeIgnoreGlobs,
})
if err != nil {
@@ -132,7 +133,8 @@ func main() {
m.SetKnowledgeDocs(k.Count())
}
svc.Start(ctx)
slog.Info("ticket processing started", "knowledge_docs", k.Count())
k.StartIncrementalSync(ctx, cfg.KnowledgeIndexScanInterval)
slog.Info("ticket processing started", "knowledge_docs", k.Count(), "knowledge_index_mode", cfg.KnowledgeIndexMode)
}()
<-ctx.Done()
shutdownCtx, c := context.WithTimeout(context.Background(), 10*time.Second)

View File

@@ -68,6 +68,9 @@ type Config struct {
KnowledgeChunkOverlapWords int
KnowledgeMaxChunksPerDoc int
KnowledgeMaxQueryChunks int
KnowledgeIndexMode string
KnowledgeEmbedBatchSize int
KnowledgeIndexScanInterval time.Duration
GLPIKBEnabled bool
GLPIKBPath string
GLPIKBFilter string
@@ -182,6 +185,9 @@ func Load() (Config, error) {
KnowledgeChunkOverlapWords: envInt("KNOWLEDGE_CHUNK_OVERLAP_WORDS", 30),
KnowledgeMaxChunksPerDoc: envInt("KNOWLEDGE_MAX_CHUNKS_PER_DOC", 24),
KnowledgeMaxQueryChunks: envInt("KNOWLEDGE_MAX_QUERY_CHUNKS", 64),
KnowledgeIndexMode: strings.ToLower(env("KNOWLEDGE_INDEX_MODE", "incremental")),
KnowledgeEmbedBatchSize: envInt("KNOWLEDGE_EMBED_BATCH_SIZE", 64),
KnowledgeIndexScanInterval: envDuration("KNOWLEDGE_INDEX_SCAN_INTERVAL", 5*time.Minute),
GLPIKBEnabled: envBool("GLPI_KB_ENABLED", false),
GLPIKBPath: env("GLPI_KB_PATH", "auto"),
GLPIKBFilter: strings.TrimSpace(os.Getenv("GLPI_KB_FILTER")),
@@ -352,6 +358,17 @@ func (c Config) Validate() error {
if c.KnowledgeMaxQueryChunks != 0 && (c.KnowledgeMaxQueryChunks < 1 || c.KnowledgeMaxQueryChunks > 200) {
return errors.New("KNOWLEDGE_MAX_QUERY_CHUNKS must be between 1 and 200")
}
switch c.KnowledgeIndexMode {
case "", "incremental", "rebuild", "readonly":
default:
return errors.New("KNOWLEDGE_INDEX_MODE must be one of: incremental, rebuild, readonly")
}
if c.KnowledgeEmbedBatchSize != 0 && (c.KnowledgeEmbedBatchSize < 1 || c.KnowledgeEmbedBatchSize > 256) {
return errors.New("KNOWLEDGE_EMBED_BATCH_SIZE must be between 1 and 256")
}
if c.KnowledgeIndexScanInterval < 0 {
return errors.New("KNOWLEDGE_INDEX_SCAN_INTERVAL must be >= 0")
}
if c.LearningEnabled {
if c.LearningMaxExamples < 1 || c.LearningMaxExamples > 10000 {
return errors.New("LEARNING_MAX_EXAMPLES must be between 1 and 10000")

View File

@@ -0,0 +1,694 @@
package knowledge
import (
"context"
"crypto/sha256"
"encoding/gob"
"encoding/hex"
"encoding/json"
"fmt"
"log/slog"
"os"
"path/filepath"
"sort"
"strings"
"time"
"github.com/example/glpi-ai-agent/internal/model"
)
const persistentSnapshotVersion = 1
type fileRecord struct {
Key string
Path string
ID string
Size int64
ModTimeUnixNano int64
RawHash string
Managed bool
Included bool
Ignored bool
Unmapped []string
}
type persistentSnapshot struct {
Version int
Fingerprint string
SavedAt time.Time
Docs []model.KnowledgeDoc
Files map[string]string
Managed map[string]bool
StaticDocs map[string]model.KnowledgeDoc
TitleVectors map[string][]float64
ChunkVectors map[string][][]float64
Chunks map[string][]string
Manifest map[string]fileRecord
LoadStats LoadStats
}
func (s *Store) indexFingerprint() string {
allowed := make([]string, 0, len(s.allowedSources))
for source := range s.allowedSources {
allowed = append(allowed, source)
}
sort.Strings(allowed)
mapKeys := make([]string, 0, len(s.categoryMap))
for k := range s.categoryMap {
mapKeys = append(mapKeys, k)
}
sort.Strings(mapKeys)
mapping := make([]struct {
Key string
IDs []int64
}, 0, len(mapKeys))
for _, k := range mapKeys {
mapping = append(mapping, struct {
Key string
IDs []int64
}{k, append([]int64(nil), s.categoryMap[k]...)})
}
payload := struct {
RAG bool
EmbeddingIdentity string
EmbeddingProfile string
ChunkWords int
ChunkOverlap int
MaxChunksPerDoc int
CategoryMode string
IgnoreGlobs []string
AllowedSources []string
CategoryMap any
}{
RAG: s.rag, EmbeddingIdentity: s.scoring.EmbeddingIdentity, EmbeddingProfile: s.scoring.EmbeddingProfile,
ChunkWords: s.scoring.ChunkWords, ChunkOverlap: s.scoring.ChunkOverlap, MaxChunksPerDoc: s.scoring.MaxChunksPerDoc,
CategoryMode: s.loadOptions.CategoryMode, IgnoreGlobs: append([]string(nil), s.loadOptions.IgnoreGlobs...), AllowedSources: allowed, CategoryMap: mapping,
}
b, _ := json.Marshal(payload)
h := sha256.Sum256(b)
return hex.EncodeToString(h[:])
}
func (s *Store) loadPersistentSnapshot() (bool, error) {
f, err := os.Open(s.snapshotPath)
if err != nil {
if os.IsNotExist(err) {
return false, nil
}
return false, fmt.Errorf("open persistent knowledge index: %w", err)
}
defer f.Close()
var snap persistentSnapshot
if err := gob.NewDecoder(f).Decode(&snap); err != nil {
return false, fmt.Errorf("decode persistent knowledge index: %w", err)
}
if snap.Version != persistentSnapshotVersion {
return false, fmt.Errorf("persistent knowledge index version %d is unsupported", snap.Version)
}
if snap.Fingerprint != s.indexFingerprint() {
return false, nil
}
if snap.Files == nil {
snap.Files = map[string]string{}
}
if snap.Managed == nil {
snap.Managed = map[string]bool{}
}
if snap.StaticDocs == nil {
snap.StaticDocs = map[string]model.KnowledgeDoc{}
}
if snap.TitleVectors == nil {
snap.TitleVectors = map[string][]float64{}
}
if snap.ChunkVectors == nil {
snap.ChunkVectors = map[string][][]float64{}
}
if snap.Chunks == nil {
snap.Chunks = map[string][]string{}
}
if snap.Manifest == nil {
snap.Manifest = map[string]fileRecord{}
}
// Paths in a snapshot may come from another host/project location. Rebind
// them to the current knowledge roots using the stable manifest key.
reboundFiles := map[string]string{}
for key, rec := range snap.Manifest {
name := filepath.Base(key)
if strings.HasPrefix(key, "managed/") {
rec.Path = filepath.Join(s.managedDir, name)
} else {
rec.Path = filepath.Join(s.dir, name)
}
snap.Manifest[key] = rec
if rec.Included && rec.ID != "" {
if rec.Managed || reboundFiles[rec.ID] == "" {
reboundFiles[rec.ID] = rec.Path
}
}
}
snap.Files = reboundFiles
s.mu.Lock()
s.docs = snap.Docs
s.files = snap.Files
s.managed = snap.Managed
s.staticDocs = snap.StaticDocs
s.titleVectors = snap.TitleVectors
s.chunkVectors = snap.ChunkVectors
s.chunks = snap.Chunks
s.manifest = snap.Manifest
s.external = map[string]string{}
s.loadStats = snap.LoadStats
s.initStatus = InitStatus{
State: "ready", Phase: "ready-cache", TotalFiles: len(snap.Manifest), ProcessedFiles: len(snap.Manifest), LoadedDocs: len(snap.Docs),
IndexedDocs: len(snap.Docs), CacheHits: len(snap.Docs), PendingEmbeddings: 0, StartedAt: time.Now(), FinishedAt: time.Now(),
SnapshotLoaded: true, SnapshotPath: s.snapshotPath, SnapshotSavedAt: snap.SavedAt,
}
s.mu.Unlock()
slog.Info("persistent knowledge index loaded", "documents", len(snap.Docs), "files", len(snap.Manifest), "saved_at", snap.SavedAt, "path", s.snapshotPath)
return true, nil
}
func (s *Store) persistSnapshot() error {
if s == nil {
return nil
}
if err := os.MkdirAll(filepath.Dir(s.snapshotPath), 0o750); err != nil {
return err
}
s.mu.RLock()
localDocs := make([]model.KnowledgeDoc, 0, len(s.docs))
files := map[string]string{}
managed := map[string]bool{}
staticDocs := make(map[string]model.KnowledgeDoc, len(s.staticDocs))
titles := map[string][]float64{}
vectors := map[string][][]float64{}
chunks := map[string][]string{}
for id, d := range s.staticDocs {
staticDocs[id] = d
}
for _, d := range s.docs {
if s.external[d.ID] != "" {
continue
}
localDocs = append(localDocs, d)
if path := s.files[d.ID]; path != "" {
files[d.ID] = path
}
if s.managed[d.ID] {
managed[d.ID] = true
}
if v := s.titleVectors[d.ID]; len(v) > 0 {
titles[d.ID] = v
}
if vv := s.chunkVectors[d.ID]; len(vv) > 0 {
vectors[d.ID] = vv
}
if cc := s.chunks[d.ID]; len(cc) > 0 {
chunks[d.ID] = cc
}
}
manifest := make(map[string]fileRecord, len(s.manifest))
for k, v := range s.manifest {
manifest[k] = v
}
stats := s.loadStats
stats.UnmappedCategories = append([]string(nil), stats.UnmappedCategories...)
s.mu.RUnlock()
snap := persistentSnapshot{Version: persistentSnapshotVersion, Fingerprint: s.indexFingerprint(), SavedAt: time.Now(), Docs: localDocs, Files: files, Managed: managed, StaticDocs: staticDocs, TitleVectors: titles, ChunkVectors: vectors, Chunks: chunks, Manifest: manifest, LoadStats: stats}
tmp := s.snapshotPath + ".tmp"
f, err := os.OpenFile(tmp, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o640)
if err != nil {
return err
}
encErr := gob.NewEncoder(f).Encode(&snap)
closeErr := f.Close()
if encErr != nil {
_ = os.Remove(tmp)
return encErr
}
if closeErr != nil {
_ = os.Remove(tmp)
return closeErr
}
if err := os.Rename(tmp, s.snapshotPath); err != nil {
_ = os.Remove(tmp)
return err
}
s.setInitStatus(func(st *InitStatus) {
st.SnapshotPath = s.snapshotPath
st.SnapshotSavedAt = snap.SavedAt
})
return nil
}
// StartIncrementalSync keeps a warm persistent index usable while source files
// are checked for changes in the background. The first delta scan runs
// immediately. A zero interval disables subsequent periodic scans.
func (s *Store) StartIncrementalSync(ctx context.Context, interval time.Duration) {
if s == nil || strings.EqualFold(s.scoring.IndexMode, "readonly") {
return
}
go func() {
s.syncLocalSafely(ctx)
if interval <= 0 {
return
}
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
s.syncLocalSafely(ctx)
}
}
}()
}
func (s *Store) syncLocalSafely(ctx context.Context) {
if err := s.SyncLocal(ctx); err != nil {
slog.Error("incremental knowledge scan failed; keeping previous index", "error", err)
s.setInitStatus(func(st *InitStatus) { st.LastScanAt = time.Now(); st.LastScanError = err.Error(); st.Phase = "ready" })
}
}
// SyncLocal performs a metadata-first delta scan. Unchanged files are never
// opened or parsed. Files whose size/mtime changed are hashed; unchanged bytes
// are reused without parsing/embedding. Only genuinely changed retrieval text
// is sent to the embedding provider.
func (s *Store) SyncLocal(ctx context.Context) error {
if s == nil {
return fmt.Errorf("knowledge store is not initialized")
}
if strings.EqualFold(s.scoring.IndexMode, "readonly") {
return nil
}
s.initMu.Lock()
defer s.initMu.Unlock()
if !s.Ready() {
return fmt.Errorf("knowledge store is not ready")
}
s.setInitStatus(func(st *InitStatus) {
st.Phase = "syncing"
st.LastScanError = ""
st.ChangedFiles = 0
st.DeletedFiles = 0
st.ReusedFiles = 0
})
s.mu.RLock()
oldManifest := make(map[string]fileRecord, len(s.manifest))
for k, v := range s.manifest {
oldManifest[k] = v
}
oldDocs := make(map[string]model.KnowledgeDoc, len(s.docs))
for _, d := range s.docs {
oldDocs[d.ID] = d
}
oldStatic := make(map[string]model.KnowledgeDoc, len(s.staticDocs))
for k, v := range s.staticDocs {
oldStatic[k] = v
}
oldManaged := make(map[string]model.KnowledgeDoc)
for _, d := range s.docs {
if s.managed[d.ID] {
oldManaged[d.ID] = d
}
}
oldTitle := s.titleVectors
oldChunkVec := s.chunkVectors
oldChunks := s.chunks
externalDocs := make([]model.KnowledgeDoc, 0)
externalMap := make(map[string]string, len(s.external))
for id, src := range s.external {
externalMap[id] = src
}
for _, d := range s.docs {
if s.external[d.ID] != "" {
externalDocs = append(externalDocs, d)
}
}
s.mu.RUnlock()
resStatic, err := s.scanDeltaDir(ctx, s.dir, "static", false, s.loadOptions, oldManifest, oldStatic)
if err != nil {
return err
}
resManaged, err := s.scanDeltaDir(ctx, s.managedDir, "managed", true, LoadOptions{CategoryMode: "strict"}, oldManifest, oldManaged)
if err != nil {
return err
}
manifest := make(map[string]fileRecord, len(resStatic.manifest)+len(resManaged.manifest))
for k, v := range resStatic.manifest {
manifest[k] = v
}
for k, v := range resManaged.manifest {
manifest[k] = v
}
// Rebuild effective local documents. Managed entries intentionally override
// static entries with the same id, matching the original startup behavior.
staticDocs := resStatic.docs
managedDocs := resManaged.docs
ids := make([]string, 0, len(staticDocs)+len(managedDocs))
merged := map[string]model.KnowledgeDoc{}
for id, d := range staticDocs {
merged[id] = d
ids = append(ids, id)
}
for id, d := range managedDocs {
if _, ok := merged[id]; !ok {
ids = append(ids, id)
}
merged[id] = d
}
sort.Strings(ids)
localDocs := make([]model.KnowledgeDoc, 0, len(ids))
for _, id := range ids {
localDocs = append(localDocs, merged[id])
}
newTitle := make(map[string][]float64, len(localDocs)+len(externalDocs))
newChunkVec := make(map[string][][]float64, len(localDocs)+len(externalDocs))
newChunks := make(map[string][]string, len(localDocs)+len(externalDocs))
needEmbed := make([]model.KnowledgeDoc, 0)
reused := 0
for _, d := range localDocs {
parts := chunkText(d.Text, s.scoring.ChunkWords, s.scoring.ChunkOverlap, s.scoring.MaxChunksPerDoc)
newChunks[d.ID] = parts
if old, ok := oldDocs[d.ID]; ok && hashDoc(old, s.scoring) == hashDoc(d, s.scoring) && len(oldTitle[d.ID]) > 0 && len(oldChunkVec[d.ID]) == len(parts) {
newTitle[d.ID] = oldTitle[d.ID]
newChunkVec[d.ID] = oldChunkVec[d.ID]
reused++
} else if s.rag {
needEmbed = append(needEmbed, d)
}
}
if s.rag && len(needEmbed) > 0 {
if s.embedder == nil {
return fmt.Errorf("RAG is enabled but no embedding provider is configured")
}
const docsPerBatch = 20
for start := 0; start < len(needEmbed); start += docsPerBatch {
if err := ctx.Err(); err != nil {
return err
}
end := start + docsPerBatch
if end > len(needEmbed) {
end = len(needEmbed)
}
emb, err := s.embedDocuments(ctx, needEmbed[start:end])
if err != nil {
return err
}
for _, d := range needEmbed[start:end] {
newTitle[d.ID] = emb[d.ID].title
newChunkVec[d.ID] = emb[d.ID].chunks
}
}
}
// Preserve connector-backed documents and vectors; their own syncers manage them.
for _, d := range externalDocs {
if _, collision := merged[d.ID]; collision {
return fmt.Errorf("local knowledge id %q collides with external source %q", d.ID, externalMap[d.ID])
}
newTitle[d.ID] = oldTitle[d.ID]
newChunkVec[d.ID] = oldChunkVec[d.ID]
newChunks[d.ID] = oldChunks[d.ID]
}
files := map[string]string{}
managedMap := map[string]bool{}
for _, rec := range manifest {
if !rec.Included || rec.ID == "" {
continue
}
if rec.Managed {
managedMap[rec.ID] = true
files[rec.ID] = rec.Path
} else if !managedMap[rec.ID] {
files[rec.ID] = rec.Path
}
}
stats := mergeLoadStats(resStatic.stats, resManaged.stats)
allDocs := append(localDocs, externalDocs...)
changed := resStatic.changed + resManaged.changed
deleted := countDeleted(oldManifest, manifest)
s.mu.Lock()
s.docs = allDocs
s.files = files
s.managed = managedMap
s.staticDocs = staticDocs
s.titleVectors = newTitle
s.chunkVectors = newChunkVec
s.chunks = newChunks
s.manifest = manifest
s.loadStats = stats
s.initStatus.State = "ready"
s.initStatus.Phase = "ready"
s.initStatus.TotalFiles = len(manifest)
s.initStatus.ProcessedFiles = len(manifest)
s.initStatus.LoadedDocs = len(localDocs)
s.initStatus.IndexedDocs = len(localDocs)
s.initStatus.PendingEmbeddings = 0
s.initStatus.LastScanAt = time.Now()
s.initStatus.ChangedFiles = changed
s.initStatus.DeletedFiles = deleted
s.initStatus.ReusedFiles = reused
s.initStatus.LastScanError = ""
s.mu.Unlock()
if changed > 0 || deleted > 0 {
if err := s.persistSnapshot(); err != nil {
return fmt.Errorf("persist incremental knowledge index: %w", err)
}
}
slog.Info("incremental knowledge scan complete", "files", len(manifest), "documents", len(localDocs), "changed_files", changed, "deleted_files", deleted, "reused_vectors", reused, "embedded_documents", len(needEmbed))
return nil
}
type deltaScanResult struct {
docs map[string]model.KnowledgeDoc
manifest map[string]fileRecord
stats LoadStats
changed int
}
func (s *Store) scanDeltaDir(ctx context.Context, dir, origin string, managed bool, opts LoadOptions, old map[string]fileRecord, oldDocs map[string]model.KnowledgeDoc) (deltaScanResult, error) {
res := deltaScanResult{docs: map[string]model.KnowledgeDoc{}, manifest: map[string]fileRecord{}}
entries, err := os.ReadDir(dir)
if err != nil {
return res, fmt.Errorf("read knowledge directory %q: %w", dir, err)
}
for _, e := range entries {
if err := ctx.Err(); err != nil {
return res, err
}
if e.IsDir() || !strings.HasSuffix(strings.ToLower(e.Name()), ".json") {
continue
}
key := origin + "/" + e.Name()
path := filepath.Join(dir, e.Name())
info, err := e.Info()
if err != nil {
return res, err
}
if matchesAnyGlob(e.Name(), opts.IgnoreGlobs) {
res.manifest[key] = fileRecord{Key: key, Path: path, Size: info.Size(), ModTimeUnixNano: info.ModTime().UnixNano(), Managed: managed, Ignored: true}
res.stats.IgnoredFiles++
continue
}
if prev, ok := old[key]; ok && prev.Size == info.Size() && prev.ModTimeUnixNano == info.ModTime().UnixNano() {
prev.Path = path
res.manifest[key] = prev
if prev.Included && prev.ID != "" {
if d, ok := oldDocs[prev.ID]; ok {
res.docs[prev.ID] = d
}
}
for _, u := range prev.Unmapped {
res.stats.UnmappedCategories = appendUniqueString(res.stats.UnmappedCategories, u)
}
if len(prev.Unmapped) > 0 {
res.stats.UnmappedCategoryFiles++
}
if prev.Ignored {
res.stats.IgnoredFiles++
}
continue
}
b, err := os.ReadFile(path)
if err != nil {
return res, err
}
h := sha256.Sum256(b)
rawHash := hex.EncodeToString(h[:])
if prev, ok := old[key]; ok && prev.RawHash != "" && prev.RawHash == rawHash {
prev.Path = path
prev.Size = info.Size()
prev.ModTimeUnixNano = info.ModTime().UnixNano()
res.manifest[key] = prev
if prev.Included && prev.ID != "" {
if d, ok := oldDocs[prev.ID]; ok {
res.docs[prev.ID] = d
}
}
continue
}
d, unmapped, skip, err := decodeKnowledgeDoc(b, opts.CategoryMode, s.categoryMap)
if err != nil {
return res, fmt.Errorf("%s: %w", e.Name(), err)
}
rec := fileRecord{Key: key, Path: path, Size: info.Size(), ModTimeUnixNano: info.ModTime().UnixNano(), RawHash: rawHash, Managed: managed, Unmapped: append([]string(nil), unmapped...)}
if len(unmapped) > 0 {
res.stats.UnmappedCategoryFiles++
res.stats.UnmappedCategories = mergeStrings(res.stats.UnmappedCategories, unmapped)
}
if skip {
rec.Ignored = true
res.stats.IgnoredFiles++
res.manifest[key] = rec
res.changed++
continue
}
if d.ID == "" || d.Title == "" {
return res, fmt.Errorf("%s: id/title required", e.Name())
}
if !safeID(d.ID) {
return res, fmt.Errorf("%s: invalid id %q", e.Name(), d.ID)
}
d.Source = strings.ToLower(strings.TrimSpace(d.Source))
if d.Source == "" {
return res, fmt.Errorf("%s: source required", e.Name())
}
if _, allowed := s.allowedSources[d.Source]; !allowed {
res.manifest[key] = rec
res.changed++
continue
}
d.Language = strings.TrimSpace(d.Language)
d.CommunicationStyle = strings.ToLower(strings.TrimSpace(d.CommunicationStyle))
rec.ID = d.ID
rec.Included = true
res.manifest[key] = rec
if _, dup := res.docs[d.ID]; dup {
return res, fmt.Errorf("duplicate knowledge id %q in %s", d.ID, dir)
}
res.docs[d.ID] = d
res.changed++
}
return res, nil
}
func countDeleted(old, cur map[string]fileRecord) int {
n := 0
for k := range old {
if _, ok := cur[k]; !ok {
n++
}
}
return n
}
func mergeLoadStats(a, b LoadStats) LoadStats {
return LoadStats{IgnoredFiles: a.IgnoredFiles + b.IgnoredFiles, UnmappedCategoryFiles: a.UnmappedCategoryFiles + b.UnmappedCategoryFiles, UnmappedCategories: mergeStrings(a.UnmappedCategories, b.UnmappedCategories)}
}
func buildManifestForDocs(docs []model.KnowledgeDoc, files []string, origin string, managed bool) (map[string]fileRecord, error) {
out := make(map[string]fileRecord, len(docs))
for i, d := range docs {
if i >= len(files) {
return nil, fmt.Errorf("knowledge manifest mismatch: %d docs, %d files", len(docs), len(files))
}
path := files[i]
info, err := os.Stat(path)
if err != nil {
return nil, err
}
b, err := os.ReadFile(path)
if err != nil {
return nil, err
}
h := sha256.Sum256(b)
key := origin + "/" + filepath.Base(path)
out[key] = fileRecord{Key: key, Path: path, ID: d.ID, Size: info.Size(), ModTimeUnixNano: info.ModTime().UnixNano(), RawHash: hex.EncodeToString(h[:]), Managed: managed, Included: true, Unmapped: append([]string(nil), d.UnmappedExternalCategories...)}
}
return out, nil
}
func (s *Store) persistExternalVectorCache(source string) error {
if s == nil || !s.rag {
return nil
}
if err := os.MkdirAll(filepath.Dir(s.externalCachePath), 0o750); err != nil {
return err
}
s.mu.RLock()
cf := cacheFile{Version: 3, Hashes: map[string]string{}, TitleVectors: map[string][]float64{}, ChunkVectors: map[string][][]float64{}}
for _, d := range s.docs {
if s.external[d.ID] != source || len(s.titleVectors[d.ID]) == 0 {
continue
}
cf.Hashes[d.ID] = hashDoc(d, s.scoring)
cf.TitleVectors[d.ID] = append([]float64(nil), s.titleVectors[d.ID]...)
cf.ChunkVectors[d.ID] = cloneChunkVectors(s.chunkVectors[d.ID])
}
s.mu.RUnlock()
b, err := json.Marshal(cf)
if err != nil {
return err
}
tmp := s.externalCachePath + ".tmp"
if err := os.WriteFile(tmp, b, 0o640); err != nil {
return err
}
if err := os.Rename(tmp, s.externalCachePath); err != nil {
_ = os.Remove(tmp)
return err
}
return nil
}
// augmentManifestAllFiles records even filtered/ignored JSON files so periodic
// delta scans do not repeatedly open files that are intentionally not part of
// the active corpus.
func augmentManifestAllFiles(dir, origin string, managed bool, opts LoadOptions, manifest map[string]fileRecord) error {
entries, err := os.ReadDir(dir)
if err != nil {
return err
}
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(strings.ToLower(e.Name()), ".json") {
continue
}
key := origin + "/" + e.Name()
if _, ok := manifest[key]; ok {
continue
}
info, err := e.Info()
if err != nil {
return err
}
path := filepath.Join(dir, e.Name())
rec := fileRecord{Key: key, Path: path, Size: info.Size(), ModTimeUnixNano: info.ModTime().UnixNano(), Managed: managed, Included: false, Ignored: matchesAnyGlob(e.Name(), opts.IgnoreGlobs)}
if !rec.Ignored {
b, err := os.ReadFile(path)
if err != nil {
return err
}
h := sha256.Sum256(b)
rec.RawHash = hex.EncodeToString(h[:])
}
manifest[key] = rec
}
return nil
}

View File

@@ -35,6 +35,9 @@ type ScoringConfig struct {
ChunkOverlap int
MaxChunksPerDoc int
MaxQueryChunks int
IndexMode string
EmbedBatchSize int
IndexScanInterval time.Duration
CategoryMode string
CategoryMapFile string
IgnoreGlobs []string
@@ -67,30 +70,41 @@ type InitStatus struct {
StartedAt time.Time `json:"started_at,omitempty"`
FinishedAt time.Time `json:"finished_at,omitempty"`
LastError string `json:"last_error,omitempty"`
SnapshotLoaded bool `json:"snapshot_loaded"`
SnapshotPath string `json:"snapshot_path,omitempty"`
SnapshotSavedAt time.Time `json:"snapshot_saved_at,omitempty"`
LastScanAt time.Time `json:"last_scan_at,omitempty"`
LastScanError string `json:"last_scan_error,omitempty"`
ChangedFiles int `json:"changed_files"`
DeletedFiles int `json:"deleted_files"`
ReusedFiles int `json:"reused_files"`
}
type Store struct {
mu sync.RWMutex
initMu sync.Mutex
initStatus InitStatus
dir string
managedDir string
docs []model.KnowledgeDoc
files map[string]string
managed map[string]bool
external map[string]string
staticDocs map[string]model.KnowledgeDoc
titleVectors map[string][]float64
chunkVectors map[string][][]float64
chunks map[string][]string
embedder Embedder
rag bool
cachePath string
allowedSources map[string]struct{}
scoring ScoringConfig
loadOptions LoadOptions
loadStats LoadStats
categoryMap map[string][]int64
mu sync.RWMutex
initMu sync.Mutex
initStatus InitStatus
dir string
managedDir string
docs []model.KnowledgeDoc
files map[string]string
managed map[string]bool
external map[string]string
staticDocs map[string]model.KnowledgeDoc
titleVectors map[string][]float64
chunkVectors map[string][][]float64
chunks map[string][]string
embedder Embedder
rag bool
cachePath string
allowedSources map[string]struct{}
scoring ScoringConfig
loadOptions LoadOptions
loadStats LoadStats
categoryMap map[string][]int64
manifest map[string]fileRecord
snapshotPath string
externalCachePath string
}
type cacheFile struct {
Version int `json:"version,omitempty"`
@@ -100,7 +114,7 @@ type cacheFile struct {
}
func DefaultScoringConfig() ScoringConfig {
return ScoringConfig{SemanticWeight: .45, TitleWeight: .20, LexicalWeight: .20, KeywordWeight: .075, CategoryWeight: .075, EmbeddingProfile: "plain", ChunkWords: 160, ChunkOverlap: 30, MaxChunksPerDoc: 24, MaxQueryChunks: 64}
return ScoringConfig{SemanticWeight: .45, TitleWeight: .20, LexicalWeight: .20, KeywordWeight: .075, CategoryWeight: .075, EmbeddingProfile: "plain", ChunkWords: 160, ChunkOverlap: 30, MaxChunksPerDoc: 24, MaxQueryChunks: 64, IndexMode: "incremental", EmbedBatchSize: 64, IndexScanInterval: 5 * time.Minute}
}
// ResolveEmbeddingProfile selects prompt formatting for the configured embedding model.
@@ -136,6 +150,16 @@ func normalizeScoring(c ScoringConfig) ScoringConfig {
if c.MaxQueryChunks <= 0 {
c.MaxQueryChunks = d.MaxQueryChunks
}
if strings.TrimSpace(c.IndexMode) == "" {
c.IndexMode = d.IndexMode
}
c.IndexMode = strings.ToLower(strings.TrimSpace(c.IndexMode))
if c.EmbedBatchSize <= 0 {
c.EmbedBatchSize = d.EmbedBatchSize
}
if c.IndexScanInterval < 0 {
c.IndexScanInterval = d.IndexScanInterval
}
return c
}
@@ -161,7 +185,8 @@ func NewStore(dir, dataDir string, embedder Embedder, rag bool, allowedSources [
titleVectors: map[string][]float64{}, chunkVectors: map[string][][]float64{}, chunks: map[string][]string{},
files: map[string]string{}, managed: map[string]bool{}, external: map[string]string{}, staticDocs: map[string]model.KnowledgeDoc{},
embedder: embedder, rag: rag, cachePath: filepath.Join(dataDir, "embeddings.json"), allowedSources: map[string]struct{}{},
scoring: scoreCfg, loadOptions: loadOpts, categoryMap: categoryMap,
scoring: scoreCfg, loadOptions: loadOpts, categoryMap: categoryMap, manifest: map[string]fileRecord{},
snapshotPath: filepath.Join(dataDir, "knowledge-index", "snapshot.gob"), externalCachePath: filepath.Join(dataDir, "knowledge-index", "external-embeddings.json"),
initStatus: InitStatus{State: "waiting", Phase: "waiting"},
}
for _, source := range allowedSources {
@@ -191,11 +216,32 @@ func (s *Store) Initialize(ctx context.Context) (err error) {
if s == nil {
return fmt.Errorf("knowledge store is not initialized")
}
mode := strings.ToLower(strings.TrimSpace(s.scoring.IndexMode))
if mode == "" {
mode = "incremental"
}
if mode != "rebuild" {
loaded, loadErr := s.loadPersistentSnapshot()
if loadErr != nil {
if mode == "readonly" {
return loadErr
}
slog.Warn("persistent knowledge index unavailable; falling back to rebuild", "error", loadErr, "path", s.snapshotPath)
} else if loaded {
return nil
} else if mode == "readonly" {
return fmt.Errorf("KNOWLEDGE_INDEX_MODE=readonly requires a compatible persistent index at %s", s.snapshotPath)
}
}
return s.fullRebuild(ctx)
}
func (s *Store) fullRebuild(ctx context.Context) (err error) {
s.initMu.Lock()
defer s.initMu.Unlock()
s.setInitStatus(func(st *InitStatus) {
*st = InitStatus{State: "loading", Phase: "scanning", StartedAt: time.Now()}
*st = InitStatus{State: "loading", Phase: "scanning", StartedAt: time.Now(), SnapshotPath: s.snapshotPath}
})
defer func() {
if err != nil {
@@ -239,6 +285,28 @@ func (s *Store) Initialize(ctx context.Context) (err error) {
stats.UnmappedCategoryFiles += managedStats.UnmappedCategoryFiles
stats.UnmappedCategories = mergeStrings(stats.UnmappedCategories, managedStats.UnmappedCategories)
staticManifest, err := buildManifestForDocs(static, staticFiles, "static", false)
if err != nil {
return fmt.Errorf("build static knowledge manifest: %w", err)
}
managedManifest, err := buildManifestForDocs(managed, managedFiles, "managed", true)
if err != nil {
return fmt.Errorf("build managed knowledge manifest: %w", err)
}
manifest := make(map[string]fileRecord, len(staticManifest)+len(managedManifest))
for k, v := range staticManifest {
manifest[k] = v
}
for k, v := range managedManifest {
manifest[k] = v
}
if err := augmentManifestAllFiles(s.dir, "static", false, s.loadOptions, manifest); err != nil {
return fmt.Errorf("complete static knowledge manifest: %w", err)
}
if err := augmentManifestAllFiles(s.managedDir, "managed", true, LoadOptions{CategoryMode: "strict"}, manifest); err != nil {
return fmt.Errorf("complete managed knowledge manifest: %w", err)
}
files := map[string]string{}
managedMap := map[string]bool{}
staticMap := map[string]model.KnowledgeDoc{}
@@ -269,8 +337,10 @@ func (s *Store) Initialize(ctx context.Context) (err error) {
s.docs = docs
s.files = files
s.managed = managedMap
s.external = map[string]string{}
s.staticDocs = staticMap
s.loadStats = stats
s.manifest = manifest
s.titleVectors = map[string][]float64{}
s.chunkVectors = map[string][][]float64{}
s.chunks = map[string][]string{}
@@ -291,6 +361,9 @@ func (s *Store) Initialize(ctx context.Context) (err error) {
return err
}
}
if err := s.persistSnapshot(); err != nil {
return fmt.Errorf("persist knowledge snapshot: %w", err)
}
s.setInitStatus(func(st *InitStatus) {
st.State = "ready"
st.Phase = "ready"
@@ -298,8 +371,11 @@ func (s *Store) Initialize(ctx context.Context) (err error) {
st.PendingEmbeddings = 0
st.FinishedAt = time.Now()
st.LastError = ""
st.SnapshotLoaded = false
st.SnapshotPath = s.snapshotPath
st.SnapshotSavedAt = time.Now()
})
slog.Info("knowledge store ready", "documents", len(docs), "rag_enabled", s.rag)
slog.Info("knowledge store ready", "documents", len(docs), "rag_enabled", s.rag, "persistent_index", s.snapshotPath)
return nil
}
@@ -723,6 +799,12 @@ func (s *Store) Upsert(ctx context.Context, d model.KnowledgeDoc) error {
_ = os.Remove(tmp)
return err
}
info, statErr := os.Stat(path)
if statErr != nil {
return statErr
}
rawSum := sha256.Sum256(b)
manifestKey := "managed/" + filepath.Base(path)
s.mu.Lock()
replaced := false
for i := range s.docs {
@@ -737,6 +819,7 @@ func (s *Store) Upsert(ctx context.Context, d model.KnowledgeDoc) error {
}
s.files[d.ID] = path
s.managed[d.ID] = true
s.manifest[manifestKey] = fileRecord{Key: manifestKey, Path: path, ID: d.ID, Size: info.Size(), ModTimeUnixNano: info.ModTime().UnixNano(), RawHash: hex.EncodeToString(rawSum[:]), Managed: true, Included: true, Unmapped: append([]string(nil), d.UnmappedExternalCategories...)}
s.chunks[d.ID] = chunks
if s.rag {
s.titleVectors[d.ID] = titleVector
@@ -780,6 +863,11 @@ func (s *Store) Delete(id string) error {
s.docs = append([]model.KnowledgeDoc(nil), out...)
delete(s.files, id)
delete(s.managed, id)
for key, rec := range s.manifest {
if rec.Managed && rec.ID == id {
delete(s.manifest, key)
}
}
delete(s.titleVectors, id)
delete(s.chunkVectors, id)
delete(s.chunks, id)
@@ -833,7 +921,10 @@ func (s *Store) ReplaceExternalSource(ctx context.Context, source string, docs [
oldDocs[d.ID] = d
}
s.mu.RUnlock()
cached := loadCache(s.cachePath)
cached := loadCache(s.externalCachePath)
if len(cached.Hashes) == 0 {
cached = loadCache(s.cachePath) // one-time migration from the legacy combined cache
}
changed := make([]model.KnowledgeDoc, 0)
seen := map[string]struct{}{}
@@ -922,33 +1013,11 @@ func (s *Store) ReplaceExternalSource(ctx context.Context, source string, docs [
}
s.docs = rebuilt
s.mu.Unlock()
return s.persistVectorCache()
return s.persistExternalVectorCache(source)
}
func (s *Store) persistVectorCache() error {
if s == nil || !s.rag {
return nil
}
s.mu.RLock()
cf := cacheFile{Version: 3, Hashes: map[string]string{}, TitleVectors: map[string][]float64{}, ChunkVectors: map[string][][]float64{}}
for _, d := range s.docs {
if len(s.titleVectors[d.ID]) == 0 {
continue
}
cf.Hashes[d.ID] = hashDoc(d, s.scoring)
cf.TitleVectors[d.ID] = append([]float64(nil), s.titleVectors[d.ID]...)
cf.ChunkVectors[d.ID] = cloneChunkVectors(s.chunkVectors[d.ID])
}
s.mu.RUnlock()
b, err := json.MarshalIndent(cf, "", " ")
if err != nil {
return err
}
tmp := s.cachePath + ".tmp"
if err := os.WriteFile(tmp, b, 0o640); err != nil {
return err
}
return os.Rename(tmp, s.cachePath)
return s.persistSnapshot()
}
func (s *Store) ManagedDir() string {
@@ -978,15 +1047,10 @@ func (s *Store) Search(ctx context.Context, text string, topK int, categorySets
return nil, fmt.Errorf("knowledge store is not initialized")
}
s.mu.RLock()
docs := append([]model.KnowledgeDoc(nil), s.docs...)
titleVecs := cloneVectorMap(s.titleVectors)
chunkVecs := cloneChunkVectorMap(s.chunkVectors)
chunks := cloneStringSliceMap(s.chunks)
scoreCfg := s.scoring
ragEnabled := s.rag
embedder := s.embedder
s.mu.RUnlock()
if len(docs) == 0 {
return nil, nil
}
var cats []model.Category
if len(categorySets) > 0 {
cats = categorySets[0]
@@ -999,16 +1063,16 @@ func (s *Store) Search(ctx context.Context, text string, topK int, categorySets
}
var queryVectors [][]float64
var queryTitleVector []float64
if s.rag && s.embedder != nil {
if ragEnabled && embedder != nil {
if len(queryChunks) > 0 {
q, err := s.embedTexts(ctx, formatQueryEmbeddings(queryChunks, scoreCfg.EmbeddingProfile), 64)
q, err := s.embedTexts(ctx, formatQueryEmbeddings(queryChunks, scoreCfg.EmbeddingProfile), scoreCfg.EmbedBatchSize)
if err != nil {
return nil, err
}
queryVectors = q
}
if strings.TrimSpace(queryTitle) != "" {
tq, err := s.embedder.Embed(ctx, formatQueryEmbeddings([]string{queryTitle}, scoreCfg.EmbeddingProfile))
tq, err := embedder.Embed(ctx, formatQueryEmbeddings([]string{queryTitle}, scoreCfg.EmbeddingProfile))
if err != nil {
return nil, err
}
@@ -1018,19 +1082,24 @@ func (s *Store) Search(ctx context.Context, text string, topK int, categorySets
}
}
hits := make([]model.KnowledgeHit, 0, len(docs))
for _, d := range docs {
s.mu.RLock()
defer s.mu.RUnlock()
if len(s.docs) == 0 {
return nil, nil
}
hits := make([]model.KnowledgeHit, 0, len(s.docs))
for _, d := range s.docs {
semantic, bestChunk, bestQueryChunk := 0.0, "", ""
semanticAvailable := false
if len(queryVectors) > 0 && len(chunkVecs[d.ID]) > 0 {
if len(queryVectors) > 0 && len(s.chunkVectors[d.ID]) > 0 {
semanticAvailable = true
for qi, qv := range queryVectors {
for di, dv := range chunkVecs[d.ID] {
for di, dv := range s.chunkVectors[d.ID] {
score := clamp01(cosine(qv, dv))
if score > semantic || bestChunk == "" {
semantic = score
if di < len(chunks[d.ID]) {
bestChunk = chunks[d.ID][di]
if di < len(s.chunks[d.ID]) {
bestChunk = s.chunks[d.ID][di]
}
if qi < len(queryChunks) {
bestQueryChunk = queryChunks[qi]
@@ -1040,7 +1109,7 @@ func (s *Store) Search(ctx context.Context, text string, topK int, categorySets
}
} else if strings.TrimSpace(d.Text) != "" {
semanticAvailable = true
docChunks := chunks[d.ID]
docChunks := s.chunks[d.ID]
if len(docChunks) == 0 {
docChunks = []string{d.Text}
}
@@ -1062,8 +1131,8 @@ func (s *Store) Search(ctx context.Context, text string, topK int, categorySets
titleQuery = text
}
title = titleSimilarity(titleQuery, d.Title)
if len(queryTitleVector) > 0 && len(titleVecs[d.ID]) > 0 {
title = math.Max(title, clamp01(cosine(queryTitleVector, titleVecs[d.ID])))
if len(queryTitleVector) > 0 && len(s.titleVectors[d.ID]) > 0 {
title = math.Max(title, clamp01(cosine(queryTitleVector, s.titleVectors[d.ID])))
}
}
lexicalScore := lexicalSimilarity(text, d)
@@ -1081,7 +1150,7 @@ func (s *Store) Search(ctx context.Context, text string, topK int, categorySets
scorePart{keyword, scoreCfg.KeywordWeight, keywordAvailable},
scorePart{category, scoreCfg.CategoryWeight, categoryAvailable},
)
hits = append(hits, model.KnowledgeHit{Doc: d, Score: total, SemanticScore: semantic, TitleScore: title, LexicalScore: lexicalScore, KeywordScore: keyword, CategoryScore: category, BestChunkExcerpt: excerpt(bestChunk, 280), BestQueryExcerpt: excerpt(bestQueryChunk, 280), QueryChunkCount: len(queryChunks), DocumentChunkCount: len(chunks[d.ID])})
hits = append(hits, model.KnowledgeHit{Doc: d, Score: total, SemanticScore: semantic, TitleScore: title, LexicalScore: lexicalScore, KeywordScore: keyword, CategoryScore: category, BestChunkExcerpt: excerpt(bestChunk, 280), BestQueryExcerpt: excerpt(bestQueryChunk, 280), QueryChunkCount: len(queryChunks), DocumentChunkCount: len(s.chunks[d.ID])})
}
sort.SliceStable(hits, func(i, j int) bool {
if hits[i].Score == hits[j].Score {
@@ -1170,7 +1239,6 @@ func (s *Store) index(ctx context.Context) error {
// thousands of files can otherwise allocate hundreds of thousands of
// embedding input strings before the first request is sent to Ollama.
const docsPerBatch = 20
const checkpointEvery = 1000
embeddedDocs := 0
for start := 0; start < len(need); start += docsPerBatch {
if err := ctx.Err(); err != nil {
@@ -1201,11 +1269,6 @@ func (s *Store) index(ctx context.Context) error {
if embeddedDocs%500 == 0 || embeddedDocs == len(need) {
slog.Info("knowledge embedding progress", "indexed_docs", indexed, "total_docs", len(s.docs), "cache_hits", cacheHits, "pending_embeddings", pending)
}
if embeddedDocs%checkpointEvery == 0 {
if err := s.persistVectorCache(); err != nil {
slog.Warn("knowledge embedding checkpoint failed", "error", err, "indexed_docs", indexed)
}
}
}
return s.persistVectorCache()
}
@@ -1233,7 +1296,7 @@ func (s *Store) embedDocuments(ctx context.Context, docs []model.KnowledgeDoc) (
refs = append(refs, ref{id: d.ID, chunk: i})
}
}
vectors, err := s.embedTexts(ctx, texts, 64)
vectors, err := s.embedTexts(ctx, texts, s.scoring.EmbedBatchSize)
if err != nil {
return nil, err
}

View File

@@ -8,6 +8,7 @@ import (
"path/filepath"
"strings"
"testing"
"time"
"github.com/example/glpi-ai-agent/internal/model"
)
@@ -422,3 +423,97 @@ func TestNewStoreSupportsBackgroundInitialization(t *testing.T) {
t.Fatalf("unexpected init status: %+v count=%d", st, s.Count())
}
}
type countingEmbedder struct{ calls int }
func (e *countingEmbedder) Embed(_ context.Context, texts []string) ([][]float64, error) {
e.calls++
out := make([][]float64, len(texts))
for i, text := range texts {
v := make([]float64, 8)
for j, b := range []byte(strings.ToLower(text)) {
v[(int(b)+j)%len(v)] += 1
}
out[i] = v
}
return out, nil
}
func TestPersistentIndexLoadsWithoutReembeddingAndSyncsDelta(t *testing.T) {
dir := t.TempDir()
data := t.TempDir()
path := filepath.Join(dir, "kb.json")
writeDoc := func(title, text string) {
t.Helper()
d := model.KnowledgeDoc{ID: "KB-1", Title: title, Text: text, Answer: "Antwort", Source: "internal-kb", Language: "de-DE", CommunicationStyle: "formal"}
b, _ := json.Marshal(d)
if err := os.WriteFile(path, b, 0o644); err != nil {
t.Fatal(err)
}
}
writeDoc("Anmeldung", "Benutzer kann sich nicht anmelden")
cfg := ScoringConfig{EmbeddingIdentity: "test-embed", EmbeddingProfile: "plain", ChunkWords: 40, ChunkOverlap: 10, MaxChunksPerDoc: 8, MaxQueryChunks: 8, IndexMode: "incremental", EmbedBatchSize: 8}
firstEmbed := &countingEmbedder{}
first, err := Load(context.Background(), dir, data, firstEmbed, true, []string{"internal-kb"}, cfg)
if err != nil {
t.Fatal(err)
}
if firstEmbed.calls == 0 {
t.Fatal("expected initial embedding calls")
}
if _, err := os.Stat(filepath.Join(data, "knowledge-index", "snapshot.gob")); err != nil {
t.Fatalf("snapshot missing: %v", err)
}
if !first.Ready() {
t.Fatal("first store not ready")
}
secondEmbed := &countingEmbedder{}
second, err := Load(context.Background(), dir, data, secondEmbed, true, []string{"internal-kb"}, cfg)
if err != nil {
t.Fatal(err)
}
if secondEmbed.calls != 0 {
t.Fatalf("snapshot startup unexpectedly re-embedded: calls=%d", secondEmbed.calls)
}
if !second.InitStatus().SnapshotLoaded {
t.Fatal("expected persistent snapshot to be loaded")
}
if second.Count() != 1 {
t.Fatalf("count=%d", second.Count())
}
// Ensure mtime changes even on filesystems with coarse timestamp resolution.
time.Sleep(20 * time.Millisecond)
writeDoc("Anmeldung geändert", "Benutzer kann sich weiterhin nicht anmelden")
if err := second.SyncLocal(context.Background()); err != nil {
t.Fatal(err)
}
if secondEmbed.calls == 0 {
t.Fatal("expected changed document to be re-embedded")
}
got, ok := second.ByID("KB-1")
if !ok || got.Title != "Anmeldung geändert" {
t.Fatalf("delta update not applied: %+v", got)
}
callsAfterChange := secondEmbed.calls
if err := second.SyncLocal(context.Background()); err != nil {
t.Fatal(err)
}
if secondEmbed.calls != callsAfterChange {
t.Fatalf("unchanged delta scan re-embedded document: before=%d after=%d", callsAfterChange, secondEmbed.calls)
}
}
func TestReadonlyIndexRequiresSnapshot(t *testing.T) {
dir := t.TempDir()
data := t.TempDir()
s, err := NewStore(dir, data, nil, false, []string{"internal-kb"}, ScoringConfig{IndexMode: "readonly"})
if err != nil {
t.Fatal(err)
}
if err := s.Initialize(context.Background()); err == nil {
t.Fatal("expected readonly mode without snapshot to fail")
}
}

View File

@@ -120,6 +120,7 @@ func (s *Server) status(w http.ResponseWriter, r *http.Request) {
"processed": s.metrics.Processed.Load(), "skipped": s.metrics.Skipped.Load(), "errors": s.metrics.Errors.Load(), "category_changes": s.metrics.CategoryChanged.Load(), "replies": s.metrics.Replies.Load(), "queue_depth": s.q.Len(),
"glpi_ok": g, "ollama_ok": o, "knowledge_docs": s.metrics.KnowledgeDocs(), "last_poll": s.metrics.LastPoll(),
"knowledge_ready": s.knowledge.Ready(), "knowledge_init_state": initStatus.State, "knowledge_init_phase": initStatus.Phase, "knowledge_init_total_files": initStatus.TotalFiles, "knowledge_init_processed_files": initStatus.ProcessedFiles, "knowledge_init_loaded_docs": initStatus.LoadedDocs, "knowledge_init_indexed_docs": initStatus.IndexedDocs, "knowledge_init_cache_hits": initStatus.CacheHits, "knowledge_init_pending_embeddings": initStatus.PendingEmbeddings, "knowledge_init_started_at": initStatus.StartedAt, "knowledge_init_finished_at": initStatus.FinishedAt, "knowledge_init_error": initStatus.LastError,
"knowledge_index_mode": s.cfg.KnowledgeIndexMode, "knowledge_embed_batch_size": s.cfg.KnowledgeEmbedBatchSize, "knowledge_index_scan_interval": s.cfg.KnowledgeIndexScanInterval.String(), "knowledge_snapshot_loaded": initStatus.SnapshotLoaded, "knowledge_snapshot_path": initStatus.SnapshotPath, "knowledge_snapshot_saved_at": initStatus.SnapshotSavedAt, "knowledge_last_scan_at": initStatus.LastScanAt, "knowledge_last_scan_error": initStatus.LastScanError, "knowledge_changed_files": initStatus.ChangedFiles, "knowledge_deleted_files": initStatus.DeletedFiles, "knowledge_reused_files": initStatus.ReusedFiles,
"communication_language": s.cfg.CommunicationLanguage, "communication_style": s.cfg.CommunicationStyle, "knowledge_allowed_sources": s.cfg.KnowledgeAllowedSources, "knowledge_auto_reply_sources": s.cfg.KnowledgeAutoReplySources, "knowledge_category_mode": s.cfg.KnowledgeCategoryMode, "knowledge_category_map_configured": strings.TrimSpace(s.cfg.KnowledgeCategoryMapFile) != "", "knowledge_ignore_globs": s.cfg.KnowledgeIgnoreGlobs, "knowledge_ignored_files": loadStats.IgnoredFiles, "knowledge_unmapped_category_files": loadStats.UnmappedCategoryFiles, "knowledge_unmapped_categories": loadStats.UnmappedCategories,
"category_confidence": s.cfg.CategoryConfidence, "reply_confidence": s.cfg.ReplyConfidence, "knowledge_min_score": s.cfg.KnowledgeMinScore, "knowledge_retrieval_floor": s.cfg.KnowledgeRetrievalFloor, "knowledge_evidence_weight_retrieval": s.cfg.KnowledgeEvidenceRetrievalWeight, "knowledge_evidence_weight_ai": s.cfg.KnowledgeEvidenceAIWeight, "knowledge_evidence_weight_category": s.cfg.KnowledgeEvidenceCategoryWeight,
"knowledge_weight_semantic": s.cfg.KnowledgeSemanticWeight, "knowledge_weight_title": s.cfg.KnowledgeTitleWeight, "knowledge_weight_lexical": s.cfg.KnowledgeLexicalWeight, "knowledge_weight_keywords": s.cfg.KnowledgeKeywordWeight, "knowledge_weight_category": s.cfg.KnowledgeCategoryWeight, "knowledge_embedding_profile": s.cfg.KnowledgeEmbeddingProfile,

File diff suppressed because one or more lines are too long