| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "log/slog" |
| 7 | "path/filepath" |
| 8 | "sort" |
| 9 | "strings" |
| 10 | "time" |
| 11 | |
| 12 | "github.com/fsnotify/fsnotify" |
| 13 | "reasonix/internal/history" |
| 14 | "reasonix/internal/sessioncatalog" |
| 15 | "reasonix/internal/store" |
| 16 | ) |
| 17 | |
| 18 | // Directory notifications admit only dirty roots. The periodic audit remains |
| 19 | // authoritative after dropped notifications and on unsupported filesystems. |
| 20 | func (a *App) watchSessionCatalog(ctx context.Context, catalog *sessioncatalog.Catalog, metadataRequests <-chan struct{}, onAdmitted func()) { |
| 21 | ctx, cancel := context.WithCancel(ctx) |
| 22 | defer cancel() |
| 23 | // Reuse the platform directory backend. On macOS the generic kqueue |
| 24 | // adapter enumerates children and opens a descriptor for every file. |
| 25 | watcher, _ := newWorkspaceWatcher() |
| 26 | var events <-chan fsnotify.Event |
| 27 | var failures <-chan error |
| 28 | if watcher != nil { |
| 29 | defer watcher.Close() |
| 30 | events, failures = watcher.Events(), watcher.Errors() |
| 31 | } |
| 32 | targets := map[string]sessioncatalog.DirectoryTarget{} |
| 33 | watched := map[string]bool{} |
| 34 | dirty := map[string]bool{} |
| 35 | discoveryPending := false |
| 36 | pacer, batch := &catalogRootPacer{}, &catalogWatchBatchTimer{} |
| 37 | defer batch.stop() |
| 38 | armBatchAfter := func(wait time.Duration) { |
| 39 | if len(dirty) > 0 || discoveryPending { |
| 40 | batch.arm(wait) |
| 41 | } |
| 42 | } |
| 43 | armBatch := func() { armBatchAfter(catalogWatchBatchDelay) } |
| 44 | // A slow registry projection must not stop receiving filesystem events. |
| 45 | // One worker coalesces refreshes and is joined before this owner returns. |
| 46 | refreshMetadata, metadataDone := startCatalogMetadataRefresh(ctx, func(ctx context.Context) { |
| 47 | if err := a.syncSessionCatalogMetadataBounded(ctx, catalog); err != nil && !errors.Is(err, context.Canceled) { |
| 48 | slog.Debug("desktop: refresh session catalog metadata", "err", err) |
| 49 | } |
| 50 | }) |
| 51 | defer func() { cancel(); <-metadataDone }() |
| 52 | refreshTargets := func() { |
| 53 | targets = refreshCatalogWatchTargets(watcher, targets, a.sessionCatalogTargets(), watched, dirty, pacer) |
| 54 | } |
| 55 | refreshTargets() |
| 56 | restored := a.tabsRestoredSignal() |
| 57 | admitted := false |
| 58 | metadata := time.NewTicker(30 * time.Second) |
| 59 | audit := time.NewTicker(5 * time.Minute) |
| 60 | defer metadata.Stop() |
| 61 | defer audit.Stop() |
| 62 | maintenance := 0 |
| 63 | for { |
| 64 | select { |
| 65 | case <-ctx.Done(): |
| 66 | return |
| 67 | case <-catalog.Invalidated(): |
| 68 | return |
| 69 | case <-metadataRequests: |
| 70 | if admitted { |
| 71 | refreshMetadata() |
| 72 | } |
| 73 | case <-restored: |
| 74 | restored = nil |
| 75 | // Restored identities can reveal additional roots. Register those |
| 76 | // watches before allowing their queued discovery to run. |
| 77 | refreshTargets() |
| 78 | history.RegisterCatalogRoots(historyCatalogRoots(a.sessionCatalogTargets())) |
| 79 | a.indexRestoredSessionPaths(ctx, catalog) |
| 80 | // Query/UI admission must not wait for journal writes. The first |
| 81 | // batch admits watched roots before releasing restored scan jobs. |
| 82 | discoveryPending = true |
| 83 | admitted = true |
| 84 | onAdmitted() |
| 85 | // First discovery covers the notification backlog before this barrier. |
| 86 | // Per-file admission would saturate the queue and schedule a second |
| 87 | // whole-root pass. Queries are already admitted. |
| 88 | catchUpCatalogWatch(watcher, events, failures, targets, watched, dirty) |
| 89 | refreshMetadata() |
| 90 | a.requestHistoricalCatalogWithContext(ctx) |
| 91 | armBatch() |
| 92 | case event, ok := <-events: |
| 93 | if !ok { |
| 94 | events = nil |
| 95 | clear(watched) |
| 96 | continue |
| 97 | } |
| 98 | admitCatalogWatchEvent(catalog, event, targets, watched, dirty, watcher, !admitted || discoveryPending, pacer) |
| 99 | armBatch() |
| 100 | case _, ok := <-failures: |
| 101 | if !ok { |
| 102 | failures = nil |
| 103 | } |
| 104 | for key := range targets { |
| 105 | dirty[key] = true |
| 106 | } |
| 107 | armBatch() |
| 108 | case <-metadata.C: |
| 109 | if admitted { |
| 110 | refreshMetadata() |
| 111 | } |
| 112 | refreshTargets() |
| 113 | armBatch() |
| 114 | case <-audit.C: |
| 115 | refreshTargets() |
| 116 | keys := make([]string, 0, len(targets)) |
| 117 | for key := range targets { |
| 118 | keys = append(keys, key) |
| 119 | } |
| 120 | sort.Strings(keys) |
| 121 | if len(keys) > 0 { |
| 122 | dirty[keys[maintenance%len(keys)]] = true |
| 123 | maintenance++ |
| 124 | } |
| 125 | armBatch() |
| 126 | case <-batch.ready: |
| 127 | batch.ready = nil |
| 128 | next := admitCatalogWatchBatch(ctx, catalog, targets, dirty, discoveryPending, pacer, time.Now()) |
| 129 | discoveryPending = false |
| 130 | // Deferred roots and failed journal writes keep their invalidation. |
| 131 | armBatchAfter(max(catalogWatchBatchDelay, time.Until(next))) |
| 132 | } |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | func catchUpCatalogWatch(watcher workspaceWatcher, events <-chan fsnotify.Event, failures <-chan error, targets map[string]sessioncatalog.DirectoryTarget, watched, dirty map[string]bool) { |
| 137 | if barrier, ok := watcher.(interface{ CatchUp() }); ok { |
| 138 | barrier.CatchUp() |
| 139 | } |
| 140 | // Drain only the notifications already queued at the barrier. New events |
| 141 | // still go through the ordinary loop; a continuously written directory |
| 142 | // cannot hold discovery hostage by keeping this drain alive indefinitely. |
| 143 | for count := len(events); count > 0; count-- { |
| 144 | if event, ok := <-events; ok { |
| 145 | admitCatalogWatchEvent(nil, event, targets, watched, dirty, watcher, true, nil) |
| 146 | } |
| 147 | } |
| 148 | for count := len(failures); count > 0; count-- { |
| 149 | if _, ok := <-failures; ok { |
| 150 | for key := range targets { |
| 151 | dirty[key] = true |
| 152 | } |
| 153 | } |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | type catalogWatchPathQueue interface { |
| 158 | TryRequestIndexSession(sessioncatalog.DirectoryTarget, string) bool |
| 159 | } |
| 160 | |
| 161 | func admitCatalogWatchEvent(catalog catalogWatchPathQueue, event fsnotify.Event, targets map[string]sessioncatalog.DirectoryTarget, watched, dirty map[string]bool, watcher workspaceWatcher, initialDiscovery bool, pacer *catalogRootPacer) { |
| 162 | key := filepath.Clean(filepath.Dir(event.Name)) |
| 163 | // A transcript/sidecar write invalidates one session, not its root. |
| 164 | if target, exists := targets[key]; exists { |
| 165 | path := catalogSessionPathForEvent(event.Name) |
| 166 | if path != "" && event.Op&(fsnotify.Remove|fsnotify.Rename) == 0 { |
| 167 | if initialDiscovery { |
| 168 | dirty[key] = true |
| 169 | return |
| 170 | } |
| 171 | // Platform events use the canonical watch path; retain the registered |
| 172 | // access spelling when publishing the exact session identity. |
| 173 | if !catalog.TryRequestIndexSession(target, filepath.Join(target.Path, filepath.Base(path))) { |
| 174 | // The event loop owns overflow, just like an unscoped directory |
| 175 | // event. Journal one coalesced root in the next batch and retain |
| 176 | // it if admission fails, rather than blocking on every notice. |
| 177 | dirty[key] = true |
| 178 | } |
| 179 | return |
| 180 | } |
| 181 | if path != "" { |
| 182 | pacer.expedite(key) |
| 183 | } else if _, nested := targets[filepath.Clean(event.Name)]; !nested { |
| 184 | // A metadata scan reads nothing else in the root: locks, temp files, |
| 185 | // other sidecars and subdirectories cannot change its projection. |
| 186 | return |
| 187 | } |
| 188 | } |
| 189 | if _, exists := targets[filepath.Clean(event.Name)]; exists { |
| 190 | key = filepath.Clean(event.Name) |
| 191 | if event.Op&(fsnotify.Remove|fsnotify.Rename) != 0 { |
| 192 | // Native backends retain subscriptions until explicitly removed. |
| 193 | // Drop the backend entry as well as our admission flag so Add can |
| 194 | // watch a newly created directory at this same path. |
| 195 | if watcher != nil { |
| 196 | _ = watcher.Remove(key) |
| 197 | } |
| 198 | delete(watched, key) |
| 199 | } |
| 200 | } |
| 201 | if _, exists := targets[key]; exists { |
| 202 | dirty[key] = true |
| 203 | } |
| 204 | } |
| 205 | |
| 206 | // catalogSessionPathForEvent names the transcript whose catalog row an event |
| 207 | // can change: the transcript, its branch sidecar, or its event log. |
| 208 | func catalogSessionPathForEvent(path string) string { |
| 209 | if store.IsSessionTranscriptName(filepath.Base(path)) { |
| 210 | return path |
| 211 | } |
| 212 | if filepath.Ext(path) == ".meta" { |
| 213 | candidate := path[:len(path)-len(".meta")] |
| 214 | if store.IsSessionTranscriptName(filepath.Base(candidate)) { |
| 215 | return candidate |
| 216 | } |
| 217 | } |
| 218 | if stem, ok := strings.CutSuffix(path, ".events.jsonl"); ok { |
| 219 | candidate := stem + ".jsonl" |
| 220 | if store.IsSessionTranscriptName(filepath.Base(candidate)) && store.SessionEventLog(candidate) == path { |
| 221 | return candidate |
| 222 | } |
| 223 | } |
| 224 | return "" |
| 225 | } |
| 226 | |
| 227 | func refreshCatalogWatchTargets(watcher workspaceWatcher, current map[string]sessioncatalog.DirectoryTarget, targets []sessioncatalog.DirectoryTarget, watched, dirty map[string]bool, pacer *catalogRootPacer) map[string]sessioncatalog.DirectoryTarget { |
| 228 | next := map[string]sessioncatalog.DirectoryTarget{} |
| 229 | for _, target := range targets { |
| 230 | key := canonicalWorkspaceRoot(target.Path) |
| 231 | next[key] = target |
| 232 | _, known := current[key] |
| 233 | if !known { |
| 234 | dirty[key] = true |
| 235 | } |
| 236 | if watcher != nil && !watched[key] { |
| 237 | watched[key] = watcher.Add(key, false) == nil |
| 238 | if watched[key] && known { |
| 239 | // Recovered watching must reconcile the unobserved interval. |
| 240 | dirty[key] = true |
| 241 | } |
| 242 | } |
| 243 | // Unavailable watches use the rotating audit, not a full-root scan |
| 244 | // on every metadata refresh. Registration itself may still retry. |
| 245 | } |
| 246 | for key := range current { |
| 247 | if _, exists := next[key]; !exists { |
| 248 | if watcher != nil && watched[key] { |
| 249 | _ = watcher.Remove(key) |
| 250 | } |
| 251 | delete(watched, key) |
| 252 | delete(dirty, key) |
| 253 | pacer.forget(key) |
| 254 | } |
| 255 | } |
| 256 | return next |
| 257 | } |
| 258 | |
| 259 | // admitCatalogWatchBatch returns when the earliest root it deferred falls due, |
| 260 | // or the zero time when nothing was deferred. Failed admissions count as |
| 261 | // attempts, so a journal that keeps failing is retried with backoff too. |
| 262 | func admitCatalogWatchBatch(ctx context.Context, catalog *sessioncatalog.Catalog, targets map[string]sessioncatalog.DirectoryTarget, dirty map[string]bool, discoveryPending bool, pacer *catalogRootPacer, now time.Time) (next time.Time) { |
| 263 | initial := []sessioncatalog.DirectoryTarget{} |
| 264 | for key := range dirty { |
| 265 | target, exists := targets[key] |
| 266 | if !exists || ctx.Err() != nil { |
| 267 | delete(dirty, key) |
| 268 | continue |
| 269 | } |
| 270 | if discoveryPending { |
| 271 | delete(dirty, key) |
| 272 | initial = append(initial, target) |
| 273 | pacer.admitted(key, now) |
| 274 | continue |
| 275 | } |
| 276 | if due := pacer.due(key); now.Before(due) { |
| 277 | if next.IsZero() || due.Before(next) { |
| 278 | next = due |
| 279 | } |
| 280 | continue |
| 281 | } |
| 282 | delete(dirty, key) |
| 283 | pacer.admitted(key, now) |
| 284 | if !catalog.RequestReconcile(target) { |
| 285 | dirty[key] = true |
| 286 | if due := pacer.due(key); next.IsZero() || due.Before(next) { |
| 287 | next = due |
| 288 | } |
| 289 | } |
| 290 | } |
| 291 | if discoveryPending { |
| 292 | for _, target := range catalog.ResumeDiscovery(initial...) { |
| 293 | dirty[canonicalWorkspaceRoot(target.Path)] = true |
| 294 | } |
| 295 | } |
| 296 | return next |
| 297 | } |
| 298 |