All checks were successful
release-tag / release-image (push) Successful in 2m43s
262 lines
9.3 KiB
Go
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()
|
|
}
|