返回 DeepSeek-Reasonix
metadata_projection_test.go
根目录 / internal / sessioncatalog / metadata_projection_test.go
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
159 lines GO