返回 DeepSeek-Reasonix
reconcile_queue.go
根目录 / internal / sessioncatalog / reconcile_queue.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "strings"
6 "time"
7 )
8
9 // RequestReconcile makes the channel a wake signal while the maps retain the
10 // newest target. Session saves never wait for catalog work.
11 func (c *Catalog) RequestReconcile(target DirectoryTarget) bool {
12 if c == nil || strings.TrimSpace(target.Path) == "" {
13 return false
14 }
15 target.Path = cleanCatalogAccessPath(target.Path)
16 key := queuePathKey(target.Path)
17 if key == "" {
18 return false
19 }
20 target.mutationSeq = c.mutationSeq.Add(1)
21 if _, loaded := c.reconcileQueued.LoadOrStore(key, target); loaded {
22 c.markReconcileDirty(target)
23 return true
24 }
25 select {
26 case c.reconcileCh <- target:
27 return true
28 case <-c.stop:
29 c.reconcileQueued.Delete(key)
30 return false
31 default:
32 c.reconcileQueued.Delete(key)
33 c.markReconcileDirty(target)
34 return false
35 }
36 }
37
38 func (c *Catalog) markReconcileDirty(target DirectoryTarget) {
39 key := queuePathKey(target.Path)
40 c.reconcileDirtyMu.Lock()
41 if queued, ok := c.reconcileQueued.Load(key); ok {
42 target = newestReconcileTarget(queued.(DirectoryTarget), target)
43 }
44 if dirty, ok := c.reconcileDirty[key]; ok {
45 target = newestReconcileTarget(dirty, target)
46 }
47 c.reconcileDirty[key] = target
48 c.reconcileQueued.Store(key, target)
49 c.reconcileDirtyMu.Unlock()
50 }
51
52 func (c *Catalog) resolveReconcileToken(target DirectoryTarget) (DirectoryTarget, bool) {
53 key := queuePathKey(target.Path)
54 c.reconcileDirtyMu.Lock()
55 defer c.reconcileDirtyMu.Unlock()
56 queued, owned := c.reconcileQueued.Load(key)
57 if !owned {
58 return DirectoryTarget{}, false
59 }
60 target = newestReconcileTarget(target, queued.(DirectoryTarget))
61 if latest, dirty := c.reconcileDirty[key]; dirty {
62 target = newestReconcileTarget(target, latest)
63 delete(c.reconcileDirty, key)
64 }
65 c.reconcileQueued.Store(key, target)
66 return target, true
67 }
68
69 func newestReconcileTarget(current, candidate DirectoryTarget) DirectoryTarget {
70 if candidate.mutationSeq > current.mutationSeq {
71 return candidate
72 }
73 return current
74 }
75
76 func (c *Catalog) takeReconcileDirty() (DirectoryTarget, bool) {
77 c.reconcileDirtyMu.Lock()
78 defer c.reconcileDirtyMu.Unlock()
79 for key, target := range c.reconcileDirty {
80 delete(c.reconcileDirty, key)
81 c.reconcileQueued.Store(key, target)
82 return target, true
83 }
84 return DirectoryTarget{}, false
85 }
86
87 func (c *Catalog) reconcileLoop() {
88 defer c.workers.Done()
89 ticker := time.NewTicker(250 * time.Millisecond)
90 defer ticker.Stop()
91 for {
92 select {
93 case token := <-c.reconcileCh:
94 if target, ok := c.resolveReconcileToken(token); ok {
95 c.runQueuedReconcile(target)
96 }
97 continue
98 default:
99 }
100 if target, ok := c.takeReconcileDirty(); ok {
101 c.runQueuedReconcile(target)
102 continue
103 }
104 select {
105 case token := <-c.reconcileCh:
106 if target, ok := c.resolveReconcileToken(token); ok {
107 c.runQueuedReconcile(target)
108 }
109 case <-ticker.C:
110 case <-c.stop:
111 return
112 }
113 }
114 }
115
116 func (c *Catalog) runQueuedReconcile(target DirectoryTarget) {
117 key := queuePathKey(target.Path)
118 for {
119 if c.testReconcileStartHook != nil {
120 c.testReconcileStartHook(target)
121 }
122 ctx, cancel := context.WithTimeout(c.workerCtx, 2*time.Minute)
123 _ = c.reconcileDirectory(ctx, target, target.mutationSeq)
124 cancel()
125
126 c.reconcileDirtyMu.Lock()
127 followUp, dirty := c.reconcileDirty[key]
128 if dirty {
129 delete(c.reconcileDirty, key)
130 c.reconcileQueued.Store(key, followUp)
131 c.reconcileDirtyMu.Unlock()
132 target = followUp
133 continue
134 }
135 c.reconcileQueued.Delete(key)
136 c.reconcileDirtyMu.Unlock()
137 return
138 }
139 }
140
140 lines GO