| 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 |