返回 DeepSeek-Reasonix
reconcile_queue.go
根目录 / internal / sessioncatalog / reconcile_queue.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "runtime"
6 "strings"
7 "time"
8 )
9
10 // RequestReconcile makes the channel a wake signal while the maps retain the
11 // newest target. Session saves never wait for catalog work.
12 func (c *Catalog) RequestReconcile(target DirectoryTarget) bool {
13 _, accepted := c.ScheduleReconcile(target)
14 return accepted
15 }
16
17 // ScheduleReconcile joins the catalog's single discovery owner. The returned
18 // signal settles after all coalesced source changes have been visited; callers
19 // must also select their own cancellation when waiting during shutdown.
20 func (c *Catalog) ScheduleReconcile(target DirectoryTarget) (<-chan struct{}, bool) {
21 if c == nil || strings.TrimSpace(target.Path) == "" {
22 return nil, false
23 }
24 target.Path = cleanCatalogAccessPath(target.Path)
25 key := queuePathKey(target.Path)
26 if key == "" {
27 return nil, false
28 }
29 target.mutationSeq = c.mutationSeq.Add(1)
30 if c.opts.OnDiscovery != nil {
31 // Keep the function symbol, never runtime file names or stack arguments.
32 // Skip only the two public queue wrappers to identify the real owner.
33 var pcs [8]uintptr
34 n := runtime.Callers(2, pcs[:])
35 frames := runtime.CallersFrames(pcs[:n])
36 for {
37 frame, more := frames.Next()
38 if !strings.HasSuffix(frame.Function, ".(*Catalog).RequestReconcile") && !strings.HasSuffix(frame.Function, ".(*Catalog).RequestIndexSession") {
39 c.observeDiscovery(target, "requested", frame.Function, "")
40 break
41 }
42 if !more {
43 break
44 }
45 }
46 }
47 if c.opts.MetadataOnly {
48 // A saturated path queue has already committed its authoritative save.
49 // Persist the root invalidation before acknowledging maintenance so a
50 // crash cannot turn an overflow into a permanently stale ready catalog.
51 if err := c.persistReconcileTarget(target); err != nil {
52 return nil, false
53 }
54 }
55 c.reconcileDirtyMu.Lock()
56 defer c.reconcileDirtyMu.Unlock()
57 select {
58 case <-c.stop:
59 return nil, false
60 default:
61 }
62 if c.reconcileDone == nil {
63 c.reconcileDone = map[string]chan struct{}{}
64 }
65 done := c.reconcileDone[key]
66 if done == nil {
67 done = make(chan struct{})
68 c.reconcileDone[key] = done
69 }
70 if queued, loaded := c.reconcileQueued.Load(key); loaded {
71 target = newestReconcileTarget(queued.(DirectoryTarget), target)
72 c.reconcileDirty[key] = target
73 c.reconcileQueued.Store(key, target)
74 return done, true
75 }
76 c.reconcileQueued.Store(key, target)
77 select {
78 case c.reconcileCh <- target:
79 return done, true
80 default:
81 c.reconcileDirty[key] = target
82 return done, true
83 }
84 }
85
86 func (c *Catalog) markReconcileDirty(target DirectoryTarget) {
87 key := queuePathKey(target.Path)
88 c.reconcileDirtyMu.Lock()
89 if queued, ok := c.reconcileQueued.Load(key); ok {
90 target = newestReconcileTarget(queued.(DirectoryTarget), target)
91 }
92 if dirty, ok := c.reconcileDirty[key]; ok {
93 target = newestReconcileTarget(dirty, target)
94 }
95 c.reconcileDirty[key] = target
96 c.reconcileQueued.Store(key, target)
97 c.reconcileDirtyMu.Unlock()
98 }
99
100 func (c *Catalog) resolveReconcileToken(target DirectoryTarget) (DirectoryTarget, bool) {
101 key := queuePathKey(target.Path)
102 c.reconcileDirtyMu.Lock()
103 defer c.reconcileDirtyMu.Unlock()
104 queued, owned := c.reconcileQueued.Load(key)
105 if !owned {
106 return DirectoryTarget{}, false
107 }
108 target = newestReconcileTarget(target, queued.(DirectoryTarget))
109 if latest, dirty := c.reconcileDirty[key]; dirty {
110 target = newestReconcileTarget(target, latest)
111 delete(c.reconcileDirty, key)
112 }
113 c.reconcileQueued.Store(key, target)
114 return target, true
115 }
116
117 func newestReconcileTarget(current, candidate DirectoryTarget) DirectoryTarget {
118 if candidate.mutationSeq > current.mutationSeq {
119 return candidate
120 }
121 return current
122 }
123
124 func (c *Catalog) takeReconcileDirty() (DirectoryTarget, bool) {
125 c.reconcileDirtyMu.Lock()
126 defer c.reconcileDirtyMu.Unlock()
127 for key, target := range c.reconcileDirty {
128 delete(c.reconcileDirty, key)
129 c.reconcileQueued.Store(key, target)
130 return target, true
131 }
132 return DirectoryTarget{}, false
133 }
134
135 func (c *Catalog) reconcileLoop() {
136 defer c.workers.Done()
137 select {
138 case <-c.discoveryStart:
139 case <-c.stop:
140 return
141 }
142 if c.opts.MetadataOnly {
143 c.metadataReconcileLoop()
144 return
145 }
146 ticker := time.NewTicker(250 * time.Millisecond)
147 defer ticker.Stop()
148 for {
149 select {
150 case token := <-c.reconcileCh:
151 if target, ok := c.resolveReconcileToken(token); ok {
152 c.runQueuedReconcile(target)
153 }
154 continue
155 default:
156 }
157 if target, ok := c.takeReconcileDirty(); ok {
158 c.runQueuedReconcile(target)
159 continue
160 }
161 select {
162 case token := <-c.reconcileCh:
163 if target, ok := c.resolveReconcileToken(token); ok {
164 c.runQueuedReconcile(target)
165 }
166 case <-ticker.C:
167 case <-c.stop:
168 return
169 }
170 }
171 }
172
173 // ResumeDiscovery separates a queryable projection from background discovery.
174 // Initial watched roots must join restored journal entries before dispatch;
175 // admitting them after resume would turn the same startup discovery into a
176 // second scan whenever an interrupted root was already pending on disk.
177 // Rejected roots remain the caller's responsibility until admission succeeds.
178 func (c *Catalog) ResumeDiscovery(initial ...DirectoryTarget) (rejected []DirectoryTarget) {
179 for _, target := range initial {
180 if !c.RequestReconcile(target) {
181 rejected = append(rejected, target)
182 }
183 }
184 c.discoveryOnce.Do(func() { close(c.discoveryStart) })
185 return rejected
186 }
187
188 func (c *Catalog) runQueuedReconcile(target DirectoryTarget) {
189 key := queuePathKey(target.Path)
190 for {
191 if c.testReconcileStartHook != nil {
192 c.testReconcileStartHook(target)
193 }
194 // A metadata scan is resumable between slices. A directory-size timeout
195 // would restart large directories forever before reaching EOF.
196 ctx, cancel := context.WithCancel(c.workerCtx)
197 if !c.opts.MetadataOnly {
198 cancel()
199 ctx, cancel = context.WithTimeout(c.workerCtx, 2*time.Minute)
200 }
201 _ = c.reconcileDirectory(ctx, target, target.mutationSeq)
202 cancel()
203
204 c.reconcileDirtyMu.Lock()
205 followUp, dirty := c.reconcileDirty[key]
206 if dirty {
207 delete(c.reconcileDirty, key)
208 c.reconcileQueued.Store(key, followUp)
209 c.reconcileDirtyMu.Unlock()
210 target = followUp
211 continue
212 }
213 c.reconcileQueued.Delete(key)
214 if done := c.reconcileDone[key]; done != nil {
215 delete(c.reconcileDone, key)
216 close(done)
217 }
218 c.reconcileDirtyMu.Unlock()
219 return
220 }
221 }
222
222 lines GO