返回 DeepSeek-Reasonix
session_catalog_watch.go
根目录 / desktop / session_catalog_watch.go
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
298 lines GO