返回 DeepSeek-Reasonix
metadata_queue.go
根目录 / internal / sessioncatalog / metadata_queue.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "errors"
6 "os"
7 "time"
8
9 "reasonix/internal/historywork"
10 )
11
12 type metadataQueueJob struct {
13 target DirectoryTarget
14 scan *metadataScan
15 ready time.Time
16 turn uint64
17 failures int
18 restartDeadline time.Time
19 }
20
21 // A known-invalid iterator must not consume an entire large root before the
22 // replacement can start. Coalesce a burst between slices, with bounded delay
23 // so a continuously changing root still gets discovery work admitted.
24 func deferMetadataRestart(job *metadataQueueJob, now time.Time) {
25 if job.restartDeadline.IsZero() {
26 job.restartDeadline = now.Add(time.Second)
27 }
28 job.ready = minTime(now.Add(historywork.PauseDuration), job.restartDeadline)
29 }
30
31 func minTime(a, b time.Time) time.Time {
32 if a.Before(b) {
33 return a
34 }
35 return b
36 }
37
38 // Keep one slot available for a newly visible workspace. Existing iterators
39 // are never evicted to admit another root: restarting them would repeatedly
40 // read the same prefix and could prevent large directories from completing.
41 const metadataIteratorLimit = 8
42
43 func selectMetadataQueueJob(jobs map[string]*metadataQueueJob, now time.Time, foreground bool, priority func(DirectoryTarget) bool) (string, *metadataQueueJob) {
44 active := 0
45 for _, job := range jobs {
46 if job.scan != nil {
47 active++
48 }
49 }
50 var selected string
51 var next *metadataQueueJob
52 for key, job := range jobs {
53 if job.ready.After(now) {
54 continue
55 }
56 visible := priority(job.target)
57 if !visible && foreground {
58 continue
59 }
60 if job.scan == nil && (active >= metadataIteratorLimit || !visible && active >= metadataIteratorLimit-1) {
61 continue
62 }
63 if next == nil || visible && !priority(next.target) || visible == priority(next.target) && (job.turn < next.turn || job.turn == next.turn && job.target.mutationSeq < next.target.mutationSeq) {
64 selected, next = key, job
65 }
66 }
67 return selected, next
68 }
69
70 // The queue retains a bounded set of directory iterators and admits only one
71 // slice. Admitted roots rotate without restarting their prefixes; waiting
72 // roots enter as earlier scans finish. Committed path updates use the
73 // independent small-metadata writer and never need an iterator slot.
74 func (c *Catalog) metadataReconcileLoop() {
75 jobs := map[string]*metadataQueueJob{}
76 defer func() {
77 for _, job := range jobs {
78 if job.scan != nil {
79 job.scan.close(c.workerCtx, context.Canceled)
80 }
81 }
82 }()
83 var turn uint64
84 for c.workerCtx.Err() == nil {
85 c.reconcileDirtyMu.Lock()
86 c.reconcileQueued.Range(func(key, value any) bool {
87 id := key.(string)
88 if jobs[id] == nil {
89 jobs[id] = &metadataQueueJob{target: value.(DirectoryTarget)}
90 delete(c.reconcileDirty, id)
91 }
92 return true
93 })
94 c.reconcileDirtyMu.Unlock()
95 now := c.opts.Now()
96 foreground := c.opts.Maintenance != nil && c.opts.Maintenance.ForegroundActive()
97 selected, next := selectMetadataQueueJob(jobs, now, foreground, c.isPriorityDirectory)
98 if next == nil {
99 timer := time.NewTimer(historywork.PauseDuration)
100 select {
101 case <-c.stop:
102 timer.Stop()
103 return
104 case <-c.reconcileCh:
105 timer.Stop()
106 case <-timer.C:
107 }
108 continue
109 }
110 turn++
111 next.turn = turn
112 c.reconcileDirtyMu.Lock()
113 latest, invalidated := c.reconcileDirty[selected]
114 if invalidated && (next.scan != nil || !next.restartDeadline.IsZero()) {
115 delete(c.reconcileDirty, selected)
116 } else {
117 invalidated = false
118 }
119 c.reconcileDirtyMu.Unlock()
120 if invalidated {
121 if next.scan != nil {
122 // Abandoning an incomplete observation cannot confirm missing
123 // rows. Keep its journal and visible prefix until the new EOF.
124 c.observeDiscovery(next.target, "superseded", "", "")
125 next.scan.close(c.workerCtx, nil)
126 next.scan = nil
127 }
128 next.target = newestReconcileTarget(next.target, latest)
129 deferMetadataRestart(next, now)
130 if next.ready.After(now) {
131 continue
132 }
133 }
134 var err error
135 if next.scan == nil {
136 // Waiting for an iterator or priority slot is not a running scan.
137 // Fold all pre-dispatch invalidations into this scan; only changes
138 // arriving after dispatch need a follow-up pass.
139 target, owned := c.resolveReconcileToken(next.target)
140 if !owned {
141 delete(jobs, selected)
142 continue
143 }
144 next.target = target
145 next.restartDeadline = time.Time{}
146 c.observeDiscovery(next.target, "started", "", "")
147 if c.testReconcileStartHook != nil {
148 c.testReconcileStartHook(next.target)
149 }
150 next.scan, err = c.startMetadataScan(c.workerCtx, next.target, next.target.mutationSeq, true)
151 }
152 var done bool
153 var bytes int64
154 if err == nil {
155 done, bytes, err = next.scan.step(c.workerCtx)
156 }
157 if errors.Is(err, errMetadataScanBusy) || errors.Is(err, historywork.ErrForegroundActive) {
158 next.ready = now.Add(historywork.PauseDuration)
159 continue
160 }
161 if done || err != nil {
162 c.finishMetadataQueueJob(jobs, selected, next, done, err)
163 }
164 // Apply a global pause, including when the next slice belongs to a
165 // different root. Per-root timers alone multiply the allowed I/O rate.
166 if c.opts.Maintenance == nil {
167 if err := historywork.Pause(c.workerCtx, bytes); err != nil {
168 return
169 }
170 }
171 }
172 }
173
174 func (c *Catalog) PrioritizeWorkspace(scope, root string) {
175 scope, root = normalizeScope(scope, root)
176 c.priorityWorkspace.Store(scope + "\x00" + c.workspaceRootKey(scope, root))
177 }
178
179 func (c *Catalog) isPriorityDirectory(target DirectoryTarget) bool {
180 scope, root := normalizeScope(target.Scope, target.WorkspaceRoot)
181 key, _ := c.priorityWorkspace.Load().(string)
182 return key == scope+"\x00"+c.workspaceRootKey(scope, root)
183 }
184
185 func (c *Catalog) finishMetadataQueueJob(jobs map[string]*metadataQueueJob, selected string, next *metadataQueueJob, done bool, err error) {
186 phase, failure := "completed", ""
187 if err != nil {
188 phase, failure = "failed", "io_or_database"
189 if errors.Is(err, context.Canceled) {
190 failure = "canceled"
191 } else if os.IsPermission(err) {
192 failure = "permission"
193 } else if os.IsNotExist(err) {
194 failure = "missing"
195 }
196 }
197 c.observeDiscovery(next.target, phase, "", failure)
198 c.observeDatabaseError(err)
199 if next.scan != nil {
200 next.scan.close(c.workerCtx, err)
201 next.scan = nil
202 }
203 if done && err == nil {
204 c.settleReconcileTarget(next.target)
205 }
206 c.reconcileDirtyMu.Lock()
207 if follow, dirty := c.reconcileDirty[selected]; dirty {
208 delete(c.reconcileDirty, selected)
209 next.target, next.failures, next.ready = follow, 0, time.Time{}
210 } else if err != nil && !errors.Is(err, context.Canceled) && !os.IsNotExist(err) && !os.IsPermission(err) && next.failures < 3 {
211 next.ready = c.opts.Now().Add([]time.Duration{time.Second, 5 * time.Second, 30 * time.Second}[next.failures])
212 next.failures++
213 } else {
214 c.reconcileQueued.Delete(selected)
215 delete(jobs, selected)
216 if done := c.reconcileDone[selected]; done != nil {
217 delete(c.reconcileDone, selected)
218 close(done)
219 }
220 }
221 c.reconcileDirtyMu.Unlock()
222 }
223
223 lines GO