Files
glpi-neural-brain/internal/graph/analysis_dashboard_test.go
jbergner 440423c5b6
All checks were successful
release-tag / release-image (push) Successful in 2m43s
RC-3
2026-08-09 11:29:13 +02:00

470 lines
25 KiB
Go

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)
}
}