| 1 | package sessioncatalog |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "fmt" |
| 7 | "sync" |
| 8 | "testing" |
| 9 | |
| 10 | "reasonix/internal/historywork" |
| 11 | ) |
| 12 | |
| 13 | func TestMetadataIncrementalWaitingObservationCanCancel(t *testing.T) { |
| 14 | c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true}) |
| 15 | if err != nil { |
| 16 | t.Fatal(err) |
| 17 | } |
| 18 | entered, release := make(chan struct{}), make(chan struct{}) |
| 19 | var once, unblock sync.Once |
| 20 | t.Cleanup(func() { unblock.Do(func() { close(release) }); _ = c.Close(context.Background()) }) |
| 21 | c.testMetadataSliceHook = func(int) { once.Do(func() { close(entered); <-release }) } |
| 22 | completed := make(chan error, 1) |
| 23 | go func() { completed <- c.SyncMetadata(t.Context(), nil, nil) }() |
| 24 | <-entered |
| 25 | ctx, cancel := context.WithCancel(t.Context()) |
| 26 | cancel() |
| 27 | if err := c.SyncMetadata(ctx, nil, nil); !errors.Is(err, context.Canceled) { |
| 28 | t.Fatalf("waiting observation ignored cancellation: %v", err) |
| 29 | } |
| 30 | unblock.Do(func() { close(release) }) |
| 31 | if err := <-completed; err != nil { |
| 32 | t.Fatal(err) |
| 33 | } |
| 34 | } |
| 35 | |
| 36 | func TestMetadataIncrementalSkipsUnchangedRegistryWrites(t *testing.T) { |
| 37 | c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true}) |
| 38 | if err != nil { |
| 39 | t.Fatal(err) |
| 40 | } |
| 41 | t.Cleanup(func() { _ = c.Close(context.Background()) }) |
| 42 | projects := []ProjectRecord{{Scope: "global", Title: "Global"}} |
| 43 | topics := []TopicMetadata{{Scope: "global", TopicID: "kept", Title: "Named", Pinned: true, CreatedAt: 42}} |
| 44 | if err := c.SyncMetadata(t.Context(), projects, topics); err != nil { |
| 45 | t.Fatal(err) |
| 46 | } |
| 47 | _, err = c.db.Exec(`CREATE TABLE registry_writes(kind TEXT); |
| 48 | CREATE TRIGGER watch_projects AFTER UPDATE ON catalog_projects BEGIN INSERT INTO registry_writes VALUES('project'); END; |
| 49 | CREATE TRIGGER watch_topics AFTER UPDATE ON catalog_topics BEGIN INSERT INTO registry_writes VALUES('topic'); END;`) |
| 50 | if err != nil { |
| 51 | t.Fatal(err) |
| 52 | } |
| 53 | revision := c.revision.Load() |
| 54 | if err := c.SyncMetadata(t.Context(), projects, topics); err != nil { |
| 55 | t.Fatal(err) |
| 56 | } |
| 57 | var count int |
| 58 | if err := c.db.QueryRow(`SELECT count(*) FROM registry_writes`).Scan(&count); err != nil || count != 0 || revision != c.revision.Load() { |
| 59 | t.Fatalf("unchanged registry republished: writes=%d revision=%d/%d err=%v", count, revision, c.revision.Load(), err) |
| 60 | } |
| 61 | // The current database, including an older writer's update, must be checked |
| 62 | // again. A remembered input hash would incorrectly skip this correction. |
| 63 | if _, err := c.db.Exec(`UPDATE catalog_topics SET title='other writer' WHERE topic_id='kept'`); err != nil { |
| 64 | t.Fatal(err) |
| 65 | } |
| 66 | if err := c.SyncMetadata(t.Context(), projects, topics); err != nil { |
| 67 | t.Fatal(err) |
| 68 | } |
| 69 | var title string |
| 70 | if err := c.db.QueryRow(`SELECT title FROM catalog_topics WHERE topic_id='kept'`).Scan(&title); err != nil || title != "Named" { |
| 71 | t.Fatalf("old writer bypassed projection validation: %q %v", title, err) |
| 72 | } |
| 73 | } |
| 74 | |
| 75 | func TestMetadataIncrementalCancellationKeepsUnvisitedMembership(t *testing.T) { |
| 76 | c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true}) |
| 77 | if err != nil { |
| 78 | t.Fatal(err) |
| 79 | } |
| 80 | t.Cleanup(func() { _ = c.Close(context.Background()) }) |
| 81 | if err := c.SyncMetadata(t.Context(), nil, []TopicMetadata{{Scope: "global", TopicID: "retire-after-success", Title: "Old"}}); err != nil { |
| 82 | t.Fatal(err) |
| 83 | } |
| 84 | topics := make([]TopicMetadata, 3*historywork.BatchEntries) |
| 85 | for i := range topics { |
| 86 | topics[i] = TopicMetadata{Scope: "global", TopicID: fmt.Sprintf("new-%04d", i), Title: "New"} |
| 87 | } |
| 88 | ctx, cancel := context.WithCancel(t.Context()) |
| 89 | defer cancel() |
| 90 | observed := 0 |
| 91 | c.testMetadataSliceHook = func(entries int) { |
| 92 | observed += entries |
| 93 | if entries > historywork.BatchEntries { |
| 94 | t.Fatalf("unbounded registry slice: %d", entries) |
| 95 | } |
| 96 | if !c.mutationMu.TryLock() { |
| 97 | t.Fatal("registry kept the writer between slices") |
| 98 | } |
| 99 | c.mutationMu.Unlock() |
| 100 | cancel() |
| 101 | } |
| 102 | if err := c.SyncMetadata(ctx, nil, topics); !errors.Is(err, context.Canceled) { |
| 103 | t.Fatalf("expected cancellation: %v", err) |
| 104 | } |
| 105 | if observed == 0 || observed > historywork.BatchEntries { |
| 106 | t.Fatalf("cancellation restarted or consumed unvisited input: %d", observed) |
| 107 | } |
| 108 | var retained int |
| 109 | if err := c.db.QueryRow(`SELECT count(*) FROM catalog_topics WHERE topic_id='retire-after-success' AND metadata_present=1`).Scan(&retained); err != nil || retained != 1 { |
| 110 | t.Fatalf("partial refresh retired unvisited membership: %d %v", retained, err) |
| 111 | } |
| 112 | c.testMetadataSliceHook = nil |
| 113 | if err := c.SyncMetadata(t.Context(), nil, topics); err != nil { |
| 114 | t.Fatal(err) |
| 115 | } |
| 116 | var count int |
| 117 | if err := c.db.QueryRow(`SELECT count(*) FROM catalog_topics WHERE metadata_present=1`).Scan(&count); err != nil || count != len(topics) { |
| 118 | t.Fatalf("retry lost registry input: %d %v", count, err) |
| 119 | } |
| 120 | if err := c.db.QueryRow(`SELECT count(*) FROM catalog_topics WHERE topic_id='retire-after-success'`).Scan(&retained); err != nil || retained != 0 { |
| 121 | t.Fatalf("completed refresh did not retire orphan: %d %v", retained, err) |
| 122 | } |
| 123 | } |
| 124 | |
| 125 | func TestMetadataIncrementalAllowsSourceCommitBetweenSlices(t *testing.T) { |
| 126 | c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true}) |
| 127 | if err != nil { |
| 128 | t.Fatal(err) |
| 129 | } |
| 130 | t.Cleanup(func() { _ = c.Close(context.Background()) }) |
| 131 | topics := make([]TopicMetadata, 2*historywork.BatchEntries) |
| 132 | for i := range topics { |
| 133 | topics[i] = TopicMetadata{Scope: "global", TopicID: fmt.Sprintf("topic-%04d", i), Title: "Registered"} |
| 134 | } |
| 135 | committed := false |
| 136 | c.testMetadataSliceHook = func(int) { |
| 137 | if committed { |
| 138 | return |
| 139 | } |
| 140 | committed = true |
| 141 | if !c.mutationMu.TryLock() { |
| 142 | t.Fatal("source commit cannot acquire the writer") |
| 143 | } |
| 144 | c.mutationMu.Unlock() |
| 145 | err := c.UpsertSession(t.Context(), SessionRecord{Path: "/fixture/saved.jsonl", Directory: "/fixture", Scope: "global", |
| 146 | TopicID: "source-only", TopicTitle: "Concurrent source", Turns: 4, TurnsState: TurnsValid, Health: HealthOK}) |
| 147 | if err != nil { |
| 148 | t.Fatal(err) |
| 149 | } |
| 150 | } |
| 151 | if err := c.SyncMetadata(t.Context(), nil, topics); err != nil { |
| 152 | t.Fatal(err) |
| 153 | } |
| 154 | var count int |
| 155 | if err := c.db.QueryRow(`SELECT count(*) FROM catalog_topics WHERE topic_id='source-only'`).Scan(&count); err != nil || count != 1 { |
| 156 | t.Fatalf("registry retirement removed concurrent source: %d %v", count, err) |
| 157 | } |
| 158 | } |
| 159 |