返回 DeepSeek-Reasonix
metadata_queue_test.go
根目录 / internal / sessioncatalog / metadata_queue_test.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "fmt"
6 "os"
7 "path/filepath"
8 "testing"
9 "time"
10 )
11
12 func TestMetadataIteratorAdmissionPreservesProgressAndVisibleSlot(t *testing.T) {
13 c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true})
14 if err != nil {
15 t.Fatal(err)
16 }
17 t.Cleanup(func() { _ = c.Close(context.Background()) })
18 jobs := map[string]*metadataQueueJob{}
19 defer func() {
20 for _, job := range jobs {
21 if job.scan != nil {
22 job.scan.close(t.Context(), context.Canceled)
23 }
24 }
25 }()
26 const roots = metadataIteratorLimit + 4
27 const records = 140
28 for i := range roots {
29 dir := t.TempDir()
30 for j := range records {
31 if err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("%03d.jsonl", j)), []byte("unread body"), 0600); err != nil {
32 t.Fatal(err)
33 }
34 }
35 jobs[dir] = &metadataQueueJob{target: DirectoryTarget{Path: dir, Scope: "project", WorkspaceRoot: dir, mutationSeq: uint64(i + 1)}}
36 }
37 now := time.Unix(1, 0)
38 visible := ""
39 priority := func(target DirectoryTarget) bool { return target.Path == visible }
40 started := map[string]*metadataScan{}
41 var turn uint64
42 run := func(key string, job *metadataQueueJob) {
43 t.Helper()
44 turn++
45 job.turn = turn
46 if job.scan == nil {
47 if started[key] != nil {
48 t.Fatal("time slice restarted a directory")
49 }
50 job.scan, err = c.startMetadataScan(t.Context(), job.target, job.target.mutationSeq, true)
51 if err != nil {
52 t.Fatal(err)
53 }
54 started[key] = job.scan
55 } else if started[key] != job.scan {
56 t.Fatal("iterator identity changed")
57 }
58 done, _, err := job.scan.step(t.Context())
59 if err != nil {
60 t.Fatal(err)
61 }
62 if done {
63 job.scan.close(t.Context(), nil)
64 count, err := c.CountDirectorySessions(t.Context(), key)
65 if err != nil || count != records || !c.DirectoryScanReady(t.Context(), key) {
66 t.Fatalf("incomplete discovery published: %d, %v", count, err)
67 }
68 delete(jobs, key)
69 }
70 active := 0
71 for _, pending := range jobs {
72 if pending.scan != nil {
73 active++
74 }
75 }
76 if active > metadataIteratorLimit {
77 t.Fatalf("iterator budget exceeded: %d", active)
78 }
79 }
80 for range metadataIteratorLimit - 1 {
81 key, job := selectMetadataQueueJob(jobs, now, false, priority)
82 if job == nil || job.scan != nil {
83 t.Fatal("waiting root did not receive its initial slice")
84 }
85 run(key, job)
86 }
87 _, job := selectMetadataQueueJob(jobs, now, false, priority)
88 if job == nil || job.scan == nil {
89 t.Fatal("background discovery consumed the visible-workspace reserve")
90 }
91 for path, pending := range jobs {
92 if pending.scan == nil {
93 visible = path
94 break
95 }
96 }
97 key, job := selectMetadataQueueJob(jobs, now, true, priority)
98 if job == nil || key != visible || job.scan != nil {
99 t.Fatal("foreground preparation blocked the visible root's reserved slot")
100 }
101 run(key, job)
102 // Changing the visible workspace while all slots are occupied must not
103 // evict a previous iterator or allocate a ninth one.
104 for path, pending := range jobs {
105 if pending.scan == nil {
106 visible = path
107 break
108 }
109 }
110 if _, job := selectMetadataQueueJob(jobs, now, true, priority); job != nil {
111 t.Fatal("priority change bypassed the iterator budget or resumed P2")
112 }
113 visible = ""
114 for len(jobs) > 0 {
115 key, job := selectMetadataQueueJob(jobs, now, false, priority)
116 if job == nil {
117 t.Fatal("waiting roots could not progress")
118 }
119 run(key, job)
120 }
121 if len(started) != roots {
122 t.Fatalf("only %d roots completed", len(started))
123 }
124 }
125
126 func TestMetadataQueueRotatesBeforeLargeRootCompletes(t *testing.T) {
127 large, small := t.TempDir(), t.TempDir()
128 for i := range 400 {
129 if err := os.WriteFile(filepath.Join(large, fmt.Sprintf("%04d.jsonl", i)), []byte("unreadable body"), 0600); err != nil {
130 t.Fatal(err)
131 }
132 }
133 if err := os.WriteFile(filepath.Join(small, "only.jsonl"), []byte("unreadable body"), 0600); err != nil {
134 t.Fatal(err)
135 }
136 c, err := Open(t.Context(), Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), MetadataOnly: true, StartPaused: true})
137 if err != nil {
138 t.Fatal(err)
139 }
140 t.Cleanup(func() { _ = c.Close(context.Background()) })
141 _, ok := c.ScheduleReconcile(DirectoryTarget{Path: large, Scope: "global"})
142 if !ok {
143 t.Fatal("large root rejected")
144 }
145 done, ok := c.ScheduleReconcile(DirectoryTarget{Path: small, Scope: "global"})
146 if !ok {
147 t.Fatal("small root rejected")
148 }
149 c.ResumeDiscovery()
150 select {
151 case <-done:
152 case <-t.Context().Done():
153 t.Fatal("small root did not settle")
154 }
155 if !c.DirectoryScanReady(t.Context(), small) {
156 t.Fatal("small root did not publish")
157 }
158 if c.DirectoryScanReady(t.Context(), large) {
159 t.Fatal("small root waited for a complete large scan")
160 }
161 count, err := c.CountDirectorySessions(t.Context(), large)
162 if err != nil || count == 0 || count >= 400 {
163 t.Fatalf("large root progress = %d: %v", count, err)
164 }
165 }
166
167 func TestMetadataQueueJournalSurvivesRestartWithoutStartingDiscovery(t *testing.T) {
168 root, database := t.TempDir(), filepath.Join(t.TempDir(), "catalog.sqlite")
169 path := filepath.Join(root, "kept.jsonl")
170 if err := os.WriteFile(path, []byte("body must not be decoded"), 0600); err != nil {
171 t.Fatal(err)
172 }
173 opts := Options{Path: database, MetadataOnly: true, StartPaused: true}
174 c, err := Open(t.Context(), opts)
175 if err != nil {
176 t.Fatal(err)
177 }
178 target := DirectoryTarget{Path: root, Scope: "global"}
179 if !c.RequestReconcile(target) {
180 t.Fatal("request rejected")
181 }
182 if err := c.Close(t.Context()); err != nil {
183 t.Fatal(err)
184 }
185 c, err = Open(t.Context(), opts)
186 if err != nil {
187 t.Fatal(err)
188 }
189 t.Cleanup(func() { _ = c.Close(context.Background()) })
190 var pending int
191 if err := c.db.QueryRow(`SELECT COUNT(*) FROM catalog_pending_roots`).Scan(&pending); err != nil || pending != 1 {
192 t.Fatalf("pending=%d: %v", pending, err)
193 }
194 if _, found, err := c.GetSession(t.Context(), path); err != nil || found {
195 t.Fatalf("paused startup scanned source: %v %v", found, err)
196 }
197 done, ok := c.ScheduleReconcile(target)
198 if !ok {
199 t.Fatal("resume rejected")
200 }
201 c.ResumeDiscovery()
202 select {
203 case <-done:
204 case <-t.Context().Done():
205 t.Fatal("resumed root did not settle")
206 }
207 if _, found, err := c.GetSession(t.Context(), path); err != nil || !found {
208 t.Fatalf("source missing: %v %v", found, err)
209 }
210 if err := c.db.QueryRow(`SELECT COUNT(*) FROM catalog_pending_roots`).Scan(&pending); err != nil || pending != 0 {
211 t.Fatalf("completed journal=%d: %v", pending, err)
212 }
213 }
214
214 lines GO