This commit is contained in:
@@ -0,0 +1,247 @@
|
||||
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{}
|
||||
for _, edge := range snapshot.Edges {
|
||||
if edge.Status == "rejected" {
|
||||
continue
|
||||
}
|
||||
switch edge.Type {
|
||||
case "research_evidence":
|
||||
if sourceIDs[edge.Target] {
|
||||
externalIDs[edge.Source] = true
|
||||
}
|
||||
case "grounded_by":
|
||||
if sourceIDs[edge.Source] {
|
||||
externalIDs[edge.Target] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(externalIDs) == 0 {
|
||||
return nil
|
||||
}
|
||||
nodes := make([]model.Node, 0, len(externalIDs))
|
||||
for _, node := range snapshot.Nodes {
|
||||
if externalIDs[node.ID] && node.Kind == "external" {
|
||||
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) {
|
||||
if !e.LearningEnabled() || len(results) == 0 || e.Ollama == nil {
|
||||
return
|
||||
}
|
||||
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 || !matchesCategories(node, e.learningCategories()) {
|
||||
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
|
||||
}
|
||||
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 {
|
||||
e.Graph.SetVector(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
|
||||
}
|
||||
for i, id := range ids {
|
||||
e.Graph.SetVector(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])}})
|
||||
}
|
||||
|
||||
func errorString(err error) string {
|
||||
if err == nil {
|
||||
return ""
|
||||
}
|
||||
return err.Error()
|
||||
}
|
||||
Reference in New Issue
Block a user