Files
glpi-neural-brain/internal/engine/article_research_cache.go
jbergner 440423c5b6
All checks were successful
release-tag / release-image (push) Successful in 2m43s
RC-3
2026-08-09 11:29:13 +02:00

262 lines
9.3 KiB
Go

package engine
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"log/slog"
"os"
"path/filepath"
"sort"
"strings"
"time"
"github.com/local/glpi-neural-brain/internal/graph"
"github.com/local/glpi-neural-brain/internal/model"
)
const researchEvidenceSchemaVersion = 1
type researchEvidenceRecord struct {
SchemaVersion int `json:"schema_version"`
SavedAt time.Time `json:"saved_at"`
ContentSHA256 string `json:"content_sha256"`
Result model.ResearchResult `json:"result"`
}
// queueResearchEvidence stores accepted full-text evidence outside graph.db.
// The graph node only keeps a small summary and a portable relative path, so
// the browser graph payload does not grow by the full page content.
func (e *Engine) queueResearchEvidence(result model.ResearchResult) (string, string, error) {
if e.Persistence == nil || strings.TrimSpace(e.Cfg.DataDir) == "" {
return "", "", fmt.Errorf("research evidence persistence is unavailable")
}
if !result.Fetched || !result.Relevant || strings.TrimSpace(result.Content) == "" || strings.TrimSpace(result.URL) == "" {
return "", "", fmt.Errorf("research evidence is not an accepted full-text result")
}
hashBytes := sha256.Sum256([]byte(result.URL + "\x00" + result.Content))
contentHash := hex.EncodeToString(hashBytes[:])
id := graph.ID("external", result.URL)
relPath := filepath.ToSlash(filepath.Join("research-evidence", strings.ToLower(id)+".json"))
record := researchEvidenceRecord{SchemaVersion: researchEvidenceSchemaVersion, SavedAt: time.Now().UTC(), ContentSHA256: contentHash, Result: result}
data, err := json.MarshalIndent(record, "", " ")
if err != nil {
return "", "", fmt.Errorf("encode research evidence: %w", err)
}
absPath := filepath.Join(e.Cfg.DataDir, filepath.FromSlash(relPath))
if _, err := e.Persistence.QueueFile(absPath, append(data, '\n'), 0o640); err != nil {
return "", "", fmt.Errorf("queue research evidence %q: %w", absPath, err)
}
e.cacheResearchEvidence(relPath, record)
return relPath, contentHash, nil
}
// researchEvidenceForSources reuses full-text evidence learned in earlier
// synthesis cycles when it is already linked to one of the selected sources.
// The knowledge model still has to confirm that the evidence closes a current
// gap; reuse never bypasses the consolidation or article quality gates.
func (e *Engine) researchEvidenceForSources(sources []articleSource) []model.ResearchResult {
if len(sources) == 0 || strings.TrimSpace(e.Cfg.DataDir) == "" {
return nil
}
sourceIDs := map[string]bool{}
for _, source := range sources {
sourceIDs[source.Node.ID] = true
}
snapshot := e.Graph.Snapshot()
externalIDs := map[string]bool{}
groundedByIDs := map[string]bool{}
for _, edge := range snapshot.Edges {
if edge.Status == "rejected" {
continue
}
switch edge.Type {
case "research_evidence":
// Legacy research_evidence edges are only candidate links. They may
// originate from versions that materialised every fetched page before
// final article review. Keep them as lookup candidates, but require the
// external node itself to carry validation_state=grounded below.
if sourceIDs[edge.Target] {
externalIDs[edge.Source] = true
}
case "grounded_by":
if sourceIDs[edge.Source] {
externalIDs[edge.Target] = true
groundedByIDs[edge.Target] = true
}
}
}
if len(externalIDs) == 0 {
return nil
}
nodes := make([]model.Node, 0, len(externalIDs))
filter := e.effectiveThinkingFilter()
for _, node := range snapshot.Nodes {
if !externalIDs[node.ID] || node.Kind != "external" || !filter.Matches(node) {
continue
}
validationState := strings.ToLower(strings.TrimSpace(metadataString(node.Metadata, "validation_state")))
if !groundedByIDs[node.ID] && validationState != "grounded" {
continue
}
nodes = append(nodes, node)
}
sort.SliceStable(nodes, func(i, j int) bool {
if nodes[i].Weight == nodes[j].Weight {
return nodes[i].ID < nodes[j].ID
}
return nodes[i].Weight > nodes[j].Weight
})
limit := e.Cfg.ArticleResearchResults * 2
if limit < 8 {
limit = 8
}
if limit > 24 {
limit = 24
}
out := make([]model.ResearchResult, 0, min(limit, len(nodes)))
for _, node := range nodes {
if len(out) >= limit {
break
}
relPath := metadataString(node.Metadata, "evidence_path")
if relPath == "" {
continue
}
record, err := e.loadResearchEvidence(relPath)
if err != nil {
slog.Warn("stored research evidence could not be reused", "node_id", node.ID, "path", relPath, "error", err)
continue
}
result := record.Result
if result.Relevance < e.Cfg.ArticleResearchMinRelevance || result.SourceQualityScore < e.Cfg.ArticleResearchMinQuality {
continue
}
out = append(out, result)
}
return uniqueResearchEvidence(out)
}
func (e *Engine) loadResearchEvidence(relPath string) (researchEvidenceRecord, error) {
var record researchEvidenceRecord
clean := filepath.Clean(filepath.FromSlash(strings.TrimSpace(relPath)))
if clean == "." || clean == "" || filepath.IsAbs(clean) || clean == ".." || strings.HasPrefix(clean, ".."+string(filepath.Separator)) {
return record, fmt.Errorf("unsafe research evidence path %q", relPath)
}
cacheKey := filepath.ToSlash(clean)
if cached, ok := e.cachedResearchEvidence(cacheKey); ok {
return cached, nil
}
root, err := filepath.Abs(e.Cfg.DataDir)
if err != nil {
return record, fmt.Errorf("resolve data directory: %w", err)
}
path := filepath.Join(root, clean)
rel, err := filepath.Rel(root, path)
if err != nil || rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) {
return record, fmt.Errorf("research evidence path escapes data directory")
}
data, err := os.ReadFile(path)
if err != nil {
return record, err
}
if err := json.Unmarshal(data, &record); err != nil {
return record, fmt.Errorf("decode research evidence: %w", err)
}
if record.SchemaVersion != researchEvidenceSchemaVersion {
return record, fmt.Errorf("unsupported research evidence schema %d", record.SchemaVersion)
}
if !record.Result.Fetched || !record.Result.Relevant || strings.TrimSpace(record.Result.Content) == "" || strings.TrimSpace(record.Result.URL) == "" {
return record, fmt.Errorf("stored research evidence is incomplete")
}
hashBytes := sha256.Sum256([]byte(record.Result.URL + "\x00" + record.Result.Content))
if record.ContentSHA256 == "" || !strings.EqualFold(record.ContentSHA256, hex.EncodeToString(hashBytes[:])) {
return record, fmt.Errorf("stored research evidence checksum mismatch")
}
e.cacheResearchEvidence(cacheKey, record)
return record, nil
}
func (e *Engine) cacheResearchEvidence(key string, record researchEvidenceRecord) {
key = filepath.ToSlash(filepath.Clean(filepath.FromSlash(strings.TrimSpace(key))))
if key == "." || key == "" {
return
}
e.researchEvidenceMu.Lock()
if e.researchEvidenceCache == nil {
e.researchEvidenceCache = map[string]researchEvidenceRecord{}
}
e.researchEvidenceCache[key] = record
e.researchEvidenceMu.Unlock()
}
func (e *Engine) cachedResearchEvidence(key string) (researchEvidenceRecord, bool) {
e.researchEvidenceMu.RLock()
record, ok := e.researchEvidenceCache[key]
e.researchEvidenceMu.RUnlock()
return record, ok
}
// learnResearchEvidence embeds newly accepted external evidence immediately.
// It is still persisted by the normal batched graph flush, but becomes usable
// for semantic placement and later relations in the current process at once.
func (e *Engine) learnResearchEvidence(ctx context.Context, results []model.ResearchResult) graph.MutationStats {
var stats graph.MutationStats
if !e.LearningEnabled() || len(results) == 0 || e.Ollama == nil {
return stats
}
ids := make([]string, 0, len(results))
texts := make([]string, 0, len(results))
for _, result := range results {
id := graph.ID("external", result.URL)
node, ok := e.Graph.GetNode(id)
if !ok {
continue
}
content := strings.TrimSpace(result.Content)
if content == "" {
content = node.Summary
}
text := strings.TrimSpace(result.Title + "\n" + clamp(content, 6000))
if text == "" {
continue
}
ids = append(ids, id)
texts = append(texts, text)
}
if len(ids) == 0 {
return stats
}
vectors, err := e.Ollama.Embed(ctx, texts)
fallback := err != nil || len(vectors) != len(ids)
if !fallback {
for _, vector := range vectors {
if len(vector) == 0 {
fallback = true
break
}
}
}
if fallback {
for i, id := range ids {
stats.Add(e.Graph.SetVectorWithStats(id, hashEmbedding(texts[i], 256)))
}
e.Broker.Publish(model.Activity{Type: "article.research.learned", Source: "brain", Phase: "embedding", NodeIDs: ids, Message: fmt.Sprintf("%d geprüfte Webbelege wurden mit lokalen Fallback-Vektoren gelernt", len(ids)), Strength: .48, Metadata: map[string]any{"result_count": len(ids), "fallback": true, "error": errorString(err)}})
return stats
}
for i, id := range ids {
stats.Add(e.Graph.SetVectorWithStats(id, vectors[i]))
}
e.Broker.Publish(model.Activity{Type: "article.research.learned", Source: "ollama", Phase: "embedding", NodeIDs: ids, Message: fmt.Sprintf("%d geprüfte Volltextbelege wurden eingebettet und sind semantisch nutzbar", len(ids)), Strength: .72, Metadata: map[string]any{"result_count": len(ids), "model": e.Cfg.EmbeddingModel, "dimensions": len(vectors[0])}})
return stats
}
func errorString(err error) string {
if err == nil {
return ""
}
return err.Error()
}