From ba7179e8b1fa9f6eebf9b2a3d4bc337d676cb9fa Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 18 Sep 2026 13:56:49 +0300 Subject: [PATCH] spindle/schedule: validate loaded schedules against the scheduler clock Signed-off-by: dawn --- spindle/schedule.go | 7 +++---- spindle/schedule_test.go | 4 ++-- 2 files changed, 5 insertions(+), 6 deletions(-) diff --git a/spindle/schedule.go b/spindle/schedule.go index c40e966bd..20abb282a 100644 --- a/spindle/schedule.go +++ b/spindle/schedule.go @@ -62,7 +62,7 @@ func (h *scheduleHeap) Pop() any { return entry } -type scheduleLoader func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) +type scheduleLoader func(context.Context, db.Repo, time.Time) ([]scheduledWorkflow, db.ScheduledRepo, error) type scheduleDispatcher func(context.Context, db.Repo, db.ScheduledRepo, []string, time.Time) (models.PipelineId, error) type pipelineScheduler struct { @@ -354,7 +354,7 @@ func (s *pipelineScheduler) RefreshRepo(ctx context.Context, repo db.Repo) error return fmt.Errorf("checking repository registration: %w", err) } repo = *current - entries, scheduledRepo, err := s.load(ctx, repo) + entries, scheduledRepo, err := s.load(ctx, repo, s.now()) if err != nil { s.removeCachedSchedules(repo.RepoDid.String()) return err @@ -584,7 +584,7 @@ func (s *pipelineScheduler) dispatchRepository(ctx context.Context, group dueRep } } -func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { +func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo, after time.Time) ([]scheduledWorkflow, db.ScheduledRepo, error) { branch, err := s.getRepoDefaultBranch(ctx, repo) if err != nil { return nil, db.ScheduledRepo{}, err @@ -606,7 +606,6 @@ func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]schedu } minimumInterval := s.cfg.Schedule.MinimumInterval - after := time.Now() var entries []scheduledWorkflow seen := make(map[struct{ name, expression, timezone string }]struct{}) for _, wf := range parsed { diff --git a/spindle/schedule_test.go b/spindle/schedule_test.go index 8e21b47b5..164cfa405 100644 --- a/spindle/schedule_test.go +++ b/spindle/schedule_test.go @@ -528,7 +528,7 @@ func TestPipelineSchedulerSerializesRepositoryRefreshes(t *testing.T) { concurrency: 1, now: time.Now, byRepo: make(map[string][]*scheduledWorkflow), - load: func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { + load: func(_ context.Context, _ db.Repo, _ time.Time) ([]scheduledWorkflow, db.ScheduledRepo, error) { callMu.Lock() calls++ call := calls @@ -596,7 +596,7 @@ func TestPipelineSchedulerRemovalWaitsForInflightRefresh(t *testing.T) { concurrency: 1, now: time.Now, byRepo: make(map[string][]*scheduledWorkflow), - load: func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { + load: func(_ context.Context, _ db.Repo, _ time.Time) ([]scheduledWorkflow, db.ScheduledRepo, error) { close(loadStarted) <-releaseLoad return []scheduledWorkflow{entry}, scheduleTestSnapshot(), nil -- 2.51.2