(interval time.Duration)
| 220 | } |
| 221 | |
| 222 | func (s *Scanner) launchFixLoop(interval time.Duration) { |
| 223 | s.wg.Add(1) |
| 224 | go func() { |
| 225 | defer s.wg.Done() |
| 226 | |
| 227 | s.lg.Info("Launching scan fix loop", "interval", interval.String()) |
| 228 | |
| 229 | fixTicker := time.NewTicker(interval) |
| 230 | defer fixTicker.Stop() |
| 231 | |
| 232 | for { |
| 233 | select { |
| 234 | case <-fixTicker.C: |
| 235 | // hold until other possible ongoing scans finish |
| 236 | if !s.acquireScanPermit() { |
| 237 | return |
| 238 | } |
| 239 | s.statsMu.Lock() |
| 240 | kvIndices := s.stats.needFix() |
| 241 | s.statsMu.Unlock() |
| 242 | s.lg.Info("Scanner fixing batch triggered", "mismatchesToFix", kvIndices) |
| 243 | err := s.worker.fixBatch(s.ctx, kvIndices, func(kvi uint64, m *scanned) { |
| 244 | s.updateStats(kvi, m) |
| 245 | }) |
| 246 | s.releaseScanPermit() |
| 247 | if err != nil { |
| 248 | s.lg.Error("Fixing batch failed", "error", err) |
| 249 | } |
| 250 | |
| 251 | case <-s.ctx.Done(): |
| 252 | return |
| 253 | } |
| 254 | } |
| 255 | }() |
| 256 | } |
| 257 | |
| 258 | func (s *Scanner) logStats() { |
| 259 | localKvCount, sum := s.worker.summaryLocalKvs() |
no test coverage detected