This commit is contained in:
@@ -105,7 +105,6 @@ func (e *Engine) synthesizeKnowledgeArticle(ctx context.Context, trigger string,
|
||||
return articleSynthesisOutcome{Skipped: true, Reason: "required_research_empty", Action: plan.Action}, nil
|
||||
}
|
||||
researchResults = append(researchResults, results...)
|
||||
e.addResearchToSources(selected, results)
|
||||
brief, err = e.buildKnowledgeBrief(ctx, selected, researchResults)
|
||||
if err != nil {
|
||||
return articleSynthesisOutcome{}, fmt.Errorf("knowledge consolidation after research failed: %w", err)
|
||||
@@ -460,16 +459,30 @@ func (e *Engine) researchKnowledgeGaps(ctx context.Context, trigger string, node
|
||||
if query == "" {
|
||||
continue
|
||||
}
|
||||
e.Broker.Publish(model.Activity{Type: "article.research.started", Source: "brain", Phase: "knowledge-research", NodeIDs: nodeIDs, Message: "Ein ungeklärter fachlicher Punkt wird recherchiert", Strength: .9, Metadata: map[string]any{"trigger": trigger, "research_query": query}})
|
||||
researchID := newResearchRunID("article-research", query)
|
||||
started := time.Now()
|
||||
startMetadata := map[string]any{"trigger": trigger, "research_id": researchID, "research_query": query, "animation_min_ms": 2000}
|
||||
e.Broker.Publish(model.Activity{Type: "article.research.started", Source: "searxng", Phase: "knowledge-research", NodeIDs: nodeIDs, Message: "Ein ungeklärter fachlicher Punkt wird mit SearXNG recherchiert", Strength: .9, Metadata: startMetadata})
|
||||
resultLimit := e.Cfg.ArticleResearchResults
|
||||
if resultLimit < 1 {
|
||||
resultLimit = 4
|
||||
}
|
||||
results, err := e.Research.Search(ctx, query, resultLimit)
|
||||
if err != nil {
|
||||
e.Broker.Publish(model.Activity{Type: "article.research.failed", Source: "brain", Phase: "knowledge-research", NodeIDs: nodeIDs, Message: "Die ergänzende Artikelrecherche ist fehlgeschlagen", Strength: .35, Metadata: map[string]any{"trigger": trigger, "research_query": query, "error": err.Error()}})
|
||||
e.Broker.Publish(model.Activity{Type: "article.research.failed", Source: "searxng", Phase: "knowledge-research", NodeIDs: nodeIDs, Message: "Die ergänzende Artikelrecherche ist fehlgeschlagen", Strength: .35, Metadata: mergeResearchMetadata(startMetadata, map[string]any{"error": err.Error(), "duration_ms": time.Since(started).Milliseconds()})})
|
||||
return nil, err
|
||||
}
|
||||
resultMetadata := researchEventMetadata(trigger, researchID, query, results, time.Since(started))
|
||||
message := fmt.Sprintf("SearXNG hat %d Quellen für den offenen Wissenspunkt geliefert", len(results))
|
||||
if len(results) == 0 {
|
||||
message = "SearXNG hat für den offenen Wissenspunkt keine verwertbare Quelle geliefert"
|
||||
}
|
||||
e.Broker.Publish(model.Activity{Type: "article.research.results", Source: "searxng", Phase: "knowledge-research-results", NodeIDs: nodeIDs, Message: message, Strength: .94, Metadata: resultMetadata})
|
||||
if len(results) > 0 {
|
||||
refs := e.addResearchToNodeIDs(nodeIDs, results)
|
||||
ingestMetadata := mergeResearchMetadata(resultMetadata, map[string]any{"result_node_ids": refs.NodeIDs, "result_edge_ids": refs.EdgeIDs, "source_node_ids": nodeIDs})
|
||||
e.Broker.Publish(model.Activity{Type: "article.research.ingested", Source: "searxng", Phase: "knowledge-research-ingest", NodeIDs: append(append([]string{}, nodeIDs...), refs.NodeIDs...), EdgeIDs: refs.EdgeIDs, Message: fmt.Sprintf("%d recherchierte Quellen wurden als neue Forschungs-Nodes verknüpft", len(refs.NodeIDs)), Strength: 1, Metadata: ingestMetadata})
|
||||
}
|
||||
for _, result := range results {
|
||||
key := strings.TrimSpace(result.URL)
|
||||
if key == "" {
|
||||
@@ -957,14 +970,34 @@ func (e *Engine) learnRuntimeArticle(ctx context.Context, articleID string) {
|
||||
e.Broker.Publish(model.Activity{Type: "article.learned", Source: "ollama", Phase: "embedding", NodeIDs: []string{nodeID}, Message: "Der neue KB-Artikel wurde eingebettet und ist sofort für Verknüpfungen verfügbar", Strength: .62, Metadata: map[string]any{"article_id": articleID, "model": e.Cfg.EmbeddingModel, "dimensions": len(vecs[0])}})
|
||||
}
|
||||
|
||||
func (e *Engine) addResearchToSources(sources []articleSource, results []model.ResearchResult) {
|
||||
func (e *Engine) addResearchToNodeIDs(nodeIDs []string, results []model.ResearchResult) researchGraphRefs {
|
||||
refs := researchGraphRefs{}
|
||||
for _, result := range results {
|
||||
id := graph.ID("external", result.URL)
|
||||
e.Graph.UpsertNode(model.Node{ID: id, Kind: "external", Label: result.Title, Summary: clamp(result.Content, 900), Status: "research", Origin: "research", ExternalID: result.URL, URI: result.URL, Weight: .8, UpdatedAt: time.Now().UTC()})
|
||||
for _, source := range sources {
|
||||
e.Graph.UpsertEdge(model.Edge{Source: id, Target: source.Node.ID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: .55, Weight: .4})
|
||||
refs.NodeIDs = append(refs.NodeIDs, id)
|
||||
for _, targetID := range nodeIDs {
|
||||
edge := model.Edge{Source: id, Target: targetID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: .55, Weight: .4}
|
||||
e.Graph.UpsertEdge(edge)
|
||||
refs.EdgeIDs = append(refs.EdgeIDs, graph.EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin))
|
||||
}
|
||||
}
|
||||
return uniqueResearchRefs(refs)
|
||||
}
|
||||
|
||||
func (e *Engine) addResearchToSources(sources []articleSource, results []model.ResearchResult) researchGraphRefs {
|
||||
refs := researchGraphRefs{}
|
||||
for _, result := range results {
|
||||
id := graph.ID("external", result.URL)
|
||||
e.Graph.UpsertNode(model.Node{ID: id, Kind: "external", Label: result.Title, Summary: clamp(result.Content, 900), Status: "research", Origin: "research", ExternalID: result.URL, URI: result.URL, Weight: .8, UpdatedAt: time.Now().UTC()})
|
||||
refs.NodeIDs = append(refs.NodeIDs, id)
|
||||
for _, source := range sources {
|
||||
edge := model.Edge{Source: id, Target: source.Node.ID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: .55, Weight: .4}
|
||||
e.Graph.UpsertEdge(edge)
|
||||
refs.EdgeIDs = append(refs.EdgeIDs, graph.EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin))
|
||||
}
|
||||
}
|
||||
return uniqueResearchRefs(refs)
|
||||
}
|
||||
|
||||
func formatArticleAnswer(draft model.KnowledgeArticleDraft) string {
|
||||
|
||||
+33
-13
@@ -627,19 +627,33 @@ func (e *Engine) enrichOne(ctx context.Context, trigger string) (EnrichOutcome,
|
||||
|
||||
var researchResults []model.ResearchResult
|
||||
if decision.NeedsResearch && e.Cfg.ResearchEnabled && e.Research != nil && strings.TrimSpace(decision.ResearchQuery) != "" {
|
||||
e.Broker.Publish(model.Activity{Type: "research.started", Source: "brain", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "Unklarheit erkannt · kontrollierte Webrecherche startet", Strength: .9, Metadata: map[string]any{"trigger": trigger, "research_query": decision.ResearchQuery, "source_label": a.Label, "target_label": b.Label}})
|
||||
researchID := newResearchRunID("relation-research", decision.ResearchQuery)
|
||||
researchStarted := time.Now()
|
||||
startMetadata := map[string]any{"trigger": trigger, "research_id": researchID, "research_query": decision.ResearchQuery, "source_label": a.Label, "target_label": b.Label, "animation_min_ms": 2000}
|
||||
e.Broker.Publish(model.Activity{Type: "research.started", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "Unklarheit erkannt · SearXNG durchsucht externe Quellen", Strength: .9, Metadata: startMetadata})
|
||||
results, err := e.Research.Search(ctx, decision.ResearchQuery, 4)
|
||||
if err != nil {
|
||||
slog.Warn("research failed", "error", err)
|
||||
} else if len(results) > 0 {
|
||||
researchResults = results
|
||||
e.addResearch(a, b, results)
|
||||
var reviewed model.RelationDecision
|
||||
reviewSystem := "Bewerte die Beziehung erneut anhand der zwei internen Wissenseinträge und der beigefügten Web-Suchergebnisse. Suchtreffer sind Hinweise, keine garantierten Fakten. Erfinde nichts, kennzeichne verbleibende Unsicherheit und gib ausschließlich JSON nach Schema zurück."
|
||||
if err := e.Ollama.ChatJSON(ctx, reviewSystem, relationContextWithResearch(a, b, sim, results), relationSchema(), &reviewed); err != nil {
|
||||
slog.Warn("research review failed; keeping pre-research decision", "error", err)
|
||||
} else {
|
||||
decision = reviewed
|
||||
e.Broker.Publish(model.Activity{Type: "research.failed", Source: "searxng", Phase: "research", NodeIDs: []string{a.ID, b.ID}, Message: "SearXNG-Recherche ist fehlgeschlagen", Strength: .35, Metadata: mergeResearchMetadata(startMetadata, map[string]any{"error": err.Error(), "duration_ms": time.Since(researchStarted).Milliseconds()})})
|
||||
} else {
|
||||
resultMetadata := researchEventMetadata(trigger, researchID, decision.ResearchQuery, results, time.Since(researchStarted))
|
||||
message := fmt.Sprintf("SearXNG hat %d verwertbare Webquellen geliefert", len(results))
|
||||
if len(results) == 0 {
|
||||
message = "SearXNG hat keine verwertbaren Webquellen geliefert"
|
||||
}
|
||||
e.Broker.Publish(model.Activity{Type: "research.results", Source: "searxng", Phase: "research-results", NodeIDs: []string{a.ID, b.ID}, Message: message, Strength: .92, Metadata: resultMetadata})
|
||||
if len(results) > 0 {
|
||||
researchResults = results
|
||||
refs := e.addResearch(a, b, results)
|
||||
ingestMetadata := mergeResearchMetadata(resultMetadata, map[string]any{"result_node_ids": refs.NodeIDs, "result_edge_ids": refs.EdgeIDs})
|
||||
e.Broker.Publish(model.Activity{Type: "research.ingested", Source: "searxng", Phase: "research-ingest", NodeIDs: append([]string{a.ID, b.ID}, refs.NodeIDs...), EdgeIDs: refs.EdgeIDs, Message: fmt.Sprintf("%d Webquellen wurden als neue Forschungs-Nodes in den Graphen übernommen", len(refs.NodeIDs)), Strength: 1, Metadata: ingestMetadata})
|
||||
var reviewed model.RelationDecision
|
||||
reviewSystem := "Bewerte die Beziehung erneut anhand der zwei internen Wissenseinträge und der beigefügten Web-Suchergebnisse. Suchtreffer sind Hinweise, keine garantierten Fakten. Erfinde nichts, kennzeichne verbleibende Unsicherheit und gib ausschließlich JSON nach Schema zurück."
|
||||
if err := e.Ollama.ChatJSON(ctx, reviewSystem, relationContextWithResearch(a, b, sim, results), relationSchema(), &reviewed); err != nil {
|
||||
slog.Warn("research review failed; keeping pre-research decision", "error", err)
|
||||
} else {
|
||||
decision = reviewed
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -676,14 +690,20 @@ func (e *Engine) enrichOne(ctx context.Context, trigger string) (EnrichOutcome,
|
||||
return outcome, nil
|
||||
}
|
||||
|
||||
func (e *Engine) addResearch(a, b model.Node, results []model.ResearchResult) {
|
||||
func (e *Engine) addResearch(a, b model.Node, results []model.ResearchResult) researchGraphRefs {
|
||||
refs := researchGraphRefs{}
|
||||
for _, r := range results {
|
||||
id := graph.ID("external", r.URL)
|
||||
n := model.Node{ID: id, Kind: "external", Label: r.Title, Summary: clamp(r.Content, 700), Status: "research", Origin: "research", ExternalID: r.URL, URI: r.URL, Weight: .8, Metadata: map[string]any{"query_pair": []string{a.ID, b.ID}}, UpdatedAt: time.Now().UTC()}
|
||||
e.Graph.UpsertNode(n)
|
||||
e.Graph.UpsertEdge(model.Edge{Source: id, Target: a.ID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: .55, Weight: .4})
|
||||
e.Graph.UpsertEdge(model.Edge{Source: id, Target: b.ID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: .55, Weight: .4})
|
||||
refs.NodeIDs = append(refs.NodeIDs, id)
|
||||
for _, targetID := range []string{a.ID, b.ID} {
|
||||
edge := model.Edge{Source: id, Target: targetID, Type: "research_evidence", Origin: "research", Status: "staging", Confidence: .55, Weight: .4}
|
||||
e.Graph.UpsertEdge(edge)
|
||||
refs.EdgeIDs = append(refs.EdgeIDs, graph.EdgeID(edge.Source, edge.Target, edge.Type, edge.Origin))
|
||||
}
|
||||
}
|
||||
return uniqueResearchRefs(refs)
|
||||
}
|
||||
func (e *Engine) Status() map[string]any {
|
||||
s := e.Graph.Snapshot()
|
||||
|
||||
@@ -246,7 +246,8 @@ func TestSynthesisResearchesUnclearKnowledgeThenLearnsAndLinksArticle(t *testing
|
||||
t.Fatal(err)
|
||||
}
|
||||
cfg := config.Config{DataDir: data, KnowledgeDirs: []string{knowledge}, StagingDirs: []string{staging}, OllamaURL: ollama.URL, ChatModel: "qwen3:8b", EmbeddingModel: "embeddinggemma", SearXNGURL: searx.URL, ResearchEnabled: true, SimilarityThreshold: .5, RelationThreshold: .7, ArticleSynthesisEnabled: true, ArticleMinSources: 3, ArticleMaxSources: 6, ArticleMinProductionRatio: .7, ArticleMaxGenerationDepth: 2, ArticleMinConfidence: .7, ArticleMinTextChars: 80, ArticleMinAnswerChars: 180, ArticleMaxResearchQueries: 3, ArticleResearchResults: 4, MaxContextChars: 16000}
|
||||
e := New(cfg, g, activity.New(100))
|
||||
broker := activity.New(100)
|
||||
e := New(cfg, g, broker)
|
||||
if err := e.Scan(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -287,6 +288,29 @@ func TestSynthesisResearchesUnclearKnowledgeThenLearnsAndLinksArticle(t *testing
|
||||
if chatCalls != 6 {
|
||||
t.Fatalf("expected relation, plan, two consolidations, article and quality calls, got %d", chatCalls)
|
||||
}
|
||||
var resultEvent, ingestEvent *model.Activity
|
||||
for _, event := range broker.Recent() {
|
||||
event := event
|
||||
switch event.Type {
|
||||
case "article.research.results":
|
||||
resultEvent = &event
|
||||
case "article.research.ingested":
|
||||
ingestEvent = &event
|
||||
}
|
||||
}
|
||||
if resultEvent == nil || ingestEvent == nil {
|
||||
t.Fatalf("expected visible SearXNG result and ingest events, results=%v ingest=%v", resultEvent != nil, ingestEvent != nil)
|
||||
}
|
||||
if count, _ := resultEvent.Metadata["result_count"].(int); count != 1 {
|
||||
t.Fatalf("unexpected result count metadata: %#v", resultEvent.Metadata["result_count"])
|
||||
}
|
||||
resultNodeIDs, _ := ingestEvent.Metadata["result_node_ids"].([]string)
|
||||
if len(resultNodeIDs) != 1 || resultNodeIDs[0] != researchID {
|
||||
t.Fatalf("unexpected ingested research nodes: %#v", ingestEvent.Metadata["result_node_ids"])
|
||||
}
|
||||
if ingestEvent.Metadata["animation_min_ms"] != 2000 {
|
||||
t.Fatalf("research animation minimum missing: %#v", ingestEvent.Metadata["animation_min_ms"])
|
||||
}
|
||||
}
|
||||
|
||||
func TestRequestEnrichDoesNotQueueDuplicateCycle(t *testing.T) {
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
package engine
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/url"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/local/glpi-neural-brain/internal/graph"
|
||||
"github.com/local/glpi-neural-brain/internal/model"
|
||||
)
|
||||
|
||||
type researchGraphRefs struct {
|
||||
NodeIDs []string
|
||||
EdgeIDs []string
|
||||
}
|
||||
|
||||
func newResearchRunID(prefix, query string) string {
|
||||
return graph.ID(prefix, strings.TrimSpace(query), fmt.Sprintf("%d", time.Now().UnixNano()))
|
||||
}
|
||||
|
||||
func researchEventMetadata(trigger, researchID, query string, results []model.ResearchResult, elapsed time.Duration) map[string]any {
|
||||
titles := make([]string, 0, min(5, len(results)))
|
||||
domains := make([]string, 0, min(5, len(results)))
|
||||
urls := make([]string, 0, min(5, len(results)))
|
||||
seenDomains := map[string]bool{}
|
||||
for _, result := range results {
|
||||
if title := strings.TrimSpace(result.Title); title != "" && len(titles) < 5 {
|
||||
titles = append(titles, title)
|
||||
}
|
||||
if rawURL := strings.TrimSpace(result.URL); rawURL != "" {
|
||||
if len(urls) < 5 {
|
||||
urls = append(urls, rawURL)
|
||||
}
|
||||
if parsed, err := url.Parse(rawURL); err == nil {
|
||||
domain := strings.TrimPrefix(strings.ToLower(parsed.Hostname()), "www.")
|
||||
if domain != "" && !seenDomains[domain] && len(domains) < 5 {
|
||||
seenDomains[domain] = true
|
||||
domains = append(domains, domain)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
sort.Strings(domains)
|
||||
return map[string]any{
|
||||
"trigger": trigger,
|
||||
"research_id": researchID,
|
||||
"research_query": query,
|
||||
"result_count": len(results),
|
||||
"result_titles": titles,
|
||||
"result_domains": domains,
|
||||
"result_urls": urls,
|
||||
"duration_ms": elapsed.Milliseconds(),
|
||||
"animation_min_ms": 2000,
|
||||
}
|
||||
}
|
||||
|
||||
func mergeResearchMetadata(base map[string]any, extra map[string]any) map[string]any {
|
||||
out := make(map[string]any, len(base)+len(extra))
|
||||
for key, value := range base {
|
||||
out[key] = value
|
||||
}
|
||||
for key, value := range extra {
|
||||
out[key] = value
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func uniqueResearchRefs(refs researchGraphRefs) researchGraphRefs {
|
||||
refs.NodeIDs = unique(refs.NodeIDs)
|
||||
refs.EdgeIDs = unique(refs.EdgeIDs)
|
||||
return refs
|
||||
}
|
||||
Reference in New Issue
Block a user