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