返回 DeepSeek-Reasonix
metadata_dispatch_test.go
根目录 / internal / sessioncatalog / metadata_dispatch_test.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "os"
6 "path/filepath"
7 "sync"
8 "testing"
9 )
10
11 func TestMetadataResumeCoalescesWatchedRootsWithInterruptedJournal(t *testing.T) {
12 root := t.TempDir()
13 if err := os.WriteFile(filepath.Join(root, "one.jsonl"), []byte("unread body"), 0600); err != nil {
14 t.Fatal(err)
15 }
16 opts := Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), MetadataOnly: true, StartPaused: true}
17 c, err := Open(t.Context(), opts)
18 if err != nil {
19 t.Fatal(err)
20 }
21 if !c.RequestReconcile(DirectoryTarget{Path: root, Scope: "project", WorkspaceRoot: root}) {
22 t.Fatal("initial journal entry rejected")
23 }
24 if err := c.Close(t.Context()); err != nil {
25 t.Fatal(err)
26 }
27 c, err = Open(t.Context(), opts)
28 if err != nil {
29 t.Fatal(err)
30 }
31 t.Cleanup(func() { _ = c.Close(context.Background()) })
32 // Attach to the restored completion before dispatch. The watcher then
33 // supplies a newer scope, as it does after restoring project identities.
34 done, accepted := c.ScheduleReconcile(DirectoryTarget{Path: root, Scope: "project", WorkspaceRoot: root})
35 if !accepted {
36 t.Fatal("restored task rejected")
37 }
38 var starts []DirectoryTarget
39 c.testReconcileStartHook = func(target DirectoryTarget) { starts = append(starts, target) }
40 if rejected := c.ResumeDiscovery(DirectoryTarget{Path: root, Scope: "global"}); len(rejected) != 0 {
41 t.Fatalf("watched root rejected: %+v", rejected)
42 }
43 select {
44 case <-done:
45 case <-t.Context().Done():
46 t.Fatal("resumed scan did not settle")
47 }
48 var pending int
49 if err := c.db.QueryRow(`SELECT COUNT(*) FROM catalog_pending_roots`).Scan(&pending); err != nil || pending != 0 {
50 t.Fatalf("completed journal=%d: %v", pending, err)
51 }
52 if err := c.Close(t.Context()); err != nil {
53 t.Fatal(err)
54 }
55 if len(starts) != 1 || starts[0].Scope != "global" {
56 t.Fatalf("startup must dispatch latest watched identity once: %+v", starts)
57 }
58 }
59
60 // Both roots have queue jobs while the first holds a slice. The second's
61 // invalidation precedes dispatch and must not become a redundant follow-up.
62 func TestMetadataQueueCoalescesUpdatesWhileWaitingForDispatch(t *testing.T) {
63 first, waiting := t.TempDir(), t.TempDir()
64 for _, dir := range []string{first, waiting} {
65 if err := os.WriteFile(filepath.Join(dir, "one.jsonl"), []byte("unread body"), 0600); err != nil {
66 t.Fatal(err)
67 }
68 }
69 c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true})
70 if err != nil {
71 t.Fatal(err)
72 }
73 entered, release := make(chan struct{}), make(chan struct{})
74 var gate sync.Once
75 t.Cleanup(func() { gate.Do(func() { close(release) }); _ = c.Close(context.Background()) })
76 var slice sync.Once
77 c.testReconcileBatchHook = func(int) {
78 slice.Do(func() {
79 close(entered)
80 <-release
81 })
82 }
83 var starts []DirectoryTarget
84 c.testReconcileStartHook = func(target DirectoryTarget) { starts = append(starts, target) }
85 c.PrioritizeWorkspace("global", "")
86 if !c.RequestReconcile(DirectoryTarget{Path: first, Scope: "global"}) {
87 t.Fatal("first root rejected")
88 }
89 done, accepted := c.ScheduleReconcile(DirectoryTarget{Path: waiting, Scope: "project", WorkspaceRoot: waiting})
90 if !accepted {
91 t.Fatal("waiting root rejected")
92 }
93 c.ResumeDiscovery()
94 select {
95 case <-entered:
96 case <-t.Context().Done():
97 t.Fatal("first scan never entered")
98 }
99 // Persisted invalidations retain their latest scope and sequence too.
100 joined, accepted := c.ScheduleReconcile(DirectoryTarget{Path: waiting, Scope: "global"})
101 if !accepted || joined != done {
102 t.Fatal("waiting invalidation did not join the existing completion signal")
103 }
104 gate.Do(func() { close(release) })
105 select {
106 case <-done:
107 case <-t.Context().Done():
108 t.Fatal("waiting root did not settle")
109 }
110 var pending int
111 if err := c.db.QueryRow(`SELECT count(*) FROM catalog_pending_roots WHERE path_key=?`, queuePathKey(waiting)).Scan(&pending); err != nil || pending != 0 {
112 t.Fatalf("latest pre-dispatch journal entry was not retired: %d %v", pending, err)
113 }
114 if err := c.Close(t.Context()); err != nil {
115 t.Fatal(err)
116 }
117 count := 0
118 for _, target := range starts {
119 if target.Path == waiting {
120 count++
121 if target.Scope != "global" {
122 t.Fatalf("dispatched stale scope: %+v", target)
123 }
124 }
125 }
126 if count != 1 {
127 t.Fatalf("pre-dispatch invalidation caused %d scans, want 1", count)
128 }
129 }
130
131 func TestMetadataQueuePreservesInvalidationAfterDispatch(t *testing.T) {
132 dir := t.TempDir()
133 c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, StartPaused: true})
134 if err != nil {
135 t.Fatal(err)
136 }
137 entered, release := make(chan struct{}), make(chan struct{})
138 var gate, slice sync.Once
139 t.Cleanup(func() { gate.Do(func() { close(release) }); _ = c.Close(context.Background()) })
140 c.testReconcileBatchHook = func(int) { slice.Do(func() { close(entered); <-release }) }
141 starts := 0
142 c.testReconcileStartHook = func(DirectoryTarget) { starts++ }
143 target := DirectoryTarget{Path: dir, Scope: "global"}
144 done, accepted := c.ScheduleReconcile(target)
145 if !accepted {
146 t.Fatal("initial discovery rejected")
147 }
148 c.ResumeDiscovery()
149 select {
150 case <-entered:
151 case <-t.Context().Done():
152 t.Fatal("scan never entered")
153 }
154 // The iterator has already reached EOF. A late file must be discovered by
155 // a follow-up, not swallowed by pre-dispatch coalescing.
156 path := filepath.Join(dir, "late.jsonl")
157 if err := os.WriteFile(path, []byte("unread body"), 0600); err != nil {
158 t.Fatal(err)
159 }
160 joined, accepted := c.ScheduleReconcile(target)
161 if !accepted || joined != done {
162 t.Fatal("late invalidation lost the shared completion signal")
163 }
164 gate.Do(func() { close(release) })
165 select {
166 case <-done:
167 case <-t.Context().Done():
168 t.Fatal("follow-up did not settle")
169 }
170 if _, found, err := c.GetSession(t.Context(), path); err != nil || !found {
171 t.Fatalf("late source was lost: %v %v", found, err)
172 }
173 if err := c.Close(t.Context()); err != nil {
174 t.Fatal(err)
175 }
176 if starts != 2 {
177 t.Fatalf("in-flight invalidation caused %d scans, want 2", starts)
178 }
179 }
180
180 lines GO