返回 DeepSeek-Reasonix
session_catalog_watch_convergence_test.go
根目录 / desktop / session_catalog_watch_convergence_test.go
1 package main
2
3 import (
4 "context"
5 "fmt"
6 "os"
7 "path/filepath"
8 "sync/atomic"
9 "testing"
10 "time"
11
12 "github.com/fsnotify/fsnotify"
13 "reasonix/internal/sessioncatalog"
14 )
15
16 type catalogWatchProbe struct {
17 catalog *sessioncatalog.Catalog
18 requested atomic.Int32
19 started atomic.Int32
20 completed atomic.Int32
21 }
22
23 // openCatalogWatchProbe opens an on-disk advisory catalog over dir. Every
24 // discovery "started" is one beginDirectoryScan, i.e. one scan_generation step.
25 func openCatalogWatchProbe(t *testing.T, dir string) (*catalogWatchProbe, sessioncatalog.DirectoryTarget) {
26 t.Helper()
27 probe := &catalogWatchProbe{}
28 catalog, err := sessioncatalog.Open(t.Context(), sessioncatalog.Options{
29 Path: filepath.Join(t.TempDir(), "v10.sqlite"), MetadataOnly: true, StartPaused: true,
30 OnDiscovery: func(event sessioncatalog.DiscoveryEvent) {
31 switch event.Phase {
32 case "requested":
33 probe.requested.Add(1)
34 case "started":
35 probe.started.Add(1)
36 case "completed":
37 probe.completed.Add(1)
38 }
39 },
40 })
41 if err != nil {
42 t.Fatal(err)
43 }
44 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
45 probe.catalog = catalog
46 target := sessioncatalog.DirectoryTarget{Path: dir, Scope: "project", WorkspaceRoot: filepath.Dir(dir)}
47 done, accepted := catalog.ScheduleReconcile(target)
48 if !accepted {
49 t.Fatal("initial discovery rejected")
50 }
51 catalog.ResumeDiscovery()
52 select {
53 case <-done:
54 case <-time.After(sessionCatalogTestDeadline):
55 t.Fatal("initial discovery did not finish")
56 }
57 return probe, target
58 }
59
60 func (p *catalogWatchProbe) settle(t *testing.T) {
61 t.Helper()
62 // Requests dispatch asynchronously, so settled means quiet for a while.
63 deadline := time.Now().Add(sessionCatalogTestDeadline)
64 quietSince, last := time.Now(), p.started.Load()
65 for {
66 started := p.started.Load()
67 if started != last || p.completed.Load() < started {
68 quietSince, last = time.Now(), started
69 } else if time.Since(quietSince) > 300*time.Millisecond {
70 return
71 }
72 if time.Now().After(deadline) {
73 t.Fatalf("discovery did not settle: started=%d completed=%d", started, p.completed.Load())
74 }
75 time.Sleep(5 * time.Millisecond)
76 }
77 }
78
79 // Reported state (#11693): an unchanged project directory of legacy rows
80 // (log_format=1, turns unknown, repair pending at zero attempts) beside files
81 // that are not transcripts and that other writers keep touching.
82 func writeLegacyCatalogDirectory(t *testing.T) string {
83 t.Helper()
84 dir := filepath.Join(t.TempDir(), "sessions")
85 if err := os.MkdirAll(filepath.Join(dir, "s0.checkpoints"), 0o755); err != nil {
86 t.Fatal(err)
87 }
88 for i := range 7 {
89 if err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("s%d.jsonl", i)), []byte("{\"role\":\"user\"}\n"), 0o600); err != nil {
90 t.Fatal(err)
91 }
92 }
93 return dir
94 }
95
96 func TestCatalogWatchNonSessionChurnDoesNotRescanUnchangedRoot(t *testing.T) {
97 dir := writeLegacyCatalogDirectory(t)
98 probe, target := openCatalogWatchProbe(t, dir)
99 record, found, err := probe.catalog.GetSession(t.Context(), filepath.Join(dir, "s0.jsonl"))
100 if err != nil || !found || record.LogFormat != 1 || record.TurnsState != sessioncatalog.TurnsUnknown {
101 t.Fatalf("fixture is not the reported legacy row state: %+v found=%v err=%v", record, found, err)
102 }
103 key := canonicalWorkspaceRoot(dir)
104 targets := map[string]sessioncatalog.DirectoryTarget{key: target}
105 watched, dirty := map[string]bool{key: true}, map[string]bool{}
106 pacer := &catalogRootPacer{}
107 baseRequested, baseStarted := probe.requested.Load(), probe.started.Load()
108 churn := []fsnotify.Event{
109 {Name: filepath.Join(key, ".titles.json.tmp"), Op: fsnotify.Create},
110 {Name: filepath.Join(key, ".titles.json.tmp"), Op: fsnotify.Rename},
111 {Name: filepath.Join(key, ".display.json"), Op: fsnotify.Write},
112 {Name: filepath.Join(key, "s0.wire.jsonl"), Op: fsnotify.Write},
113 {Name: filepath.Join(key, "s0.context.json"), Op: fsnotify.Write},
114 {Name: filepath.Join(key, "s0.checkpoints"), Op: fsnotify.Write},
115 {Name: filepath.Join(key, "writer.lock"), Op: fsnotify.Create},
116 {Name: filepath.Join(key, "writer.lock"), Op: fsnotify.Remove},
117 }
118 now := time.Now()
119 for cycle := range 20 {
120 for _, event := range churn {
121 admitCatalogWatchEvent(probe.catalog, event, targets, watched, dirty, nil, false, nil)
122 }
123 admitCatalogWatchBatch(t.Context(), probe.catalog, targets, dirty, false, pacer, now.Add(time.Duration(cycle)*3*time.Second))
124 }
125 probe.settle(t)
126 if requested, started := probe.requested.Load()-baseRequested, probe.started.Load()-baseStarted; requested != 0 || started != 0 || len(dirty) != 0 {
127 t.Fatalf("unchanged root rescanned on non-session churn: requested=%d scans=%d dirty=%v", requested, started, dirty)
128 }
129 }
130
131 func TestCatalogWatchSessionSourcesStillReachTheCatalog(t *testing.T) {
132 dir := writeLegacyCatalogDirectory(t)
133 probe, target := openCatalogWatchProbe(t, dir)
134 key := canonicalWorkspaceRoot(dir)
135 targets := map[string]sessioncatalog.DirectoryTarget{key: target}
136 watched, dirty := map[string]bool{key: true}, map[string]bool{}
137 for _, name := range []string{"s1.jsonl", "s1.jsonl.meta", "s1.events.jsonl"} {
138 admitCatalogWatchEvent(probe.catalog, fsnotify.Event{Name: filepath.Join(key, name), Op: fsnotify.Write}, targets, watched, dirty, nil, false, nil)
139 if len(dirty) != 0 {
140 t.Fatalf("%s write scheduled a root scan instead of an exact index: %v", name, dirty)
141 }
142 }
143 admitCatalogWatchEvent(probe.catalog, fsnotify.Event{Name: filepath.Join(key, "s1.events.jsonl"), Op: fsnotify.Write}, targets, watched, dirty, nil, true, nil)
144 if !dirty[key] {
145 t.Fatal("pre-discovery session event lost its root invalidation")
146 }
147 clear(dirty)
148 admitCatalogWatchEvent(probe.catalog, fsnotify.Event{Name: filepath.Join(key, "s2.jsonl"), Op: fsnotify.Remove}, targets, watched, dirty, nil, false, nil)
149 if !dirty[key] {
150 t.Fatal("transcript removal must reconcile the root to mark the row missing")
151 }
152 }
153
154 func TestCatalogWatchRootPacingConvergesUnderPersistentTrigger(t *testing.T) {
155 dir := writeLegacyCatalogDirectory(t)
156 probe, target := openCatalogWatchProbe(t, dir)
157 key := canonicalWorkspaceRoot(dir)
158 targets := map[string]sessioncatalog.DirectoryTarget{key: target}
159 dirty := map[string]bool{}
160 pacer := &catalogRootPacer{}
161 baseRequested := probe.requested.Load()
162 start := time.Now()
163 // A root-level invalidation every 3s for ten minutes: the shape of the
164 // reported flapping row, whatever produces it.
165 var last time.Time
166 for tick := range 200 {
167 now := start.Add(time.Duration(tick) * 3 * time.Second)
168 dirty[key] = true
169 before := probe.requested.Load()
170 next := admitCatalogWatchBatch(t.Context(), probe.catalog, targets, dirty, false, pacer, now)
171 if probe.requested.Load() != before {
172 last = now
173 } else if !dirty[key] || next.IsZero() || !next.After(now) {
174 t.Fatalf("deferred root lost its invalidation or due time: dirty=%v next=%v", dirty, next)
175 }
176 }
177 probe.settle(t)
178 admissions := probe.requested.Load() - baseRequested
179 if admissions > 20 {
180 t.Fatalf("persistent trigger was not backed off: %d full scans in ten minutes", admissions)
181 }
182 if admissions < 5 {
183 t.Fatalf("pacing starved a dirty root: %d full scans in ten minutes", admissions)
184 }
185 // A root quiet for longer than the backoff window is admitted at once again.
186 quiet := last.Add(3 * catalogRootPaceMax)
187 dirty[key] = true
188 before := probe.requested.Load()
189 admitCatalogWatchBatch(t.Context(), probe.catalog, targets, dirty, false, pacer, quiet)
190 if probe.requested.Load() == before || dirty[key] {
191 t.Fatal("a settled root must be admitted without inherited backoff")
192 }
193 }
194
195 func TestCatalogWatchSessionRemovalBypassesChurnBackoff(t *testing.T) {
196 dir := writeLegacyCatalogDirectory(t)
197 probe, target := openCatalogWatchProbe(t, dir)
198 key := canonicalWorkspaceRoot(dir)
199 targets := map[string]sessioncatalog.DirectoryTarget{key: target}
200 watched, dirty := map[string]bool{key: true}, map[string]bool{}
201 pacer := &catalogRootPacer{}
202 now := time.Now()
203 // Drive the root to the backoff cap with root-level invalidations.
204 for range 40 {
205 now = now.Add(3 * time.Second)
206 dirty[key] = true
207 admitCatalogWatchBatch(t.Context(), probe.catalog, targets, dirty, false, pacer, now)
208 }
209 if due := pacer.due(key); due.Sub(now) > catalogRootPaceMax {
210 t.Fatalf("root-level change deferred past the documented maximum: %v", due.Sub(now))
211 }
212 if dirty[key] && !now.Add(catalogRootPaceMax).After(pacer.due(key)) {
213 t.Fatal("deferred root has no due time within the maximum")
214 }
215 if err := os.Remove(filepath.Join(dir, "s3.jsonl")); err != nil {
216 t.Fatal(err)
217 }
218 admitCatalogWatchEvent(probe.catalog, fsnotify.Event{Name: filepath.Join(key, "s3.jsonl"), Op: fsnotify.Remove}, targets, watched, dirty, nil, false, pacer)
219 before := probe.requested.Load()
220 now = now.Add(time.Millisecond)
221 admitCatalogWatchBatch(t.Context(), probe.catalog, targets, dirty, false, pacer, now)
222 if probe.requested.Load() == before || dirty[key] {
223 t.Fatal("a removed transcript waited behind unrelated churn backoff")
224 }
225 probe.settle(t)
226 if record, found, err := probe.catalog.GetSession(t.Context(), filepath.Join(dir, "s3.jsonl")); err != nil || (found && record.Health != sessioncatalog.HealthMissing) {
227 t.Fatalf("removed transcript still listed as present: %+v found=%v err=%v", record, found, err)
228 }
229 // The expedited scan must not reset or grow the churn backoff either way.
230 dirty[key] = true
231 if next := admitCatalogWatchBatch(t.Context(), probe.catalog, targets, dirty, false, pacer, now.Add(time.Second)); next.IsZero() || next.Sub(now) > catalogRootPaceMax {
232 t.Fatalf("churn after an expedited scan lost its bounded backoff: next=%v", next)
233 }
234 }
235
235 lines GO