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