package crontab import ( "context" "log" "sync" "sync/atomic" "time" "senlinai-agent/backend/internal/models" ) const datasetInterval = 5 * time.Minute type DatasetCollector interface { SyncAllSources(context.Context) ([]models.SaDatasetCron, 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) { crons, err := s.collector.SyncAllSources(ctx) completed := 0 failed := 0 for _, cron := range crons { switch cron.Status { case "completed": completed++ case "failed": failed++ } } if err != nil { s.logger.Printf("dataset collection failed: tasks=%d completed=%d failed=%d error=%v", len(crons), completed, failed, err) return } s.logger.Printf("dataset collection finished: tasks=%d completed=%d failed=%d", len(crons), completed, failed) }