package crontab import ( "context" "log" "sync" "sync/atomic" "time" "senlinai-agent/backend/internal/logic/dataset" ) var datasetInterval = 5 * time.Minute type DatasetCollector interface { SyncAllSources(context.Context) ([]dataset.SyncResult, error) } type DatasetScheduler struct { collector DatasetCollector interval time.Duration logger *log.Logger } func NewDatasetScheduler(collector DatasetCollector, logger *log.Logger) *DatasetScheduler { if logger == nil { logger = log.Default() } return &DatasetScheduler{ collector: collector, interval: datasetInterval, logger: logger, } } // Run collects once at startup, then on a fixed five-minute interval. // A tick is skipped when the previous collection is still running. func (s *DatasetScheduler) Run(ctx context.Context) { if ctx.Err() != nil { return } ticker := time.NewTicker(s.interval) defer ticker.Stop() var running atomic.Bool var workers sync.WaitGroup start := func() { if !running.CompareAndSwap(false, true) { s.logger.Print("dataset collection skipped: previous run is still active") return } workers.Add(1) go func() { defer workers.Done() defer running.Store(false) s.collect(ctx) }() } start() for { select { case <-ctx.Done(): workers.Wait() return case <-ticker.C: start() } } } func (s *DatasetScheduler) collect(ctx context.Context) { results, err := s.collector.SyncAllSources(ctx) completed := 0 failed := 0 for _, result := range results { switch result.Status { case "completed": completed++ case "failed": failed++ s.logger.Printf( "dataset source collection failed: source=%s error=%s", result.SourceIdentity, result.Result, ) } } if err != nil { s.logger.Printf("dataset collection failed: sources=%d completed=%d failed=%d error=%v", len(results), completed, failed, err) return } s.logger.Printf("dataset collection finished: sources=%d completed=%d failed=%d", len(results), completed, failed) }