| 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 |