package engine import ( "context" "fmt" "log/slog" "math" "regexp" "strings" "time" "github.com/local/glpi-neural-brain/internal/graph" "github.com/local/glpi-neural-brain/internal/model" "github.com/local/glpi-neural-brain/internal/research" "github.com/local/glpi-neural-brain/internal/sourceagent" ) const SourceInboxClassifierVersion = 3 type sourceInboxClassification struct { Status string Priority float64 Semantic float64 Freshness float64 TaskContext float64 EventSignal float64 Security bool Reason string MatchedNodeID string } type securityInboxAssessment struct { Materialize bool `json:"materialize"` SecurityRelevant bool `json:"security_relevant"` Confidence float64 `json:"confidence"` EventType string `json:"event_type"` Severity string `json:"severity"` Vendor string `json:"vendor"` Products []string `json:"products"` CVEs []string `json:"cves"` AffectedVersions []string `json:"affected_versions"` FixedVersions []string `json:"fixed_versions"` Summary string `json:"summary"` Facts []string `json:"facts"` RecommendedActions []string `json:"recommended_actions"` ResearchNeeded bool `json:"research_needed"` ResearchQueries []string `json:"research_queries"` Reason string `json:"reason"` } func (e *Engine) SetSourceInbox(store *sourceagent.Store) { e.SourceInbox = store } func (e *Engine) evidenceAcquisitionEnabled() bool { return (e.SourceInbox != nil && e.Cfg.SourceInboxEnabled) || e.ResearchEnabledForRuntime() } func (e *Engine) sourceInboxLoop(ctx context.Context) { e.reconcileProactiveSecurityMaterialization(ctx) e.queueExistingProactiveSecurityCandidates(ctx) run := func() int { processed := e.processSourceInbox(ctx) e.queueExistingProactiveSecurityCandidates(ctx) processed += e.processProactiveSecurityInbox(ctx) return processed } drain := func() { for { processed := run() if !e.SpeedModeEnabled() || processed == 0 || ctx.Err() != nil { return } } } // Always classify once immediately. In Speed mode the loop keeps claiming // batches until the queues are empty instead of sleeping between batches. drain() ticker := time.NewTicker(e.Cfg.SourceInboxInterval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: drain() case <-e.sourceInboxWake: drain() } } } func (e *Engine) WakeSourceInbox() { if e == nil || e.sourceInboxWake == nil { return } select { case e.sourceInboxWake <- struct{}{}: default: } } func (e *Engine) processSourceInbox(ctx context.Context) int { if e.SourceInbox == nil || !e.Cfg.SourceInboxEnabled { return 0 } items, err := e.SourceInbox.ClaimInbox(ctx, e.Cfg.SourceInboxBatchSize) if err != nil { slog.Warn("source inbox claim failed", "error", err) return 0 } if len(items) == 0 { return 0 } texts := make([]string, 0, len(items)) for _, item := range items { text := item.Document.Title + "\n" + strings.Join(item.Document.Categories, " · ") + "\n" + item.Document.Text if len([]rune(text)) > 12000 { text = string([]rune(text)[:12000]) } texts = append(texts, text) } cctx, cancel := context.WithTimeout(e.backgroundOllamaContext(ctx), 4*time.Minute) vectors, embedErr := e.Ollama.Embed(cctx, texts) cancel() if embedErr != nil || len(vectors) != len(items) { reason := fmt.Sprint(embedErr) if embedErr == nil { reason = fmt.Sprintf("embedding response count mismatch: got %d vectors for %d inbox documents", len(vectors), len(items)) } for _, item := range items { if releaseErr := e.SourceInbox.ReleaseInbox(ctx, item.ID, reason); releaseErr != nil { slog.Error("source inbox release after embedding failure failed", "inbox_id", item.ID, "error", releaseErr) } } return 0 } candidateCount := 0 for i, item := range items { hits, stats := e.similarKnowledge(vectors[i], 8, e.effectiveLearningFilter(), 0) hit, found := e.sourceInboxKnowledgeHit(hits) if !found { if releaseErr := e.SourceInbox.ReleaseInbox(ctx, item.ID, "knowledge vectors are not ready"); releaseErr != nil { slog.Error("source inbox release without knowledge hit failed", "inbox_id", item.ID, "error", releaseErr) } continue } classification := e.classifySourceInboxItem(item, hit.Score, hit.NodeID, time.Now().UTC()) if classification.Status == "candidate" { candidateCount++ } meta := make(map[string]any, len(item.Metadata)+12) for key, value := range item.Metadata { meta[key] = value } meta["processing_mode"] = e.RuntimeSettings().ProcessingMode meta["exact_comparisons"] = stats.ExactComparisons meta["coarse_comparisons"] = stats.CoarseComparisons meta["candidate_pool"] = stats.CandidatePool meta["classification_version"] = SourceInboxClassifierVersion meta["classification_reason"] = classification.Reason meta["semantic_similarity"] = classification.Semantic meta["priority_score"] = classification.Priority meta["freshness_score"] = classification.Freshness meta["task_context_score"] = classification.TaskContext meta["event_signal_score"] = classification.EventSignal meta["security_candidate"] = classification.Security if completeErr := e.SourceInbox.CompleteClassification(ctx, item.ID, classification.Status, classification.Priority, classification.MatchedNodeID, meta); completeErr != nil { slog.Error("source inbox classification completion failed", "inbox_id", item.ID, "error", completeErr) if releaseErr := e.SourceInbox.ReleaseInbox(ctx, item.ID, completeErr.Error()); releaseErr != nil { slog.Error("source inbox classification recovery failed", "inbox_id", item.ID, "error", releaseErr) } continue } if classification.Security && e.Cfg.SourceInboxSecurityProactiveEnabled && classification.Priority >= e.Cfg.SourceInboxSecurityMinPriority { if queueErr := e.SourceInbox.QueueProactiveSecurity(ctx, item.ID); queueErr != nil { slog.Error("source inbox security queue failed", "inbox_id", item.ID, "error", queueErr) } } } if e.Broker != nil { e.Broker.Publish(model.Activity{Type: "source.inbox.classified", Source: "brain", Phase: "source-inbox", Message: fmt.Sprintf("Source-Inbox: %d Dokumente geprüft · %d als Wissenskandidaten vorgemerkt", len(items), candidateCount), Strength: .36, Metadata: map[string]any{"documents": len(items), "candidates": candidateCount, "minimum_similarity": e.Cfg.SourceInboxMinSimilarity, "minimum_priority": e.Cfg.SourceInboxMinPriority, "novelty_floor": e.Cfg.SourceInboxNoveltyFloor, "classifier_version": SourceInboxClassifierVersion}}) } return len(items) } func (e *Engine) sourceInboxKnowledgeHit(hits []model.Hit) (model.Hit, bool) { for _, hit := range hits { node, ok := e.Graph.GetNode(hit.NodeID) if !ok || (node.Kind != "knowledge" && node.Kind != "ai-think") { continue } return hit, true } return model.Hit{}, false } func (e *Engine) classifySourceInboxItem(item sourceagent.InboxDocument, semantic float64, matchedNodeID string, now time.Time) sourceInboxClassification { freshness := sourceInboxFreshness(item, now) taskContext := sourceInboxTaskContext(item) eventSignal := sourceInboxEventSignal(item) priority := clampInboxScore(0.65*semantic + 0.15*taskContext + 0.10*freshness + 0.10*eventSignal) status := "archived" reason := "zu geringe Nähe zur bestehenden Wissensbasis und keine ausreichend starken Aktualitäts-/Quellensignale" if semantic >= e.Cfg.SourceInboxMinSimilarity { status = "candidate" reason = "direkte semantische Nähe zur bestehenden Wissensbasis" } else if semantic >= e.Cfg.SourceInboxNoveltyFloor && priority >= e.Cfg.SourceInboxMinPriority { status = "candidate" reason = "neues, ausreichend KB-nahes Material mit zusätzlichem Quellen-, Aktualitäts- oder Advisory-Signal" } security := sourceInboxSecurityProactive(item, taskContext, eventSignal) return sourceInboxClassification{Status: status, Priority: priority, Semantic: semantic, Freshness: freshness, TaskContext: taskContext, EventSignal: eventSignal, Security: security, Reason: reason, MatchedNodeID: matchedNodeID} } func sourceInboxSecurityProactive(item sourceagent.InboxDocument, taskContext, eventSignal float64) bool { mode := strings.ToLower(strings.TrimSpace(sourceInboxMetadataString(item.Metadata, "security_proactive"))) switch mode { case "false", "0", "off", "passive", "disabled": return false case "true", "1", "on", "proactive", "enabled": return true } return taskContext >= .9 && eventSignal >= .75 } func sourceInboxMetadataString(meta map[string]any, key string) string { if meta == nil { return "" } value, ok := meta[key] if !ok || value == nil { return "" } return fmt.Sprint(value) } func sourceInboxFreshness(item sourceagent.InboxDocument, now time.Time) float64 { t := item.Document.PublishedAt if t.IsZero() { t = item.Document.DiscoveredAt } if t.IsZero() { t = item.ReceivedAt } if t.IsZero() { return .25 } age := now.Sub(t) if age < 0 { age = 0 } switch { case age <= 24*time.Hour: return 1 case age <= 7*24*time.Hour: return .85 case age <= 30*24*time.Hour: return .55 case age <= 90*24*time.Hour: return .25 default: return .10 } } func sourceInboxTaskContext(item sourceagent.InboxDocument) float64 { contextText := strings.ToLower(strings.Join(append(append([]string{}, item.Document.Categories...), item.Document.SourceName, item.Document.SourceBaseURL), " ")) for _, marker := range []string{"security", "sicherheit", "advisory", "alert", "vulnerability", "vulnerabil", "cve", "incident", "patch", "release", "threat", "bedroh", "exploit"} { if strings.Contains(contextText, marker) { return 1 } } if len(item.Document.Categories) > 0 { return .65 } if strings.TrimSpace(item.Document.SourceName) != "" { return .40 } return .20 } func sourceInboxEventSignal(item sourceagent.InboxDocument) float64 { text := strings.ToLower(item.Document.Title + " " + strings.Join(item.Document.Categories, " ") + " " + item.Document.SourceName) for _, marker := range []string{"0-day", "0day", "zero-day", "cve-", "kritisch", "critical", "actively exploited", "aktiv ausgenutzt", "backdoor", "authentifizierung umgehen", "authentication bypass", "remote code", "rce", "schadcode", "malware", "kompromitt"} { if strings.Contains(text, marker) { return 1 } } for _, marker := range []string{"angreifer", "attack", "sicherheitslücke", "sicherheitsleck", "vulnerability", "security update", "sicherheitsupdate", "patch", "update", "exploit", "datenleck", "security", "firewall", "advisory"} { if strings.Contains(text, marker) { return .75 } } return .20 } func clampInboxScore(v float64) float64 { return math.Max(0, math.Min(1, v)) } func (e *Engine) reconcileProactiveSecurityMaterialization(ctx context.Context) { if e.SourceInbox == nil || e.Graph == nil { return } lifecycles, err := e.SourceInbox.SecurityLifecycles(ctx, time.Unix(0, 0).UTC()) if err != nil { slog.Error("source security startup reconciliation failed", "error", err) return } requeued := 0 repairedVectors := 0 for _, lifecycle := range lifecycles { state := strings.ToLower(strings.TrimSpace(lifecycle.ProactiveState)) status := strings.ToLower(strings.TrimSpace(lifecycle.Status)) if state != "done" || (status != "materialized" && status != "used") { continue } node, exists := e.Graph.GetNode(lifecycle.MaterializedNodeID) if lifecycle.MaterializedNodeID == "" || !exists || node.Origin != "source-agent-security" { reason := "materialized Source-Inbox record has no persisted source-agent-security graph node after restart" if requeueErr := e.SourceInbox.RequeueMissingMaterializedSecurity(ctx, lifecycle.InboxID, reason); requeueErr != nil { slog.Error("source security startup requeue failed", "inbox_id", lifecycle.InboxID, "error", requeueErr) continue } requeued++ continue } if vector, ok := e.Graph.Vector(node.ID); !ok || len(vector) == 0 { stats := e.learnProactiveSecurityNode(ctx, node) if !stats.Empty() { repairedVectors++ } } } if (requeued > 0 || repairedVectors > 0) && e.Broker != nil { e.Broker.Publish(model.Activity{Type: "source.security.reconciled", Source: "brain", Phase: "source-inbox-security", Message: fmt.Sprintf("Security-Startup-Reconciliation: %d fehlende Materialisierungen neu eingeplant · %d fehlende Vektoren repariert", requeued, repairedVectors), Strength: .72, Metadata: map[string]any{"requeued": requeued, "vectors_repaired": repairedVectors}}) } } func (e *Engine) queueExistingProactiveSecurityCandidates(ctx context.Context) { if e.SourceInbox == nil || !e.Cfg.SourceInboxEnabled || !e.Cfg.SourceInboxSecurityProactiveEnabled { return } items, err := e.SourceInbox.ListInbox(ctx, "candidate", 1000) if err != nil { return } for _, item := range items { priority := item.Relevance if v, ok := item.Metadata["priority_score"].(float64); ok { priority = v } taskContext := sourceInboxTaskContext(item) eventSignal := sourceInboxEventSignal(item) if item.ProactiveState == "" && priority >= e.Cfg.SourceInboxSecurityMinPriority && sourceInboxSecurityProactive(item, taskContext, eventSignal) { if queueErr := e.SourceInbox.QueueProactiveSecurity(ctx, item.ID); queueErr != nil { slog.Error("existing source inbox security queue failed", "inbox_id", item.ID, "error", queueErr) } } } } func (e *Engine) processProactiveSecurityInbox(ctx context.Context) int { if e.SourceInbox == nil || !e.Cfg.SourceInboxEnabled || !e.Cfg.SourceInboxSecurityProactiveEnabled || e.Ollama == nil { return 0 } items, err := e.SourceInbox.ClaimProactiveSecurity(ctx, e.Cfg.SourceInboxSecurityBatchSize) if err != nil { slog.Warn("source inbox proactive security claim failed", "error", err) return 0 } if len(items) == 0 { return 0 } materialized := 0 for _, item := range items { if ctx.Err() != nil { bg := context.Background() runID, started, startErr := e.SourceInbox.StartProactiveSecurityRun(bg, item.ID) if startErr == nil { if releaseErr := e.SourceInbox.ReleaseProactiveSecurity(bg, item.ID, ctx.Err().Error()); releaseErr != nil { slog.Error("source security cancellation release failed", "inbox_id", item.ID, "error", releaseErr) } if e.Broker != nil { meta := withRunMutations(map[string]any{"run_id": runID, "inbox_id": item.ID, "title": item.Document.Title, "error": ctx.Err().Error(), "duration_ms": time.Since(started).Milliseconds(), "result": "cancelled"}, graph.MutationStats{}) e.Broker.Publish(model.Activity{Type: "source.security.failed", Source: "brain", Phase: "source-inbox-security", Message: "Proaktive Security-Auswertung wurde vor dem Start abgebrochen", Strength: .22, Metadata: meta}) } } else { if releaseErr := e.SourceInbox.ReleaseProactiveSecurity(bg, item.ID, ctx.Err().Error()); releaseErr != nil { slog.Error("source security cancellation recovery failed", "inbox_id", item.ID, "start_error", startErr, "release_error", releaseErr) } } continue } runID, started, startErr := e.SourceInbox.StartProactiveSecurityRun(ctx, item.ID) if startErr != nil { if releaseErr := e.SourceInbox.ReleaseProactiveSecurity(ctx, item.ID, startErr.Error()); releaseErr != nil { slog.Error("source security start recovery failed", "inbox_id", item.ID, "start_error", startErr, "release_error", releaseErr) } continue } if item.Metadata == nil { item.Metadata = map[string]any{} } item.Metadata["proactive_run_id"] = runID item.Metadata["proactive_started_at"] = started if e.Broker != nil { e.Broker.Publish(model.Activity{Type: "source.security.started", Source: "brain", Phase: "source-inbox-security", Message: "Priorisierte Security-Meldung wird gegen die vorhandene KB geprüft", Strength: .58, Metadata: map[string]any{"run_id": runID, "inbox_id": item.ID, "title": item.Document.Title, "agent_id": item.AgentID, "task_id": item.TaskID, "priority": item.Relevance, "matched_node_id": item.MatchedNodeID}}) } ok, mutations, processErr := e.processProactiveSecurityItem(ctx, item, started) if processErr != nil { if releaseErr := e.SourceInbox.ReleaseProactiveSecurity(ctx, item.ID, processErr.Error()); releaseErr != nil { slog.Error("source security retry release failed", "inbox_id", item.ID, "process_error", processErr, "release_error", releaseErr) } if e.Broker != nil { meta := withRunMutations(map[string]any{"run_id": runID, "inbox_id": item.ID, "title": item.Document.Title, "error": processErr.Error(), "duration_ms": time.Since(started).Milliseconds()}, mutations) e.Broker.Publish(model.Activity{Type: "source.security.failed", Source: "brain", Phase: "source-inbox-security", Message: "Proaktive Security-Auswertung wurde vertagt", Strength: .28, Metadata: meta}) } continue } if ok { materialized++ } } if materialized > 0 { e.RequestEnrich("source-security") } return len(items) } func (e *Engine) processProactiveSecurityItem(ctx context.Context, item sourceagent.InboxDocument, started time.Time) (bool, graph.MutationStats, error) { var mutations graph.MutationStats primary, directFetch := e.securityInboxPrimaryEvidence(ctx, item) matched, _ := e.Graph.GetNode(item.MatchedNodeID) assessment, err := e.assessSecurityInbox(ctx, item, matched, primary, nil) if err != nil { return false, mutations, err } supplemental := []model.ResearchResult{} if assessment.ResearchNeeded && e.ResearchEnabledForRuntime() && e.Research != nil { supplemental = e.securityInboxSupplementalResearch(ctx, item, assessment) if len(supplemental) > 0 { assessment, err = e.assessSecurityInbox(ctx, item, matched, primary, supplemental) if err != nil { return false, mutations, err } } } assessment.Confidence = clamp01(assessment.Confidence) applicability, applicabilityReason := securityInboxDirectApplicability(assessment, matched) if !assessment.Materialize || !assessment.SecurityRelevant || assessment.Confidence < e.Cfg.SourceInboxSecurityMinConfidence || strings.TrimSpace(assessment.Summary) == "" || len(assessment.Facts) == 0 { meta := securityAssessmentMetadata(assessment) meta["direct_fetch_attempted"] = directFetch meta["supplemental_sources"] = researchURLs(supplemental) meta["proactive_run_id"] = sourceInboxMetadataString(item.Metadata, "proactive_run_id") meta = withRunMutations(meta, mutations) if rejectErr := e.SourceInbox.RejectProactiveSecurity(ctx, item.ID, nonempty(assessment.Reason, "kein ausreichend belastbarer Security-Faktensatz"), meta); rejectErr != nil { return false, mutations, rejectErr } if e.Broker != nil { terminalMeta := mergeResearchMetadata(meta, map[string]any{"run_id": sourceInboxMetadataString(item.Metadata, "proactive_run_id"), "inbox_id": item.ID, "title": item.Document.Title, "duration_ms": time.Since(started).Milliseconds(), "direct_fetch_attempted": directFetch, "supplemental_sources": len(supplemental)}) e.Broker.Publish(model.Activity{Type: "source.security.rejected", Source: "brain", Phase: "source-inbox-security", Message: "Security-Candidate bleibt passiv in der Inbox; Gemma sah noch keinen belastbaren Graph-Faktensatz", Strength: .34, Metadata: terminalMeta}) } return false, mutations, nil } allEvidence := append([]model.ResearchResult{primary}, supplemental...) evidencePaths := e.persistResearchMaterial(allEvidence) nodeID := graph.ID("external", primary.URL) now := time.Now().UTC() meta := securityAssessmentMetadata(assessment) meta["security_applicability"] = applicability meta["security_applicability_reason"] = applicabilityReason meta["validation_state"] = "security_verified" meta["proactive_security"] = true meta["source_inbox_id"] = item.ID meta["source_agent_id"] = item.AgentID meta["source_task_id"] = item.TaskID meta["source"] = graph.SourceFromURL(primary.URL) meta["priority_score"] = item.Relevance meta["matched_node_id"] = item.MatchedNodeID meta["published_at"] = item.Document.PublishedAt meta["evidence_paths"] = evidencePaths meta["supplemental_sources"] = researchURLs(supplemental) meta["direct_fetch_attempted"] = directFetch meta["verified_by_model"] = e.Cfg.ArticleSynthesisModel meta["content_sha256"] = item.Document.ContentSHA256 summary := strings.TrimSpace(assessment.Summary) if len(assessment.Facts) > 0 { summary += "\n\nFakten:\n- " + strings.Join(unique(assessment.Facts), "\n- ") } if len(assessment.RecommendedActions) > 0 { summary += "\n\nEmpfohlene Maßnahmen:\n- " + strings.Join(unique(assessment.RecommendedActions), "\n- ") } categories := unique(append(append([]string{}, item.Document.Categories...), "Security Update", "Source Agent")) node := model.Node{ID: nodeID, Kind: "external", Label: nonempty(item.Document.Title, primary.Title), Summary: clamp(summary, 6000), Status: "research", Origin: "source-agent-security", ExternalID: primary.URL, URI: primary.URL, Categories: categories, Keywords: unique(append(append([]string{}, assessment.Products...), assessment.CVEs...)), Weight: 1.15 + item.Relevance*.35 + assessment.Confidence*.25, Metadata: meta, UpdatedAt: now} mutations.Add(e.Graph.UpsertNodeWithStats(node)) if item.MatchedNodeID != "" { if _, ok := e.Graph.GetNode(item.MatchedNodeID); ok { edgeType := "security_context_for" edgeConfidence := math.Min(assessment.Confidence, math.Max(.35, item.Relevance*.80)) edgeWeight := math.Max(.35, item.Relevance*.55) if applicability == "direct" { edgeType = "security_update_for" edgeConfidence = assessment.Confidence edgeWeight = math.Max(.65, item.Relevance) } edge := model.Edge{Source: nodeID, Target: item.MatchedNodeID, Type: edgeType, Origin: "source-agent-security", Status: "staging", Confidence: edgeConfidence, Weight: edgeWeight, Explanation: applicabilityReason, Metadata: map[string]any{"event_type": assessment.EventType, "severity": assessment.Severity, "cves": unique(assessment.CVEs), "products": unique(assessment.Products), "source_inbox_id": item.ID, "applicability": applicability}} mutations.Add(e.Graph.UpsertEdgeWithStats(edge)) } } mutations.Add(e.learnProactiveSecurityNode(ctx, node)) meta["proactive_run_id"] = sourceInboxMetadataString(item.Metadata, "proactive_run_id") meta = withRunMutations(meta, mutations) if err := e.SourceInbox.CompleteProactiveSecurity(ctx, item.ID, nodeID, meta); err != nil { return false, mutations, err } if e.Broker != nil { terminalMeta := withRunMutations(map[string]any{"run_id": sourceInboxMetadataString(item.Metadata, "proactive_run_id"), "inbox_id": item.ID, "node_id": nodeID, "title": node.Label, "confidence": assessment.Confidence, "severity": assessment.Severity, "event_type": assessment.EventType, "cves": unique(assessment.CVEs), "products": unique(assessment.Products), "matched_node_id": item.MatchedNodeID, "supplemental_sources": len(supplemental), "direct_fetch_attempted": directFetch, "research_needed": assessment.ResearchNeeded, "security_applicability": applicability, "duration_ms": time.Since(started).Milliseconds()}, mutations) e.Broker.Publish(model.Activity{Type: "source.security.materialized", Source: "brain", Phase: "source-inbox-security", NodeIDs: []string{nodeID}, Message: "Priorisierte Security-Meldung wurde als verifizierter externer Graph-Node materialisiert", Strength: .82, Metadata: terminalMeta}) } return true, mutations, nil } func securityInboxDirectApplicability(assessment securityInboxAssessment, matched model.Node) (string, string) { if matched.ID == "" { return "standalone", "Security-Meldung wird als eigenständige Evidenz materialisiert; es existiert kein direktes KB-Ziel." } target := strings.ToLower(strings.Join(append([]string{matched.Label}, matched.Keywords...), " ")) targetTokens := securityApplicabilityTokens(target) for _, cve := range assessment.CVEs { cve = strings.ToLower(strings.TrimSpace(cve)) if cve != "" && strings.Contains(target, cve) { return "direct", "Direkte Applicability: dieselbe CVE ist im KB-Ziel explizit referenziert." } } for _, product := range assessment.Products { normalized := strings.TrimSpace(strings.ToLower(product)) if normalized == "" { continue } if len([]rune(normalized)) >= 4 && strings.Contains(target, normalized) { return "direct", fmt.Sprintf("Direkte Applicability: Produkt %q ist im Titel/Keyword-Kontext des KB-Ziels enthalten.", product) } productTokens := securityApplicabilityTokens(normalized) if len(productTokens) == 1 { for token := range productTokens { if len(token) >= 5 && targetTokens[token] { return "direct", fmt.Sprintf("Direkte Applicability: Produktbegriff %q stimmt mit dem KB-Ziel überein.", product) } } } if len(productTokens) >= 2 && tokenJaccard(productTokens, targetTokens) >= .75 { return "direct", fmt.Sprintf("Direkte Applicability: Produktbegriffe %q stimmen weitgehend mit dem KB-Ziel überein.", product) } } return "contextual", "Nur kontextuelle Security-Nähe: Produkt/CVE stimmt nicht direkt mit Titel oder Keywords des KB-Ziels überein; daher keine security_update_for-Relation." } func securityApplicabilityTokens(value string) map[string]bool { value = strings.ToLower(value) parts := regexp.MustCompile(`[^a-z0-9äöüß._+-]+`).Split(value, -1) stop := map[string]bool{"security": true, "sicherheit": true, "update": true, "advisory": true, "multiple": true, "mehrere": true, "the": true, "und": true, "for": true, "für": true, "unter": true} out := map[string]bool{} for _, part := range parts { part = strings.TrimSpace(part) if len(part) < 3 || stop[part] { continue } out[part] = true } return out } func tokenJaccard(a, b map[string]bool) float64 { if len(a) == 0 || len(b) == 0 { return 0 } intersection := 0 union := map[string]bool{} for value := range a { union[value] = true if b[value] { intersection++ } } for value := range b { union[value] = true } return float64(intersection) / float64(len(union)) } func (e *Engine) securityInboxPrimaryEvidence(ctx context.Context, item sourceagent.InboxDocument) (model.ResearchResult, bool) { d := item.Document result := model.ResearchResult{Title: d.Title, URL: nonempty(d.CanonicalURL, d.URL), Snippet: clamp(d.Text, 1200), Content: d.Text, ContentType: d.ContentType, Language: d.Language, Fetched: true, Relevant: true, Relevance: item.Relevance, SourceQuality: "curated_agent", SourceQualityScore: .82, AssessmentReason: "Vom konfigurierten Source-Agent als priorisierte Security-Meldung geliefert."} canonicalURL := result.URL textShort := len([]rune(strings.TrimSpace(d.Text))) < e.Cfg.SourceInboxSecurityFetchMinChars fetchError := strings.TrimSpace(sourceInboxMetadataString(item.Metadata, "fetch_error_kind")) != "" if !textShort && !fetchError { return result, false } client := e.Research if client == nil { client = research.New("") } var page research.FetchedPage err := e.withSharedResearchWork(ctx, "web.fetch", func() error { var fetchErr error page, _, fetchErr = client.FetchPage(ctx, result.URL, research.FetchOptions{MaxBytes: e.Cfg.ArticleResearchPageMaxBytes, MaxChars: maxInt(e.Cfg.ArticleResearchPageMaxChars, e.Cfg.SourceInboxSecurityFetchMinChars*2), Timeout: e.Cfg.ArticleResearchFetchTimeout, AllowPrivate: e.Cfg.ArticleResearchAllowPrivate}) return fetchErr }) if err == nil && len([]rune(strings.TrimSpace(page.Content))) > len([]rune(strings.TrimSpace(result.Content))) { result.Content = page.Content result.Snippet = clamp(page.Content, 1200) result.ContentType = page.ContentType // Keep the Agent's canonical URL stable so later Claim-Grounding can mark // the same inbox record as used even when the publisher redirects. result.URL = canonicalURL result.Title = nonempty(page.Title, result.Title) } return result, true } func (e *Engine) assessSecurityInbox(ctx context.Context, item sourceagent.InboxDocument, matched model.Node, primary model.ResearchResult, supplemental []model.ResearchResult) (securityInboxAssessment, error) { var b strings.Builder fmt.Fprintf(&b, "SOURCE-INBOX-ID: %s\nAGENT: %s\nTASK: %s\nPRIORITÄT: %.3f\n\n", item.ID, item.AgentID, item.TaskID, item.Relevance) if matched.ID != "" { fmt.Fprintf(&b, "PASSENDSTER KB-KONTEXT:\nID: %s\nTitel: %s\nKategorien: %s\nInhalt: %s\n\n", matched.ID, matched.Label, strings.Join(matched.Categories, ", "), clamp(matched.Summary, 3000)) } fmt.Fprintf(&b, "P1 PRIMÄRQUELLE:\nTitel: %s\nURL: %s\nVeröffentlicht: %s\nInhalt:\n%s\n", primary.Title, primary.URL, item.Document.PublishedAt.UTC().Format(time.RFC3339), clamp(primary.Content, 12000)) for i, result := range supplemental { fmt.Fprintf(&b, "\nS%d ERGÄNZUNGSQUELLE:\nTitel: %s\nURL: %s\nInhalt:\n%s\n", i+1, result.Title, result.URL, clamp(result.Content, 7000)) } var out securityInboxAssessment modelName := strings.TrimSpace(e.Cfg.ArticleSynthesisModel) if modelName == "" { modelName = e.Cfg.ChatModel } cctx, cancel := context.WithTimeout(e.backgroundOllamaContext(ctx), 6*time.Minute) err := e.Ollama.ChatJSONModel(cctx, modelName, securityInboxSystemPrompt(), b.String(), securityInboxSchema(), &out) cancel() out = normalizeSecurityAssessment(out, item.Document.Title) return out, err } var cvePattern = regexp.MustCompile(`(?i)^CVE-[0-9]{4}-[0-9]{4,}$`) func normalizeSecurityAssessment(out securityInboxAssessment, sourceTitle string) securityInboxAssessment { out.Products = unique(out.Products) validCVEs := make([]string, 0, len(out.CVEs)) for _, cve := range unique(out.CVEs) { cve = strings.ToUpper(strings.TrimSpace(cve)) if cvePattern.MatchString(cve) { validCVEs = append(validCVEs, cve) } } out.CVEs = validCVEs out.AffectedVersions = unique(out.AffectedVersions) out.FixedVersions = unique(out.FixedVersions) out.Facts = unique(out.Facts) out.RecommendedActions = unique(out.RecommendedActions) out.ResearchQueries = unique(out.ResearchQueries) out.Severity = normalizeSecuritySeverity(out.Severity, sourceTitle) out.EventType = normalizeSecurityEventType(out.EventType) return out } func normalizeSecuritySeverity(value, sourceTitle string) string { v := strings.ToLower(strings.TrimSpace(value)) switch v { case "critical", "kritisch", "sehr hoch", "very high": return "critical" case "high", "hoch": return "high" case "medium", "mittel", "moderate", "moderat": return "medium" case "low", "niedrig": return "low" case "informational", "information", "info": return "informational" } title := strings.ToLower(sourceTitle) for _, candidate := range []struct { markers []string value string }{ {[]string{"[kritisch]", "[critical]"}, "critical"}, {[]string{"[hoch]", "[high]"}, "high"}, {[]string{"[mittel]", "[medium]"}, "medium"}, {[]string{"[niedrig]", "[low]"}, "low"}, } { for _, marker := range candidate.markers { if strings.Contains(title, marker) { return candidate.value } } } return "unknown" } func normalizeSecurityEventType(value string) string { v := strings.ToLower(strings.TrimSpace(value)) v = strings.NewReplacer("-", "_", " ", "_", "/", "_").Replace(v) for strings.Contains(v, "__") { v = strings.ReplaceAll(v, "__", "_") } switch v { case "vulnerability", "security_vulnerability", "schwachstelle", "vulnerabilities", "multiple_vulnerabilities": return "vulnerability" case "security_advisory", "advisory": return "security_advisory" case "security_update", "update": return "security_update" case "denial_of_service", "dos": return "denial_of_service" case "cross_site_scripting", "xss": return "xss" case "privilege_escalation", "privilegieneskalation": return "privilege_escalation" case "remote_code_execution", "code_execution", "rce", "codeausführung", "codeausfuehrung": return "code_execution" case "authentication_bypass", "auth_bypass": return "authentication_bypass" case "information_disclosure", "information_leak": return "information_disclosure" } if v == "" || v == "unknown" { return "vulnerability" } return v } func securityInboxSystemPrompt() string { return `Du bist der Security-Intelligence-Extractor des Neural Brain. Prüfe eine vom Betreiber kuratierte Security-/Advisory-Quelle gegen den angegebenen KB-Kontext. WICHTIG: - Der Source-Agent ist ein bestätigter Discovery-Kanal, aber du darfst AUSSCHLIESSLICH Fakten aus P1 und den S-Quellen übernehmen. - Nutze kein parametrisches Modellwissen für CVEs, Versionen, Schweregrade, Produkte oder Maßnahmen. - materialize=true nur wenn die Quellen eine konkrete neue oder aktualisierte Security-Information enthalten, die zum KB-Kontext passt oder diesen sinnvoll erweitert. - facts müssen konkrete, quellenbelegte Aussagen sein. Keine allgemeinen Security-Floskeln. - severity nur setzen, wenn sie aus den Quellen hervorgeht; sonst "unknown". - CVE-IDs, betroffene/fixe Versionen und Maßnahmen nur übernehmen, wenn explizit belegt. - research_needed=true, wenn die Meldung erkennbar relevant ist, aber für einen belastbaren Faktensatz wichtige Angaben fehlen. Formuliere dann höchstens zwei präzise Suchanfragen. - Wenn die Meldung nur allgemeine Meinung/Marketing ohne konkretes Security-Ereignis ist, materialize=false. Ziel ist ein kleiner, belastbarer Security-Faktensatz für einen externen Graph-Node; noch kein KB-Artikel.` } func securityInboxSchema() map[string]any { return map[string]any{ "type": "object", "properties": map[string]any{ "materialize": map[string]any{"type": "boolean"}, "security_relevant": map[string]any{"type": "boolean"}, "confidence": map[string]any{"type": "number"}, "event_type": map[string]any{"type": "string"}, "severity": map[string]any{"type": "string"}, "vendor": map[string]any{"type": "string"}, "products": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "cves": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "affected_versions": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "fixed_versions": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "summary": map[string]any{"type": "string"}, "facts": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "recommended_actions": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "research_needed": map[string]any{"type": "boolean"}, "research_queries": map[string]any{"type": "array", "items": map[string]any{"type": "string"}}, "reason": map[string]any{"type": "string"}, }, "required": []string{"materialize", "security_relevant", "confidence", "event_type", "severity", "vendor", "products", "cves", "affected_versions", "fixed_versions", "summary", "facts", "recommended_actions", "research_needed", "research_queries", "reason"}, } } func (e *Engine) securityInboxSupplementalResearch(ctx context.Context, item sourceagent.InboxDocument, assessment securityInboxAssessment) []model.ResearchResult { if e.Research == nil || !e.ResearchEnabledForRuntime() { return nil } queries := append([]string{}, assessment.ResearchQueries...) if len(queries) == 0 { queries = []string{strings.TrimSpace(item.Document.Title + " " + assessment.Vendor + " security advisory")} } if len(queries) > 2 { queries = queries[:2] } seen := map[string]bool{canonicalResearchURL(nonempty(item.Document.CanonicalURL, item.Document.URL)): true} out := []model.ResearchResult{} for _, query := range queries { if strings.TrimSpace(query) == "" || len(out) >= e.Cfg.SourceInboxSecurityResearchResults { continue } var results []model.ResearchResult err := e.withSharedResearchWork(ctx, "searxng.search", func() error { var searchErr error results, _, searchErr = e.Research.SearchDetailedLanguage(ctx, sanitizeSearchQuerySiteFilters(query), e.Cfg.SourceInboxSecurityResearchResults, nonempty(item.Document.Language, e.Cfg.ArticleLanguage)) return searchErr }) if err != nil { continue } for _, result := range results { if len(out) >= e.Cfg.SourceInboxSecurityResearchResults { break } key := canonicalResearchURL(result.URL) if key == "" || seen[key] { continue } seen[key] = true var page research.FetchedPage fetchErr := e.withSharedResearchWork(ctx, "web.fetch", func() error { var inner error page, _, inner = e.Research.FetchPage(ctx, result.URL, research.FetchOptions{MaxBytes: e.Cfg.ArticleResearchPageMaxBytes, MaxChars: e.Cfg.ArticleResearchPageMaxChars, Timeout: e.Cfg.ArticleResearchFetchTimeout, AllowPrivate: e.Cfg.ArticleResearchAllowPrivate}) return inner }) if fetchErr != nil || strings.TrimSpace(page.Content) == "" { continue } result.URL = nonempty(page.URL, result.URL) result.Title = nonempty(page.Title, result.Title) result.Content = page.Content result.Snippet = clamp(page.Content, 1200) result.ContentType = page.ContentType result.Fetched = true result.Relevant = true result.Relevance = .65 result.SourceQuality = "supplemental_security" result.SourceQualityScore = .72 result.AssessmentReason = "Gezielt nachgeladen, weil der Security-Candidate für einen belastbaren Faktensatz zusätzliche Details benötigte." out = append(out, result) } } if len(out) > 0 && e.Broker != nil { e.Broker.Publish(model.Activity{Type: "source.security.research", Source: "searxng", Phase: "source-inbox-security", Message: fmt.Sprintf("Security-Candidate wurde mit %d gezielt nachgeladenen Volltextquellen ergänzt", len(out)), Strength: .62, Metadata: map[string]any{"run_id": sourceInboxMetadataString(item.Metadata, "proactive_run_id"), "inbox_id": item.ID, "title": item.Document.Title, "results": len(out), "queries": queries}}) } return out } func securityAssessmentMetadata(a securityInboxAssessment) map[string]any { return map[string]any{"security_relevant": a.SecurityRelevant, "security_confidence": a.Confidence, "security_event_type": a.EventType, "security_severity": a.Severity, "security_vendor": a.Vendor, "security_products": unique(a.Products), "security_cves": unique(a.CVEs), "security_affected_versions": unique(a.AffectedVersions), "security_fixed_versions": unique(a.FixedVersions), "security_summary": a.Summary, "security_facts": unique(a.Facts), "security_recommended_actions": unique(a.RecommendedActions), "security_research_needed": a.ResearchNeeded, "security_research_queries": unique(a.ResearchQueries), "security_assessment_reason": a.Reason} } func researchURLs(results []model.ResearchResult) []string { values := make([]string, 0, len(results)) for _, result := range results { if strings.TrimSpace(result.URL) != "" { values = append(values, result.URL) } } return unique(values) } func (e *Engine) learnProactiveSecurityNode(ctx context.Context, node model.Node) graph.MutationStats { if e.Ollama == nil || strings.TrimSpace(node.ID) == "" { return graph.MutationStats{} } text := strings.TrimSpace(node.Label + "\n" + clamp(node.Summary, 6000)) if text == "" { return graph.MutationStats{} } cctx, cancel := context.WithTimeout(e.backgroundOllamaContext(ctx), 3*time.Minute) vectors, err := e.Ollama.Embed(cctx, []string{text}) cancel() if err != nil || len(vectors) != 1 || len(vectors[0]) == 0 { return e.Graph.SetVectorWithStats(node.ID, hashEmbedding(text, 256)) } return e.Graph.SetVectorWithStats(node.ID, vectors[0]) } func (e *Engine) sourceInboxResearch(ctx context.Context, query string, limit int, freshnessSensitive bool) []model.ResearchResult { if e.SourceInbox == nil || !e.Cfg.SourceInboxEnabled || limit < 1 { return nil } maxAge := time.Duration(0) if freshnessSensitive { maxAge = e.Cfg.SourceInboxFreshMaxAge } items, err := e.SourceInbox.SearchCandidates(ctx, query, limit, maxAge) if err != nil { return nil } out := make([]model.ResearchResult, 0, len(items)) for _, item := range items { if item.QueryScore < .18 { continue } d := item.Document out = append(out, model.ResearchResult{SourceInboxID: item.ID, Title: d.Title, URL: d.CanonicalURL, Snippet: clamp(d.Text, 1000), Content: d.Text, ContentType: d.ContentType, Query: query, Language: d.Language, Fetched: true, Relevant: true, Relevance: item.QueryScore, SourceQuality: "source_inbox", SourceQualityScore: .68, AssessmentReason: "Vorab durch Source-Agent gesammelt und vom Brain als thematisch passend zur Knowledgebase klassifiziert."}) } if len(out) > 0 && e.Broker != nil { e.Broker.Publish(model.Activity{Type: "article.research.inbox", Source: "brain", Phase: "knowledge-research-routing", Message: fmt.Sprintf("Source-Inbox liefert %d bereits gecrawlte Kandidaten vor SearXNG", len(out)), Strength: .56, Metadata: map[string]any{"query": query, "results": len(out), "freshness_sensitive": freshnessSensitive}}) } return out }