返回 DeepSeek-Reasonix
index_queue.go
根目录 / internal / sessioncatalog / index_queue.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "path/filepath"
6 "time"
7 )
8
9 // RequestIndexSession coalesces an authoritative session write by path. A full
10 // queue journals a root invalidation before returning; notification consumers
11 // must use TryRequestIndexSession and coalesce overflow at their own boundary.
12 func (c *Catalog) RequestIndexSession(target DirectoryTarget, path string) bool {
13 if c == nil || cleanCatalogAccessPath(path) == "" {
14 return false
15 }
16 if c.TryRequestIndexSession(target, path) {
17 return true
18 }
19 if target.Path == "" {
20 target.Path = filepath.Dir(cleanCatalogAccessPath(path))
21 }
22 // Losing a wake-up must not lose the authoritative write. The root
23 // queue coalesces overflow and retains it until reconciliation.
24 c.RequestReconcile(target)
25 return false
26 }
27
28 // TryRequestIndexSession never waits for filesystem or database work. On false
29 // the caller owns the root invalidation, including retrying failed journal
30 // admission. This lets filesystem consumers drain bursts without doing one
31 // synchronous journal transaction for every overflowing file notification.
32 func (c *Catalog) TryRequestIndexSession(target DirectoryTarget, path string) bool {
33 if c == nil {
34 return false
35 }
36 path = cleanCatalogAccessPath(path)
37 if path == "" {
38 return false
39 }
40 target.Path = cleanCatalogAccessPath(target.Path)
41 if target.Path == "" {
42 target.Path = filepath.Dir(path)
43 }
44 key := queuePathKey(path)
45 request := sessionPathRequest{target: target, path: path, queueKey: key, sequence: c.mutationSeq.Add(1)}
46 c.pathQueueMu.Lock()
47 if _, loaded := c.pathQueued.Load(key); loaded {
48 c.pathQueued.Store(key, request)
49 c.pathQueueMu.Unlock()
50 return true
51 }
52 c.pathQueued.Store(key, request)
53 select {
54 case c.pathCh <- request:
55 c.pathQueueMu.Unlock()
56 return true
57 case <-c.stop:
58 c.pathQueued.Delete(key)
59 c.pathQueueMu.Unlock()
60 return false
61 default:
62 c.pathQueued.Delete(key)
63 c.pathQueueMu.Unlock()
64 return false
65 }
66 }
67
68 func (c *Catalog) sessionPathLoop() {
69 defer c.workers.Done()
70 for {
71 select {
72 case token := <-c.pathCh:
73 c.pathQueueMu.Lock()
74 queued, ok := c.pathQueued.LoadAndDelete(token.queueKey)
75 c.pathQueueMu.Unlock()
76 if !ok {
77 continue
78 }
79 request := queued.(sessionPathRequest)
80 ctx, cancel := context.WithTimeout(c.workerCtx, 30*time.Second)
81 _ = c.indexSessionPath(ctx, request.target, request.path, request.sequence)
82 cancel()
83 case <-c.stop:
84 return
85 }
86 }
87 }
88
88 lines GO