package db import ( "path/filepath" "sync" "testing" ) func tmpDB(t *testing.T) *KnoxDB { t.Helper() kdb, err := Open(filepath.Join(t.TempDir(), "index.db")) if err != nil { t.Fatalf("open db: %v", err) } t.Cleanup(func() { kdb.Close() }) return kdb } func record(t *testing.T, k *KnoxDB, fp, title string) { t.Helper() _, _, err := k.RecordObservation(ObservationRecord{ Fingerprint: fp, SourceID: "test", SourcePath: fp, Project: "itest", ContentType: "test", Title: title, Summary: "summary", CreatedAt: "2026-08-29T00:00:00Z", LineEnd: 0, Confidence: 0.9, IngesterVersion: "test/v1", }) if err != nil { t.Errorf("record %s: %v", fp, err) } } func obsCount(t *testing.T, k *KnoxDB) int { t.Helper() stats, err := k.Stats() if err != nil { t.Fatalf("stats: %v", err) } n, _ := stats["total_observations"].(int) return n } // TestRecordObservationConcurrentDedup: N goroutines ingesting identical // content must produce exactly one observation row. This exercises the // check-then-insert dedup under the writeMu + BEGIN IMMEDIATE serialization. func TestRecordObservationConcurrentDedup(t *testing.T) { k := tmpDB(t) var wg sync.WaitGroup for i := 0; i < 20; i++ { wg.Add(1) go func() { defer wg.Done() record(t, k, "fp:concurrent", "same title") }() } wg.Wait() if n := obsCount(t, k); n != 1 { t.Fatalf("expected exactly 1 observation after concurrent identical ingests, got %d", n) } } // TestRecordObservationDistinctFingerprints: different content must never be // deduped away (the constraint is per-fingerprint signature, not global). func TestRecordObservationDistinctFingerprints(t *testing.T) { k := tmpDB(t) record(t, k, "fp:a", "alpha") record(t, k, "fp:b", "beta") record(t, k, "fp:a", "alpha changed") if n := obsCount(t, k); n != 3 { t.Fatalf("expected 3 observations, got %d", n) } } // TestPushObservationsValidation: forged/malformed rows are rejected without // error — self node_id (poisoning vector), empty node_id, negative HCL. func TestPushObservationsValidation(t *testing.T) { k := tmpDB(t) foreign := GossipObservation{ NodeID: "0123456789abcdef0123456789abcdef", HCL: 42, Fingerprint: "fp:foreign", SourceID: "test", Title: "t", Summary: "s", CollectedAt: "2026-08-29T00:00:00Z", } rows := []GossipObservation{ foreign, {NodeID: k.NodeID(), HCL: 100, Fingerprint: "fp:self"}, // spoof poisoning attempt {NodeID: "", HCL: 1, Fingerprint: "fp:empty"}, {NodeID: "other", HCL: -5, Fingerprint: "fp:neg"}, } n, err := k.PushObservations(rows) if err != nil { t.Fatalf("push: %v", err) } if n != 1 { t.Fatalf("expected exactly the one valid row inserted, got %d", n) } if got := obsCount(t, k); got != 1 { t.Fatalf("expected 1 observation in the log, got %d", got) } } // TestOpenReopenHCLMonotonicAcrossRestart: reopening a DB must resume the HLC // from its persisted max (clock seeding), keep the node identity, and order the // new observation above every previous one. func TestOpenReopenHCLMonotonicAcrossRestart(t *testing.T) { path := filepath.Join(t.TempDir(), "index.db") k1, err := Open(path) if err != nil { t.Fatalf("open: %v", err) } record(t, k1, "fp:r1", "one") record(t, k1, "fp:r2", "two") record(t, k1, "fp:r3", "three") nodeID1 := k1.NodeID() maxBefore := int64(0) rows, err := k1.ObservationsAfter(nodeID1, 0, 100) if err != nil { t.Fatalf("obs after: %v", err) } for _, r := range rows { if r.HCL > maxBefore { maxBefore = r.HCL } } if err := k1.Close(); err != nil { t.Fatalf("close: %v", err) } k2, err := Open(path) if err != nil { t.Fatalf("reopen: %v", err) } defer k2.Close() if k2.NodeID() != nodeID1 { t.Errorf("node id changed across reopen: %q -> %q", nodeID1, k2.NodeID()) } record(t, k2, "fp:r4", "four") after, err := k2.ObservationsAfter(nodeID1, maxBefore, 100) if err != nil { t.Fatalf("obs after (reopened): %v", err) } if len(after) != 1 { t.Fatalf("expected exactly the new observation above the pre-restart max, got %d rows", len(after)) } if after[0].Fingerprint != "fp:r4" { t.Errorf("unexpected row above max: %s", after[0].Fingerprint) } } // TestCreateThreadClusterConcurrentIdempotent: two concurrent creators of the // same cluster key must converge on one thread (one created=true, one false) // with no UNIQUE constraint error — the INSERT OR IGNORE path. func TestCreateThreadClusterConcurrentIdempotent(t *testing.T) { k := tmpDB(t) const key = "cluster:race" var wg sync.WaitGroup ids := make([]int64, 2) created := make([]bool, 2) errs := make([]error, 2) for i := 0; i < 2; i++ { wg.Add(1) go func(i int) { defer wg.Done() ids[i], created[i], errs[i] = k.CreateThreadCluster("Race thread", "m", "medium", "[]", `{}`, key) }(i) } wg.Wait() for i := range errs { if errs[i] != nil { t.Fatalf("creator %d: %v", i, errs[i]) } } if ids[0] != ids[1] { t.Errorf("concurrent creators got different ids: %d vs %d", ids[0], ids[1]) } if created[0] == created[1] { t.Errorf("exactly one creator should report created=true, got %v %v", created[0], created[1]) } threads, err := k.ListThreads("") if err != nil { t.Fatal(err) } n := 0 for _, th := range threads { if th.Title == "Race thread" { n++ } } if n != 1 { t.Errorf("want exactly 1 race thread, got %d", n) } } // TestRebuildEntriesDeterministicFold: two nodes holding the same observations // in different merge orders must rebuild identical derived state — content // fields come from the highest-HCL observation per fingerprint (spec 6.2), // never from an arbitrary (merge-order-dependent) row. func TestRebuildEntriesDeterministicFold(t *testing.T) { const foreign = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" rows := []GossipObservation{ {NodeID: foreign, HCL: 10, Fingerprint: "fp:fold", SourceID: "test", Title: "old title", Summary: "old", CollectedAt: "2026-08-29T00:00:00Z", Confidence: 0.5}, {NodeID: foreign, HCL: 20, Fingerprint: "fp:fold", SourceID: "test", Title: "new title", Summary: "new", CollectedAt: "2026-08-29T01:00:00Z", Confidence: 0.9}, } results := make([]Entry, 2) for i, order := range [][]GossipObservation{{rows[0], rows[1]}, {rows[1], rows[0]}} { k := tmpDB(t) if _, err := k.PushObservations(order); err != nil { t.Fatalf("push: %v", err) } if _, _, err := k.RebuildEntriesFromObservations(); err != nil { t.Fatalf("rebuild: %v", err) } e, err := k.FindEntry("fp:fold") if err != nil || e == nil { t.Fatalf("find entry: %v", err) } results[i] = *e } if results[0].Title != "new title" { t.Errorf("order [old,new]: title = %q, want %q", results[0].Title, "new title") } if results[1].Title != "new title" { t.Errorf("order [new,old]: title = %q, want %q", results[1].Title, "new title") } if results[0] != results[1] { t.Errorf("derived state diverged between merge orders:\n%+v\n%+v", results[0], results[1]) } if results[0].RefCount != 2 { t.Errorf("ref_count = %d, want 2", results[0].RefCount) } if results[0].FirstSeen != "2026-08-29T00:00:00Z" || results[0].LastSeen != "2026-08-29T01:00:00Z" { t.Errorf("first/last seen = %q/%q", results[0].FirstSeen, results[0].LastSeen) } }