package graph import ( "context" "sort" "testing" "time" "github.com/local/glpi-neural-brain/internal/activity" "github.com/local/glpi-neural-brain/internal/model" ) func TestMutationStatsAndAnalysisHistory(t *testing.T) { store, err := Open(t.TempDir()) if err != nil { t.Fatal(err) } t.Cleanup(func() { _ = store.Close() }) broker := activity.New(20) broker.SetSink(store.RecordActivity) store.UpsertNode(model.Node{ID: "a", Kind: "knowledge", Label: "A", Origin: "test", Metadata: map[string]any{"source": "internal"}}) store.UpsertNode(model.Node{ID: "b", Kind: "knowledge", Label: "B", Origin: "test", Metadata: map[string]any{"source": "internal"}}) store.SetVector("a", []float64{1, 0}) store.SetVector("b", []float64{.9, .1}) broker.Publish(model.Activity{Type: "learning.scan.completed", Source: "brain", Message: "done", Metadata: map[string]any{"run_id": "scan-1", "result": "updated"}}) deadline := time.Now().Add(3 * time.Second) for { history, err := store.AnalysisHistory(context.Background(), time.Now().Add(-time.Hour), 50) if err != nil { t.Fatal(err) } if len(history.Events) > 0 { if history.Totals.NodesCreated != 2 || history.Totals.VectorsCreated != 2 { t.Fatalf("unexpected mutation totals: %+v", history.Totals) } if history.ChangeCount != 4 || len(history.Changes) != 4 { t.Fatalf("expected four detailed changes, count=%d changes=%+v", history.ChangeCount, history.Changes) } if history.Events[0].ChangeCount != 4 { t.Fatalf("expected event to own four changes: %+v", history.Events[0]) } break } if time.Now().After(deadline) { t.Fatal("analysis writer did not persist event") } time.Sleep(20 * time.Millisecond) } } func TestBuildAnalysisRunsExplainsPositiveThinkingResult(t *testing.T) { start := time.Now().UTC().Add(-2 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "s", Type: "think.cycle.started", Timestamp: start, Metadata: map[string]any{"trigger": "manual"}}}, {Activity: model.Activity{ID: "r", Type: "think.relation.created", Timestamp: start.Add(time.Second), Message: "Relation erstellt", Metadata: map[string]any{"mutation_attribution": "explicit", "run_edges_created": 1}}, Point: AnalysisPoint{Delta: MutationStats{EdgesCreated: 999}}}, {Activity: model.Activity{ID: "e", Type: "think.cycle.completed", Timestamp: start.Add(1500 * time.Millisecond), Metadata: map[string]any{"duration_ms": 1500, "relations_created": 1}}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 { t.Fatalf("expected one run, got %d: %+v", len(runs), runs) } if runs[0].Status != "success" || runs[0].Mutations.EdgesCreated != 1 || runs[0].Verdict != "positives Ergebnis" { t.Fatalf("unexpected run: %+v", runs[0]) } } func TestEdgeUpdateChangeExplainsSimilarityRecalculation(t *testing.T) { previous := model.Edge{ID: "e", Type: "related", Origin: "ai-inference", Status: "active", Confidence: .72, Metadata: map[string]any{"semantic_similarity": .81}} current := model.Edge{ID: "e", Type: "related", Origin: "ai-inference", Status: "active", Confidence: .88, Metadata: map[string]any{"semantic_similarity": .93}} change := edgeUpdateChange(previous, current) if change.Action != "updated" || change.Details["previous_semantic_similarity"] != .81 || change.Details["semantic_similarity"] != .93 { t.Fatalf("similarity delta missing: %+v", change) } if change.Details["previous_confidence"] != .72 || change.Details["confidence"] != .88 { t.Fatalf("confidence delta missing: %+v", change) } } func TestBuildAnalysisRunsTreatsResearchFetchAsNestedEvents(t *testing.T) { start := time.Now().UTC().Add(-3 * time.Second) meta := map[string]any{"research_id": "r-1", "research_query": "evidence tools"} events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "s", Type: "article.research.started", Timestamp: start, Metadata: meta}}, {Activity: model.Activity{ID: "f1s", Type: "article.research.fetch.started", Timestamp: start.Add(100 * time.Millisecond), Metadata: meta}}, {Activity: model.Activity{ID: "f1e", Type: "article.research.fetch.completed", Timestamp: start.Add(200 * time.Millisecond), Metadata: meta}}, {Activity: model.Activity{ID: "f2s", Type: "article.research.fetch.started", Timestamp: start.Add(300 * time.Millisecond), Metadata: meta}}, {Activity: model.Activity{ID: "f2e", Type: "article.research.fetch.completed", Timestamp: start.Add(400 * time.Millisecond), Metadata: meta}}, {Activity: model.Activity{ID: "e", Type: "article.research.completed", Timestamp: start.Add(time.Second), Metadata: meta}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 { t.Fatalf("expected exactly one research run, got %d: %+v", len(runs), runs) } if runs[0].EventCount != len(events) || runs[0].Status == "running" { t.Fatalf("nested fetch lifecycle was not grouped correctly: %+v", runs[0]) } } func TestBuildAnalysisRunsDoesNotDuplicateActiveOrderForRepeatedStart(t *testing.T) { start := time.Now().UTC().Add(-time.Second) meta := map[string]any{"research_id": "r-dup", "research_query": "same query"} events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "s1", Type: "article.research.started", Timestamp: start, Metadata: meta}}, {Activity: model.Activity{ID: "s2", Type: "article.research.started", Timestamp: start.Add(10 * time.Millisecond), Metadata: meta}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 { t.Fatalf("expected one active run for duplicate start ID, got %d: %+v", len(runs), runs) } if runs[0].EventCount != 2 || runs[0].Status != "running" { t.Fatalf("unexpected duplicate-start aggregation: %+v", runs[0]) } } func TestBuildAnalysisRunsClosesStaleLegacyRelationResearchFromResults(t *testing.T) { start := time.Now().UTC().Add(-5 * time.Minute) meta := map[string]any{"research_id": "legacy-r", "research_query": "legacy query"} events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "s", Type: "research.started", Timestamp: start, Metadata: meta}}, {Activity: model.Activity{ID: "r", Type: "research.results", Timestamp: start.Add(time.Second), Metadata: meta}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 { t.Fatalf("expected one legacy run, got %d: %+v", len(runs), runs) } if runs[0].Status == "running" || runs[0].Verdict != "abgeschlossen (Legacy)" { t.Fatalf("stale legacy research should not remain running: %+v", runs[0]) } } func TestBuildAnalysisRunsKeepsRelationResearchOpenUntilExplicitCompletion(t *testing.T) { start := time.Now().UTC().Add(-3 * time.Second) meta := map[string]any{"research_id": "rel-r", "research_query": "relation query"} events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "s", Type: "research.started", Timestamp: start, Metadata: meta}}, {Activity: model.Activity{ID: "res", Type: "research.results", Timestamp: start.Add(200 * time.Millisecond), Metadata: meta}}, {Activity: model.Activity{ID: "ing", Type: "research.ingested", Timestamp: start.Add(300 * time.Millisecond), Metadata: meta}}, {Activity: model.Activity{ID: "end", Type: "research.completed", Timestamp: start.Add(time.Second), Metadata: meta}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 { t.Fatalf("expected one relation research run, got %d: %+v", len(runs), runs) } if runs[0].EventCount != 4 || runs[0].DurationMS < 900 { t.Fatalf("research.ingested ended the run too early: %+v", runs[0]) } } func TestBuildAnalysisRunsSynthesizesLearningRunWithoutStart(t *testing.T) { completed := time.Now().UTC() events := []AnalysisEventRecord{{Activity: model.Activity{ID: "done", Type: "learning.scan.completed", Timestamp: completed, Metadata: map[string]any{"run_id": "scan-x", "result": "updated", "duration_ms": 2450}}}} runs := buildAnalysisRuns(events) if len(runs) != 1 || runs[0].Kind != "learning" || runs[0].Status == "running" { t.Fatalf("expected synthetic completed learning run, got %+v", runs) } if runs[0].DurationMS != 2450 || completed.Sub(runs[0].StartedAt) < 2400*time.Millisecond { t.Fatalf("duration/start reconstruction failed: %+v", runs[0]) } } func TestBuildAnalysisRunsGroupsSecurityLifecycle(t *testing.T) { start := time.Now().UTC().Add(-4 * time.Second) meta := map[string]any{"inbox_id": "inbox-1", "title": "Critical vendor advisory"} events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "s", Type: "source.security.started", Timestamp: start, Metadata: meta}}, {Activity: model.Activity{ID: "r", Type: "source.security.research", Timestamp: start.Add(time.Second), Metadata: map[string]any{"inbox_id": "inbox-1", "results": 2}}}, {Activity: model.Activity{ID: "e", Type: "source.security.materialized", Timestamp: start.Add(3 * time.Second), Metadata: map[string]any{"inbox_id": "inbox-1", "title": "Critical vendor advisory", "duration_ms": 3000, "confidence": .91, "mutation_attribution": "explicit", "run_nodes_created": 1, "run_edges_created": 1, "run_vectors_created": 1}}, Point: AnalysisPoint{Delta: MutationStats{NodesCreated: 100, EdgesCreated: 100, VectorsCreated: 5000}}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 || runs[0].Kind != "security-source" || runs[0].EventCount != 3 || runs[0].DurationMS != 3000 { t.Fatalf("security lifecycle was not grouped: %+v", runs) } if runs[0].Mutations.NodesCreated != 1 || runs[0].Mutations.EdgesCreated != 1 || runs[0].Mutations.VectorsCreated != 1 { t.Fatalf("security mutations missing: %+v", runs[0].Mutations) } } func TestBuildAnalysisRunsGroupsArticlePipeline(t *testing.T) { start := time.Now().UTC().Add(-8 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "plan", Type: "article.plan.started", Timestamp: start, Metadata: map[string]any{"trigger": "automatic"}}}, {Activity: model.Activity{ID: "draft", Type: "article.draft.started", Timestamp: start.Add(time.Second)}}, {Activity: model.Activity{ID: "review", Type: "article.review.completed", Timestamp: start.Add(5 * time.Second), Metadata: map[string]any{"supported_claims": 8}}}, {Activity: model.Activity{ID: "created", Type: "article.created", Timestamp: start.Add(7 * time.Second), Metadata: map[string]any{"title": "DNS Hardening"}}, Point: AnalysisPoint{Delta: MutationStats{NodesCreated: 1, EdgesCreated: 3, VectorsCreated: 1}}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 || runs[0].Kind != "article" || runs[0].Status == "running" || runs[0].EventCount != 4 { t.Fatalf("article lifecycle was not grouped: %+v", runs) } if runs[0].DurationMS < 6900 { t.Fatalf("article duration not measured from plan to terminal event: %+v", runs[0]) } } func TestBuildAnalysisRunStatsCalculatesPercentiles(t *testing.T) { runs := []AnalysisRun{ {Kind: "article", Status: "success", DurationMS: 1000, EventCount: 4}, {Kind: "article", Status: "success", DurationMS: 2000, EventCount: 5}, {Kind: "article", Status: "warning", DurationMS: 9000, EventCount: 6}, } stats := buildAnalysisRunStats(runs) if len(stats) != 1 { t.Fatalf("unexpected stats: %+v", stats) } if stats[0].Duration.P50MS != 2000 || stats[0].Duration.P95MS != 9000 || stats[0].Duration.TotalMS != 12000 { t.Fatalf("unexpected duration stats: %+v", stats[0].Duration) } } func TestLearningScanAggregateRetainsExactRunCountAndDuration(t *testing.T) { store := &Store{} base := time.Now().UTC().Add(-time.Minute) for i, duration := range []int64{2400, 2500, 2600} { store.addLearningScanAggregateLocked(analysisRecord{activity: model.Activity{Timestamp: base.Add(time.Duration(i) * 20 * time.Second), Metadata: map[string]any{"duration_ms": duration, "knowledge_elements": 59184, "ollama_ok": true}}, point: AnalysisPoint{NodeCount: 59296, EdgeCount: 179938}}) } record := store.flushLearningScanAggregateLocked() if record == nil { t.Fatal("expected compacted record") } if got := int(metadataNumber(record.activity.Metadata, "scan_count")); got != 3 { t.Fatalf("scan count=%d metadata=%+v", got, record.activity.Metadata) } if got := int64(metadataNumber(record.activity.Metadata, "avg_duration_ms")); got != 2500 { t.Fatalf("avg duration=%d metadata=%+v", got, record.activity.Metadata) } if got := int(metadataNumber(record.activity.Metadata, "equivalent_raw_events")); got != 6 { t.Fatalf("equivalent raw events=%d", got) } } func TestBuildAnalysisPipelineSummaryCountsSecurityAndClaims(t *testing.T) { events := []AnalysisEventRecord{ {Activity: model.Activity{Type: "source.security.materialized", Metadata: map[string]any{"confidence": .8, "severity": "high", "event_type": "vulnerability", "supplemental_sources": 2}}}, {Activity: model.Activity{Type: "article.review.completed", Metadata: map[string]any{"supported_claims": 5, "unsupported_claims_count": 1}}}, {Activity: model.Activity{Type: "article.research.inbox", Metadata: map[string]any{"results": 3}}}, } counts := map[string]int{"source.security.materialized": 1, "article.review.completed": 1, "article.created": 1} summary := buildAnalysisPipelineSummary(events, counts) if summary.Security.Materialized != 1 || summary.Security.Severities["high"] != 1 || summary.Security.SupplementalSources != 2 { t.Fatalf("unexpected security summary: %+v", summary.Security) } if summary.Articles.Created != 1 || summary.Articles.SupportedClaims != 5 || summary.Articles.UnsupportedClaims != 1 || summary.Articles.InboxResearchHits != 3 { t.Fatalf("unexpected article summary: %+v", summary.Articles) } } func TestEmbeddingBatchAggregatePreservesMutationsAndDetails(t *testing.T) { store := &Store{} base := time.Now().UTC().Add(-time.Minute) for i := 0; i < 3; i++ { store.addEmbeddingBatchAggregateLocked(analysisRecord{ activity: model.Activity{Timestamp: base.Add(time.Duration(i) * time.Second), NodeIDs: []string{string(rune('a' + i))}, Metadata: map[string]any{"batch_count": 16, "model": "embeddinggemma", "mutation_attribution": "explicit", "run_vectors_created": 16}}, point: AnalysisPoint{Delta: MutationStats{VectorsCreated: 999}, VectorCount: 100 + i*16}, changes: []GraphChange{{EntityKind: "vector", Action: "created", EntityID: string(rune('a' + i))}}, }) } record := store.flushEmbeddingBatchAggregateLocked() if record == nil || record.activity.Type != "embedding.batch.aggregate" { t.Fatalf("expected embedding aggregate, got %#v", record) } if got := int(metadataNumber(record.activity.Metadata, "batch_events")); got != 3 { t.Fatalf("batch events=%d metadata=%+v", got, record.activity.Metadata) } if got := int(metadataNumber(record.activity.Metadata, "batch_count")); got != 48 { t.Fatalf("element count=%d metadata=%+v", got, record.activity.Metadata) } if record.point.Delta.VectorsCreated != 48 || len(record.changes) != 3 { t.Fatalf("embedding mutations/details were not preserved: delta=%+v changes=%+v", record.point.Delta, record.changes) } } func TestBuildAnalysisRunsDoesNotAttachStandaloneEmbeddingToSecurity(t *testing.T) { start := time.Now().UTC().Add(-5 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "sec", Type: "source.security.started", Timestamp: start, Metadata: map[string]any{"run_id": "sec-run", "inbox_id": "inbox"}}}, {Activity: model.Activity{ID: "emb", Type: "embedding.batch.aggregate", Timestamp: start.Add(time.Second), Metadata: map[string]any{"batch_count": 256, "mutation_attribution": "explicit", "run_vectors_created": 256}}}, {Activity: model.Activity{ID: "done", Type: "source.security.materialized", Timestamp: start.Add(2 * time.Second), Metadata: map[string]any{"run_id": "sec-run", "inbox_id": "inbox", "duration_ms": 2000, "mutation_attribution": "explicit", "run_nodes_created": 1, "run_edges_created": 1, "run_vectors_created": 1}}}, } runs := buildAnalysisRuns(events) if len(runs) != 2 { t.Fatalf("expected security + standalone embedding run, got %+v", runs) } for _, run := range runs { if run.Kind == "security-source" { if run.EventCount != 2 || run.Mutations.VectorsCreated != 1 { t.Fatalf("security run contaminated by embedding telemetry: %+v", run) } } } } func TestReconcileSecurityLifecycleClosesMissingTerminal(t *testing.T) { started := time.Now().UTC().Add(-10 * time.Second) h := AnalysisHistory{Runs: []AnalysisRun{{ID: "security:run-1", Kind: "security-source", Title: "Advisory", Status: "running", Verdict: "läuft", StartedAt: started, Metrics: map[string]any{}}}, Pipelines: AnalysisPipelineSummary{Security: AnalysisSecuritySummary{Severities: map[string]int{}, EventTypes: map[string]int{}}}} ReconcileSecurityLifecycles(&h, []AnalysisSecurityLifecycle{{InboxID: "inbox-1", RunID: "run-1", Title: "Advisory", Status: "materialized", ProactiveState: "done", Outcome: "materialized", StartedAt: started, CompletedAt: started.Add(6 * time.Second), DurationMS: 6000, MaterializedNodeID: "node-1", Confidence: .9, Severity: "high", EventType: "vulnerability", Mutations: MutationStats{NodesCreated: 1, EdgesCreated: 1, VectorsCreated: 1}}}) if len(h.Runs) != 1 || h.Runs[0].Status != "success" || h.Runs[0].DurationMS != 6000 { t.Fatalf("lifecycle reconciliation failed: %+v", h.Runs) } if h.Runs[0].Mutations.VectorsCreated != 1 || !h.Runs[0].MutationsKnown { t.Fatalf("authoritative mutations missing: %+v", h.Runs[0]) } if h.Pipelines.Security.Materialized != 1 || h.Pipelines.Security.AuthoritativeRecords != 1 { t.Fatalf("security summary not authoritative: %+v", h.Pipelines.Security) } } func TestBuildAnalysisRunsKeepsConcurrentQueriesSeparateByRunID(t *testing.T) { start := time.Now().UTC().Add(-3 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "a-start", Type: "query.started", Timestamp: start, Query: "A", Metadata: map[string]any{"run_id": "qa"}}}, {Activity: model.Activity{ID: "b-start", Type: "query.started", Timestamp: start.Add(100 * time.Millisecond), Query: "B", Metadata: map[string]any{"run_id": "qb"}}}, {Activity: model.Activity{ID: "a-hit", Type: "node.activated", Timestamp: start.Add(200 * time.Millisecond), Query: "A", Metadata: map[string]any{"run_id": "qa"}}}, {Activity: model.Activity{ID: "b-done", Type: "query.completed", Timestamp: start.Add(500 * time.Millisecond), Query: "B", Metadata: map[string]any{"run_id": "qb", "duration_ms": 400}}}, {Activity: model.Activity{ID: "a-done", Type: "query.completed", Timestamp: start.Add(900 * time.Millisecond), Query: "A", Metadata: map[string]any{"run_id": "qa", "duration_ms": 900}}}, } runs := buildAnalysisRuns(events) if len(runs) != 2 { t.Fatalf("expected two query runs, got %+v", runs) } counts := map[string]int{} for _, run := range runs { counts[run.ID] = run.EventCount if run.Status == "running" { t.Fatalf("query remained running: %+v", run) } } if counts["query:qa"] != 3 || counts["query:qb"] != 2 { t.Fatalf("query events crossed run boundaries: %+v", counts) } } func TestAnalysisAggregateIDsRemainUniqueAtSameWindowsClockTick(t *testing.T) { store := &Store{processID: "process-test"} timestamp := time.Unix(123, 456).UTC() store.analysisMu.Lock() first := store.nextAnalysisAggregateIDLocked("embedding-batch-aggregate", timestamp) second := store.nextAnalysisAggregateIDLocked("embedding-batch-aggregate", timestamp) third := store.nextAnalysisAggregateIDLocked("learning-scan-aggregate", timestamp) store.analysisMu.Unlock() if first == second || first == third || second == third { t.Fatalf("aggregate IDs must remain unique even with identical timestamps: %q %q %q", first, second, third) } } func TestAnalysisEventOrderKeepsReviewBeforeArticleTerminalAtSameTimestamp(t *testing.T) { stamp := time.Now().UTC() events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "plan", Type: "article.plan.started", Timestamp: stamp.Add(-time.Second)}}, {Activity: model.Activity{ID: "reject", Type: "article.draft.rejected", Timestamp: stamp}}, {Activity: model.Activity{ID: "review", Type: "article.review.completed", Timestamp: stamp, Metadata: map[string]any{"supported_claims": 3}}}, } sort.SliceStable(events, func(i, j int) bool { a, b := events[i].Activity, events[j].Activity if a.Timestamp.Equal(b.Timestamp) { pa, pb := analysisEventOrder(a), analysisEventOrder(b) if pa != pb { return pa < pb } return a.ID < b.ID } return a.Timestamp.Before(b.Timestamp) }) runs := buildAnalysisRuns(events) if len(runs) != 1 || runs[0].EventCount != 3 || runs[0].Status == "running" { t.Fatalf("same-timestamp review/terminal split article run: %+v", runs) } } func TestGraphUpdatedWithLearningRunIDIsNotStandalone(t *testing.T) { stamp := time.Now().UTC() events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "g", Type: "graph.updated", Timestamp: stamp, Metadata: map[string]any{"run_id": "learning-1"}}}, {Activity: model.Activity{ID: "done", Type: "learning.scan.completed", Timestamp: stamp.Add(time.Second), Metadata: map[string]any{"run_id": "learning-1", "duration_ms": 1000}}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 || runs[0].Kind != "learning" { t.Fatalf("graph.updated duplicated learning run: %+v", runs) } } func TestBuildAnalysisRunsKeepsConcurrentArticlesSeparateByRunID(t *testing.T) { start := time.Now().UTC().Add(-10 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "a-plan", Type: "article.plan.started", Timestamp: start, Metadata: map[string]any{"run_id": "article-a"}}}, {Activity: model.Activity{ID: "b-plan", Type: "article.plan.started", Timestamp: start.Add(time.Millisecond), Metadata: map[string]any{"run_id": "article-b"}}}, {Activity: model.Activity{ID: "a-review", Type: "article.review.completed", Timestamp: start.Add(time.Second), Metadata: map[string]any{"run_id": "article-a", "supported_claims": 4}}}, {Activity: model.Activity{ID: "b-review", Type: "article.review.completed", Timestamp: start.Add(2 * time.Second), Metadata: map[string]any{"run_id": "article-b", "supported_claims": 5}}}, {Activity: model.Activity{ID: "b-created", Type: "article.created", Timestamp: start.Add(3 * time.Second), Metadata: map[string]any{"run_id": "article-b", "title": "B"}}}, {Activity: model.Activity{ID: "a-rejected", Type: "article.draft.rejected", Timestamp: start.Add(4 * time.Second), Metadata: map[string]any{"run_id": "article-a"}}}, } runs := buildAnalysisRuns(events) if len(runs) != 2 { t.Fatalf("expected two independent article runs, got %d: %+v", len(runs), runs) } byKey := map[string]AnalysisRun{} for _, run := range runs { byKey[run.ID] = run } a, okA := byKey["article:article-a"] b, okB := byKey["article:article-b"] if !okA || !okB { t.Fatalf("native article run IDs were not preserved: %+v", runs) } if a.EventCount != 3 || b.EventCount != 3 || a.Status == "running" || b.Status == "running" { t.Fatalf("concurrent article events crossed run boundaries: a=%+v b=%+v", a, b) } } func TestBuildAnalysisRunsDoesNotAttachClusterSchedulingToOpenArticle(t *testing.T) { start := time.Now().UTC().Add(-5 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "plan", Type: "article.plan.started", Timestamp: start, Metadata: map[string]any{"run_id": "article-a"}}}, {Activity: model.Activity{ID: "cluster", Type: "article.cluster.started", Timestamp: start.Add(time.Second), Metadata: map[string]any{"topic_guard": "strict-v2"}}}, {Activity: model.Activity{ID: "created", Type: "article.created", Timestamp: start.Add(2 * time.Second), Metadata: map[string]any{"run_id": "article-a"}}}, } runs := buildAnalysisRuns(events) var article AnalysisRun for _, run := range runs { if run.Kind == "article" && run.ID == "article:article-a" { article = run } } if article.EventCount != 2 { t.Fatalf("cluster scheduler event leaked into article lifecycle: %+v / all=%+v", article, runs) } } func TestBuildAnalysisRunsUsesLearningTerminalMutationsOnce(t *testing.T) { start := time.Now().UTC().Add(-3 * time.Second) stats := map[string]any{"mutation_attribution": "explicit", "run_nodes_created": 10, "run_edges_created": 20, "run_vectors_created": 30} events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "start", Type: "learning.scan.started", Timestamp: start, Metadata: map[string]any{"run_id": "learning-1"}}}, {Activity: model.Activity{ID: "graph", Type: "graph.updated", Timestamp: start.Add(time.Second), Metadata: mergeTestMetadata(map[string]any{"run_id": "learning-1"}, stats)}}, {Activity: model.Activity{ID: "done", Type: "learning.scan.completed", Timestamp: start.Add(2 * time.Second), Metadata: mergeTestMetadata(map[string]any{"run_id": "learning-1", "duration_ms": 2000}, stats)}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 { t.Fatalf("expected one learning run, got %+v", runs) } got := runs[0].Mutations if got.NodesCreated != 10 || got.EdgesCreated != 20 || got.VectorsCreated != 30 { t.Fatalf("learning mutations were double-counted: %+v", got) } } func mergeTestMetadata(parts ...map[string]any) map[string]any { out := map[string]any{} for _, part := range parts { for key, value := range part { out[key] = value } } return out } func TestBuildAnalysisRunsDoesNotAttachUnkeyedTelemetryByTime(t *testing.T) { start := time.Now().UTC().Add(-3 * time.Second) events := []AnalysisEventRecord{ {Activity: model.Activity{ID: "plan", Type: "article.plan.started", Timestamp: start, Metadata: map[string]any{"run_id": "article-a"}}}, {Activity: model.Activity{ID: "inbox", Type: "article.research.inbox", Timestamp: start.Add(time.Second), Metadata: map[string]any{"results": 3}}}, {Activity: model.Activity{ID: "done", Type: "article.created", Timestamp: start.Add(2 * time.Second), Metadata: map[string]any{"run_id": "article-a"}}}, } runs := buildAnalysisRuns(events) if len(runs) != 1 || runs[0].EventCount != 2 { t.Fatalf("unkeyed telemetry was attached by time: %+v", runs) } }