返回 DeepSeek-Reasonix
session_catalog_watch_backend_test.go
根目录 / desktop / session_catalog_watch_backend_test.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "os"
8 "path/filepath"
9 "sync/atomic"
10 "testing"
11 "time"
12
13 "github.com/fsnotify/fsnotify"
14 "reasonix/internal/sessioncatalog"
15 )
16
17 type recoveringCatalogWatch struct {
18 workspaceWatcher
19 fail bool
20 adds int
21 removes []string
22 }
23
24 type catchUpCatalogWatcher struct {
25 workspaceWatcher
26 catchUp func()
27 }
28
29 func (w *catchUpCatalogWatcher) CatchUp() { w.catchUp() }
30
31 type saturatedCatalogPathQueue struct {
32 attempts int
33 }
34
35 func (q *saturatedCatalogPathQueue) TryRequestIndexSession(sessioncatalog.DirectoryTarget, string) bool {
36 q.attempts++
37 return false
38 }
39
40 func TestCatalogWatchPathOverflowCoalescesAtEventOwner(t *testing.T) {
41 root := canonicalWorkspaceRoot(t.TempDir())
42 target := sessioncatalog.DirectoryTarget{Path: root, Scope: "global"}
43 targets := map[string]sessioncatalog.DirectoryTarget{root: target}
44 watching, dirty := map[string]bool{root: true}, map[string]bool{}
45 queue := &saturatedCatalogPathQueue{}
46 for i := range 4096 {
47 admitCatalogWatchEvent(queue, fsnotify.Event{Name: filepath.Join(root, fmt.Sprintf("%d.jsonl", i)), Op: fsnotify.Write}, targets, watching, dirty, nil, false, nil)
48 }
49 if queue.attempts != 4096 || len(dirty) != 1 || !dirty[root] {
50 t.Fatalf("overflow lost its coalesced root: attempts=%d dirty=%v", queue.attempts, dirty)
51 }
52 }
53
54 func TestCatalogWatchInitialBacklogCoalescesBeforeDiscovery(t *testing.T) {
55 dir, other := t.TempDir(), t.TempDir()
56 path := filepath.Join(dir, "old.jsonl")
57 if err := os.WriteFile(path, []byte("unreadable body"), 0600); err != nil {
58 t.Fatal(err)
59 }
60 var requested, started atomic.Int32
61 catalog, err := sessioncatalog.Open(t.Context(), sessioncatalog.Options{InMemory: true, MetadataOnly: true, StartPaused: true,
62 OnDiscovery: func(event sessioncatalog.DiscoveryEvent) {
63 if event.Phase == "requested" {
64 requested.Add(1)
65 }
66 if event.Phase == "started" {
67 started.Add(1)
68 }
69 },
70 })
71 if err != nil {
72 t.Fatal(err)
73 }
74 defer catalog.Close(context.Background())
75 key, otherKey := canonicalWorkspaceRoot(dir), canonicalWorkspaceRoot(other)
76 target := sessioncatalog.DirectoryTarget{Path: dir, Scope: "global"}
77 targets := map[string]sessioncatalog.DirectoryTarget{key: target, otherKey: {Path: other, Scope: "global"}}
78 watched, dirty := map[string]bool{key: true, otherKey: true}, map[string]bool{}
79 // Force more notices than the exact-path queue can retain. Before the
80 // first iterator, all of them belong to the same root invalidation.
81 for range 4096 {
82 admitCatalogWatchEvent(catalog, fsnotify.Event{Name: filepath.Join(key, "old.jsonl"), Op: fsnotify.Create}, targets, watched, dirty, nil, true, nil)
83 }
84 events, failures := make(chan fsnotify.Event, 2), make(chan error, 1)
85 barrierReached := false
86 watcher := &catchUpCatalogWatcher{catchUp: func() {
87 barrierReached = true
88 events <- fsnotify.Event{Name: filepath.Join(key, "old.jsonl"), Op: fsnotify.Write}
89 failures <- fsnotify.ErrEventOverflow
90 }}
91 catchUpCatalogWatch(watcher, events, failures, targets, watched, dirty)
92 if !barrierReached || len(events) != 0 || len(failures) != 0 || !dirty[key] || !dirty[otherKey] {
93 t.Fatalf("initial barrier lost notifications: reached=%v dirty=%v", barrierReached, dirty)
94 }
95 if requested.Load() != 0 {
96 t.Fatalf("pre-discovery backlog admitted %d redundant reconciles", requested.Load())
97 }
98 done, accepted := catalog.ScheduleReconcile(target)
99 if !accepted {
100 t.Fatal("initial discovery rejected")
101 }
102 catalog.ResumeDiscovery()
103 select {
104 case <-done:
105 case <-t.Context().Done():
106 t.Fatal("initial discovery did not finish")
107 }
108 if started.Load() != 1 {
109 t.Fatalf("initial backlog started %d scans", started.Load())
110 }
111 if _, found, err := catalog.GetSession(t.Context(), path); err != nil || !found {
112 t.Fatalf("first scan omitted historical source: found=%v err=%v", found, err)
113 }
114 clear(dirty)
115 admitCatalogWatchEvent(catalog, fsnotify.Event{Name: key, Op: fsnotify.Write}, targets, watched, dirty, nil, false, nil)
116 if !dirty[key] {
117 t.Fatal("post-discovery directory invalidation was swallowed")
118 }
119 }
120
121 func (w *recoveringCatalogWatch) Remove(path string) error {
122 w.removes = append(w.removes, path)
123 return nil
124 }
125
126 func (w *recoveringCatalogWatch) Add(string, bool) error {
127 w.adds++
128 if w.fail {
129 return errors.New("watch temporarily unavailable")
130 }
131 return nil
132 }
133
134 func TestCatalogWatchRecoveryReconcilesOnlyTheUnobservedInterval(t *testing.T) {
135 root := canonicalWorkspaceRoot(t.TempDir())
136 targets := []sessioncatalog.DirectoryTarget{{Path: root, Scope: "global"}}
137 watched, dirty := map[string]bool{}, map[string]bool{}
138 watcher := &recoveringCatalogWatch{fail: true}
139 current := refreshCatalogWatchTargets(watcher, nil, targets, watched, dirty, nil)
140 if !dirty[root] || watched[root] {
141 t.Fatal("failed initial watch omitted initial discovery")
142 }
143 clear(dirty)
144 current = refreshCatalogWatchTargets(watcher, current, targets, watched, dirty, nil)
145 if len(dirty) != 0 || watcher.adds != 2 {
146 t.Fatal("watch retry must not schedule full discovery on every metadata tick")
147 }
148 watcher.fail = false
149 current = refreshCatalogWatchTargets(watcher, current, targets, watched, dirty, nil)
150 if !dirty[root] || !watched[root] {
151 t.Fatal("recovered watch failed to reconcile its observation gap")
152 }
153 clear(dirty)
154 refreshCatalogWatchTargets(watcher, current, targets, watched, dirty, nil)
155 if len(dirty) != 0 || watcher.adds != 3 {
156 t.Fatal("settled watch restarted registration or discovery")
157 }
158 }
159
160 func TestCatalogWatchCanonicalEventKeepsRegisteredAccessPath(t *testing.T) {
161 dir := t.TempDir()
162 source := filepath.Join(dir, "session.jsonl")
163 if err := os.WriteFile(source, []byte("not a valid transcript"), 0600); err != nil {
164 t.Fatal(err)
165 }
166 updated := make(chan struct{}, 1)
167 catalog, err := sessioncatalog.Open(t.Context(), sessioncatalog.Options{InMemory: true, MetadataOnly: true, StartPaused: true,
168 OnRevision: func(uint64, []string, string) {
169 select {
170 case updated <- struct{}{}:
171 default:
172 }
173 },
174 })
175 if err != nil {
176 t.Fatal(err)
177 }
178 defer catalog.Close(context.Background())
179 watched, dirty := map[string]bool{}, map[string]bool{}
180 targets := refreshCatalogWatchTargets(nil, nil, []sessioncatalog.DirectoryTarget{{Path: dir, Scope: "global"}}, watched, dirty, nil)
181 clear(dirty)
182 key := canonicalWorkspaceRoot(dir)
183 admitCatalogWatchEvent(catalog, fsnotify.Event{Name: filepath.Join(key, "session.jsonl.meta"), Op: fsnotify.Write}, targets, watched, dirty, nil, false, nil)
184 if len(dirty) != 0 {
185 t.Fatalf("exact metadata event scheduled a root scan: %v", dirty)
186 }
187 select {
188 case <-updated:
189 case <-time.After(sessionCatalogTestDeadline):
190 t.Fatal("canonical event failed to update the exact path")
191 }
192 record, exists, err := catalog.GetSession(t.Context(), source)
193 if err != nil || !exists || record.Path != source {
194 t.Fatalf("watch event changed access identity: %+v %v %v", record, exists, err)
195 }
196 // A removed root must drop its watch and schedule reconciliation even
197 // when there is no current filesystem object to canonicalize.
198 watched[key] = true
199 watcher := &recoveringCatalogWatch{}
200 admitCatalogWatchEvent(catalog, fsnotify.Event{Name: key, Op: fsnotify.Remove}, targets, watched, dirty, watcher, false, nil)
201 if watched[key] || !dirty[key] {
202 t.Fatalf("root removal lost invalidation: watched=%v dirty=%v", watched, dirty)
203 }
204 if len(watcher.removes) != 1 || watcher.removes[0] != key {
205 t.Fatalf("root removal retained the native subscription: %v", watcher.removes)
206 }
207 clear(dirty)
208 refreshCatalogWatchTargets(watcher, targets, []sessioncatalog.DirectoryTarget{{Path: dir, Scope: "global"}}, watched, dirty, nil)
209 if watcher.adds != 1 || !watched[key] || !dirty[key] {
210 t.Fatal("replacement directory did not re-register and reconcile")
211 }
212 }
213
213 lines GO