From 3aa7885bb048cce22b34106fa9cedbc15724fa95 Mon Sep 17 00:00:00 2001 From: n3oney Date: Mon, 10 Aug 2026 16:35:50 +0200 Subject: [PATCH] workflow, spindle: enforce configurable schedule intervals Signed-off-by: n3oney --- docs/DOCS.md | 10 ++- spindle/config/config.go | 6 +- spindle/config/config_test.go | 3 + spindle/db/schedule_runs.go | 18 ++++- spindle/schedule.go | 65 ++++++++++++++---- spindle/schedule_test.go | 123 +++++++++++++++++++++++----------- workflow/def.go | 86 ++++++++++++++++++++++++ workflow/def_test.go | 41 ++++++++++++ 8 files changed, 295 insertions(+), 57 deletions(-) diff --git a/docs/DOCS.md b/docs/DOCS.md index 5b664d8de..f86d80c5b 100644 --- a/docs/DOCS.md +++ b/docs/DOCS.md @@ -1259,7 +1259,15 @@ or `H/15`. Unbounded day-of-month hashes use days 1–28. Exact expressions remain exact; spindle logs a warning recommending `H` rather than silently changing their timing. -Spindle checks schedules once per minute. +Spindle checks schedules once per minute. The minimum interval +defaults to five minutes and is configured with +`SPINDLE_SCHEDULE_MINIMUM_INTERVAL` (for example `10m`). +It must be a positive whole number of minutes. Spindle rejects +workflows whose combined schedule entries have distinct +occurrences closer than this interval. Validation is best-effort: +it examines up to 400 days or 4096 merged occurrences, whichever +comes first. Simultaneous entries count as one run. Runtime claims +also enforce the minimum interval. `SPINDLE_SCHEDULE_CONCURRENCY` controls parallel repository dispatches and startup schedule refreshes (default `4`). diff --git a/spindle/config/config.go b/spindle/config/config.go index f19fd1bd7..b933c642d 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -208,7 +208,8 @@ type Logging struct { } type Schedule struct { - Concurrency int `env:"CONCURRENCY, default=4"` + MinimumInterval time.Duration `env:"MINIMUM_INTERVAL, default=5m"` + Concurrency int `env:"CONCURRENCY, default=4"` } type Config struct { @@ -249,6 +250,9 @@ func (c *Config) validate() error { if c.Mill.DrainTimeout <= 0 { return fmt.Errorf("SPINDLE_MILL_DRAIN_TIMEOUT must be greater than zero") } + if c.Schedule.MinimumInterval < time.Minute || c.Schedule.MinimumInterval%time.Minute != 0 { + return fmt.Errorf("SPINDLE_SCHEDULE_MINIMUM_INTERVAL must be a positive whole number of minutes") + } if c.Schedule.Concurrency < 1 { return fmt.Errorf("SPINDLE_SCHEDULE_CONCURRENCY must be greater than zero") } diff --git a/spindle/config/config_test.go b/spindle/config/config_test.go index 53df45710..b5ce452bf 100644 --- a/spindle/config/config_test.go +++ b/spindle/config/config_test.go @@ -72,6 +72,9 @@ func TestLoadRejectsInvalidScheduleLimits(t *testing.T) { t.Setenv("SPINDLE_SERVER_HOSTNAME", "spindle.example.com") t.Setenv("SPINDLE_SERVER_OWNER", "did:web:spindle.example.com") for _, tc := range []struct{ name, key, value string }{ + {"disabled interval", "SPINDLE_SCHEDULE_MINIMUM_INTERVAL", "0s"}, + {"sub-minute interval", "SPINDLE_SCHEDULE_MINIMUM_INTERVAL", "30s"}, + {"fractional minute", "SPINDLE_SCHEDULE_MINIMUM_INTERVAL", "90s"}, {"no workers", "SPINDLE_SCHEDULE_CONCURRENCY", "0"}, {"negative workers", "SPINDLE_SCHEDULE_CONCURRENCY", "-1"}, } { diff --git a/spindle/db/schedule_runs.go b/spindle/db/schedule_runs.go index 442f0f310..018d4557f 100644 --- a/spindle/db/schedule_runs.go +++ b/spindle/db/schedule_runs.go @@ -5,11 +5,23 @@ import ( "time" ) -func (d *DB) ClaimScheduleRun(ctx context.Context, repoDid, workflow string, scheduledAt time.Time) (bool, error) { +func (d *DB) ClaimScheduleRun(ctx context.Context, repoDid, workflow string, scheduledAt time.Time, minimumInterval time.Duration) (bool, error) { + scheduledAtUnix := scheduledAt.Unix() result, err := d.ExecContext(ctx, ` insert or ignore into schedule_runs (repo_did, workflow, scheduled_at) - values (?, ?, ?) - `, repoDid, workflow, scheduledAt.Unix()) + select ?, ?, ? + where not exists ( + select 1 + from schedule_runs + where repo_did = ? + and workflow = ? + and scheduled_at > ? + and scheduled_at <= ? + ) + `, + repoDid, workflow, scheduledAtUnix, + repoDid, workflow, scheduledAt.Add(-minimumInterval).Unix(), scheduledAtUnix, + ) if err != nil { return false, err } diff --git a/spindle/schedule.go b/spindle/schedule.go index 18f3bdab4..f0718a70c 100644 --- a/spindle/schedule.go +++ b/spindle/schedule.go @@ -69,7 +69,8 @@ type pipelineScheduler struct { dispatch scheduleDispatcher now func() time.Time - concurrency int + minimumInterval time.Duration + concurrency int mu sync.Mutex byRepo map[string][]*scheduledWorkflow @@ -81,18 +82,20 @@ type pipelineScheduler struct { } func newPipelineScheduler(s *Spindle) *pipelineScheduler { + minimumInterval := s.cfg.Schedule.MinimumInterval concurrency := s.cfg.Schedule.Concurrency if concurrency < 1 { concurrency = 1 } return &pipelineScheduler{ - db: s.db, - l: log.SubLogger(s.l, "schedule"), - load: s.loadRepoSchedules, - dispatch: s.dispatchScheduledPipeline, - now: time.Now, - concurrency: concurrency, - byRepo: make(map[string][]*scheduledWorkflow), + db: s.db, + l: log.SubLogger(s.l, "schedule"), + load: s.loadRepoSchedules, + dispatch: s.dispatchScheduledPipeline, + now: time.Now, + minimumInterval: minimumInterval, + concurrency: concurrency, + byRepo: make(map[string][]*scheduledWorkflow), } } @@ -108,7 +111,7 @@ func (s *pipelineScheduler) run(ctx context.Context) { dispatches := make(chan time.Time, scheduleDispatchBacklog) go s.runDispatches(ctx, dispatches) - if err := s.db.PruneScheduleRuns(ctx, s.now().UTC().Add(-24*time.Hour)); err != nil { + if err := s.db.PruneScheduleRuns(ctx, s.schedulePruneBefore(s.now())); err != nil { s.l.Error("failed to prune schedule claims", "err", err) } pruneTicker := time.NewTicker(scheduleRunPruneInterval) @@ -130,13 +133,21 @@ func (s *pipelineScheduler) run(ctx context.Context) { } timer.Reset(time.Until(nextScheduleMinute(s.now()))) case <-pruneTicker.C: - if err := s.db.PruneScheduleRuns(ctx, s.now().UTC().Add(-24*time.Hour)); err != nil { + if err := s.db.PruneScheduleRuns(ctx, s.schedulePruneBefore(s.now())); err != nil { s.l.Error("failed to prune schedule claims", "err", err) } } } } +func (s *pipelineScheduler) schedulePruneBefore(now time.Time) time.Time { + retention := 24 * time.Hour + if s.minimumInterval > retention { + retention = s.minimumInterval + } + return now.UTC().Add(-retention) +} + func nextScheduleMinute(now time.Time) time.Time { return now.UTC().Truncate(time.Minute).Add(time.Minute) } @@ -170,6 +181,7 @@ func (s *pipelineScheduler) restoreSchedules(ctx context.Context) error { return fmt.Errorf("listing workflow schedules: %w", err) } grouped := make(map[string][]scheduledWorkflow) + groupedSchedules := make(map[string]map[string][]cron.Schedule) seen := make(map[struct { repo string workflow string @@ -219,10 +231,31 @@ func (s *pipelineScheduler) restoreSchedules(ctx context.Context) error { timezone: timezone, schedule: schedule, }) + if groupedSchedules[definition.RepoDid] == nil { + groupedSchedules[definition.RepoDid] = make(map[string][]cron.Schedule) + } + groupedSchedules[definition.RepoDid][definition.Workflow] = append( + groupedSchedules[definition.RepoDid][definition.Workflow], + schedule, + ) } now := s.now() for repoDid, entries := range grouped { - s.replaceCachedSchedules(repoDid, entries, now) + invalidWorkflows := make(map[string]struct{}) + for workflowName, schedules := range groupedSchedules[repoDid] { + if err := workflow.ValidateScheduleInterval(schedules, s.minimumInterval, now); err != nil { + s.l.Warn("ignoring invalid persisted workflow schedules", "repo", repoDid, "workflow", workflowName, "minimumInterval", s.minimumInterval, "err", err) + invalidWorkflows[workflowName] = struct{}{} + } + } + valid := make([]scheduledWorkflow, 0, len(entries)) + for _, entry := range entries { + if _, invalid := invalidWorkflows[entry.name]; invalid { + continue + } + valid = append(valid, entry) + } + s.replaceCachedSchedules(repoDid, valid, now) } s.l.Info("restored persisted schedules", "repositories", len(grouped), "schedules", len(definitions)) return nil @@ -492,7 +525,7 @@ func (s *pipelineScheduler) dispatchRepository(ctx context.Context, group dueRep claimed := names[:0] for _, name := range names { - ok, err := s.db.ClaimScheduleRun(ctx, repoDid, name, at) + ok, err := s.db.ClaimScheduleRun(ctx, repoDid, name, at, s.minimumInterval) if err != nil { s.l.Error("failed to claim scheduled workflow", "repo", repoDid, "workflow", name, "scheduledAt", at, "err", err) continue @@ -558,9 +591,12 @@ func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]schedu s.l.Warn("scheduled workflow manifest is invalid", "repo", repo.RepoDid, "diagnostic", diagnostic.String()) } + minimumInterval := s.cfg.Schedule.MinimumInterval + after := time.Now() var entries []scheduledWorkflow seen := make(map[struct{ name, expression, timezone string }]struct{}) for _, wf := range parsed { + var workflowSchedules []cron.Schedule var workflowEntries []scheduledWorkflow invalid := false for _, constraint := range wf.When { @@ -588,6 +624,7 @@ func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]schedu if !workflow.HasHashedCron(expression) { s.l.Warn("scheduled workflow uses exact cron fields; consider H for load spreading", "repo", repo.RepoDid, "workflow", wf.Name, "expression", expression) } + workflowSchedules = append(workflowSchedules, schedule) workflowEntries = append(workflowEntries, scheduledWorkflow{ repo: repo, scheduledRepo: scheduledRepo, @@ -602,6 +639,10 @@ func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]schedu s.l.Warn("ignoring workflow with invalid schedule", "repo", repo.RepoDid, "workflow", wf.Name) continue } + if err := workflow.ValidateScheduleInterval(workflowSchedules, minimumInterval, after); err != nil { + s.l.Warn("ignoring workflow with schedules that are too close", "repo", repo.RepoDid, "workflow", wf.Name, "minimumInterval", minimumInterval, "err", err) + continue + } entries = append(entries, workflowEntries...) } return entries, scheduledRepo, nil diff --git a/spindle/schedule_test.go b/spindle/schedule_test.go index e2ca73be2..97cc984a7 100644 --- a/spindle/schedule_test.go +++ b/spindle/schedule_test.go @@ -73,8 +73,9 @@ func TestPipelineSchedulerDueUsesUTCMinuteAndDeduplicatesWorkflows(t *testing.T) noon.name = "noon.yml" scheduler := &pipelineScheduler{ - byRepo: make(map[string][]*scheduledWorkflow), - concurrency: 1, + byRepo: make(map[string][]*scheduledWorkflow), + minimumInterval: 5 * time.Minute, + concurrency: 1, } scheduler.replaceCachedSchedules( scheduleTestRepo().RepoDid.String(), @@ -113,11 +114,12 @@ func TestPipelineSchedulerDispatchDueClaimsOccurrenceOnce(t *testing.T) { } for range 2 { scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - byRepo: make(map[string][]*scheduledWorkflow), - concurrency: 1, - dispatch: dispatch, + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + byRepo: make(map[string][]*scheduledWorkflow), + minimumInterval: 5 * time.Minute, + concurrency: 1, + dispatch: dispatch, } scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) @@ -148,10 +150,11 @@ func TestPipelineSchedulerBoundsConcurrentRepositoryDispatches(t *testing.T) { var mu sync.Mutex active, maxActive := 0, 0 scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - concurrency: 2, - byRepo: make(map[string][]*scheduledWorkflow), + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 5 * time.Minute, + concurrency: 2, + byRepo: make(map[string][]*scheduledWorkflow), dispatch: func(_ context.Context, repo db.Repo, _ db.ScheduledRepo, _ []string, _ time.Time) (models.PipelineId, error) { mu.Lock() active++ @@ -205,10 +208,11 @@ func TestPipelineSchedulerReleasesClaimWhenDispatchIsCanceled(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - concurrency: 1, - byRepo: make(map[string][]*scheduledWorkflow), + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 5 * time.Minute, + concurrency: 1, + byRepo: make(map[string][]*scheduledWorkflow), dispatch: func(ctx context.Context, _ db.Repo, _ db.ScheduledRepo, _ []string, _ time.Time) (models.PipelineId, error) { close(started) <-ctx.Done() @@ -228,7 +232,7 @@ func TestPipelineSchedulerReleasesClaimWhenDispatchIsCanceled(t *testing.T) { case <-time.After(time.Second): t.Fatal("canceled dispatch did not finish") } - reclaimed, err := d.ClaimScheduleRun(context.Background(), entry.repo.RepoDid.String(), entry.name, at) + reclaimed, err := d.ClaimScheduleRun(context.Background(), entry.repo.RepoDid.String(), entry.name, at, scheduler.minimumInterval) if err != nil { t.Fatal(err) } @@ -237,6 +241,39 @@ func TestPipelineSchedulerReleasesClaimWhenDispatchIsCanceled(t *testing.T) { } } +func TestPipelineSchedulerEnforcesMinimumRunIntervalAcrossSchedules(t *testing.T) { + at := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) + entries := make([]scheduledWorkflow, 0, 3) + for _, expression := range []string{"0 9 * * *", "5 9 * * *", "10 9 * * *"} { + entry := mustCron(t, expression) + entry.name = "build.yml" + entries = append(entries, entry) + } + + d := newScheduleTestDB(t) + var dispatched []time.Time + scheduler := &pipelineScheduler{ + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 10 * time.Minute, + concurrency: 1, + byRepo: make(map[string][]*scheduledWorkflow), + dispatch: func(_ context.Context, repo db.Repo, _ db.ScheduledRepo, _ []string, scheduledAt time.Time) (models.PipelineId, error) { + dispatched = append(dispatched, scheduledAt) + return models.PipelineId{Knot: repo.Knot, Rkey: "pipeline"}, nil + }, + } + scheduler.replaceCachedSchedules(scheduleTestRepo().RepoDid.String(), entries, at.Add(-time.Minute)) + scheduler.dispatchDue(context.Background(), at) + scheduler.dispatchDue(context.Background(), at.Add(5*time.Minute)) + scheduler.dispatchDue(context.Background(), at.Add(10*time.Minute)) + + want := []time.Time{at, at.Add(10 * time.Minute)} + if !slices.Equal(dispatched, want) { + t.Fatalf("scheduled dispatches = %v, want %v", dispatched, want) + } +} + func TestPipelineSchedulerDispatchFailureClaimLifecycle(t *testing.T) { tests := []struct { name string @@ -258,10 +295,11 @@ func TestPipelineSchedulerDispatchFailureClaimLifecycle(t *testing.T) { at := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - concurrency: 1, - byRepo: make(map[string][]*scheduledWorkflow), + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 5 * time.Minute, + concurrency: 1, + byRepo: make(map[string][]*scheduledWorkflow), dispatch: func(context.Context, db.Repo, db.ScheduledRepo, []string, time.Time) (models.PipelineId, error) { return test.pipelineID, errors.New("temporary failure") }, @@ -269,7 +307,7 @@ func TestPipelineSchedulerDispatchFailureClaimLifecycle(t *testing.T) { scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) scheduler.dispatchDue(context.Background(), at) - reclaimed, err := d.ClaimScheduleRun(context.Background(), entry.repo.RepoDid.String(), entry.name, at) + reclaimed, err := d.ClaimScheduleRun(context.Background(), entry.repo.RepoDid.String(), entry.name, at, scheduler.minimumInterval) if err != nil { t.Fatal(err) } @@ -316,11 +354,12 @@ func TestPipelineSchedulerRestoresPersistedSchedulesWithoutLoadingRepository(t * at := time.Date(2026, time.August, 10, 9, 30, 0, 0, time.UTC) scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - concurrency: 1, - now: func() time.Time { return at.Add(-time.Minute) }, - byRepo: make(map[string][]*scheduledWorkflow), + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 5 * time.Minute, + concurrency: 1, + now: func() time.Time { return at.Add(-time.Minute) }, + byRepo: make(map[string][]*scheduledWorkflow), } if err := scheduler.restoreSchedules(ctx); err != nil { t.Fatalf("restoreSchedules: %v", err) @@ -360,8 +399,9 @@ func TestPipelineSchedulerDoesNotEvaluateNonDueSchedulesEachMinute(t *testing.T) } scheduler := &pipelineScheduler{ - byRepo: make(map[string][]*scheduledWorkflow), - concurrency: 1, + byRepo: make(map[string][]*scheduledWorkflow), + minimumInterval: 5 * time.Minute, + concurrency: 1, } scheduler.replaceCachedSchedules(scheduleTestRepo().RepoDid.String(), entries, now) for _, schedule := range schedules { @@ -410,8 +450,9 @@ func TestPipelineSchedulerDropsSchedulesThatDoNotAdvance(t *testing.T) { schedule: &sequenceSchedule{next: test.next}, } scheduler := &pipelineScheduler{ - byRepo: make(map[string][]*scheduledWorkflow), - concurrency: 1, + byRepo: make(map[string][]*scheduledWorkflow), + minimumInterval: 5 * time.Minute, + concurrency: 1, } scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) due := scheduler.due(at) @@ -446,11 +487,12 @@ func TestPipelineSchedulerSerializesRepositoryRefreshes(t *testing.T) { var callMu sync.Mutex calls := 0 scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - concurrency: 1, - now: time.Now, - byRepo: make(map[string][]*scheduledWorkflow), + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 5 * time.Minute, + concurrency: 1, + now: time.Now, + byRepo: make(map[string][]*scheduledWorkflow), load: func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { callMu.Lock() calls++ @@ -513,11 +555,12 @@ func TestPipelineSchedulerRemovalWaitsForInflightRefresh(t *testing.T) { release := func() { releaseOnce.Do(func() { close(releaseLoad) }) } defer release() scheduler := &pipelineScheduler{ - db: d, - l: slog.New(slog.NewTextHandler(io.Discard, nil)), - concurrency: 1, - now: time.Now, - byRepo: make(map[string][]*scheduledWorkflow), + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + minimumInterval: 5 * time.Minute, + concurrency: 1, + now: time.Now, + byRepo: make(map[string][]*scheduledWorkflow), load: func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { close(loadStarted) <-releaseLoad diff --git a/workflow/def.go b/workflow/def.go index 0bdf735f5..10d56591d 100644 --- a/workflow/def.go +++ b/workflow/def.go @@ -302,6 +302,92 @@ func validateCronHasOccurrence(schedule cron.Schedule) error { return nil } +const ( + // ValidateScheduleInterval examines a bounded future window of roughly 400 + // days and at most 4096 merged occurrences total. This best-effort bound + // covers dense and ordinary sparse calendars while preventing pathological + // schedules from hanging validation. + scheduleValidationHorizon = 400 * 24 * time.Hour + scheduleValidationMaxOccurrences = 4096 +) + +type scheduleCursor struct { + index int + schedule cron.Schedule + next time.Time +} + +// ValidateScheduleInterval merges actual future occurrences from all schedules, +// collapsing identical instants, and rejects distinct occurrences that are +// closer than minimum. It scans roughly 400 days and at most 4096 merged +// occurrences, returning nil when either bound is reached. +func ValidateScheduleInterval(schedules []cron.Schedule, minimum time.Duration, after time.Time) error { + if minimum <= 0 { + return fmt.Errorf("minimum schedule interval must be positive") + } + + horizon := after.Add(scheduleValidationHorizon) + cursors := make([]scheduleCursor, 0, len(schedules)) + for index, schedule := range schedules { + if schedule == nil { + return fmt.Errorf("schedule[%d] is nil", index) + } + next := schedule.Next(after) + if next.IsZero() || next.After(horizon) { + continue + } + if !next.After(after) { + return fmt.Errorf("schedule[%d] did not advance beyond %s", index, after) + } + cursors = append(cursors, scheduleCursor{index: index, schedule: schedule, next: next}) + } + + var previous time.Time + havePrevious := false + for range scheduleValidationMaxOccurrences { + if len(cursors) == 0 { + break + } + + earliestIndex := 0 + for index := 1; index < len(cursors); index++ { + if cursors[index].next.Before(cursors[earliestIndex].next) { + earliestIndex = index + } + } + current := cursors[earliestIndex].next + if !havePrevious || !current.Equal(previous) { + if havePrevious { + interval := current.Sub(previous) + if interval < minimum { + return fmt.Errorf("schedule occurrences %s and %s are %s apart, below minimum %s", previous, current, interval, minimum) + } + } + previous = current + havePrevious = true + } + + for index := 0; index < len(cursors); { + if !cursors[index].next.Equal(current) { + index++ + continue + } + prior := cursors[index].next + next := cursors[index].schedule.Next(prior) + if next.IsZero() || next.After(horizon) { + cursors = append(cursors[:index], cursors[index+1:]...) + continue + } + if !next.After(prior) { + return fmt.Errorf("schedule[%d] did not advance beyond %s", cursors[index].index, prior) + } + cursors[index].next = next + index++ + } + } + return nil +} + // matchesPattern checks if a name matches any of the given patterns. // Patterns can be exact matches or glob patterns using * and **. // * matches any sequence of non-separator characters diff --git a/workflow/def_test.go b/workflow/def_test.go index 62646f260..108342cb1 100644 --- a/workflow/def_test.go +++ b/workflow/def_test.go @@ -4,6 +4,7 @@ import ( "testing" "time" + "github.com/robfig/cron/v3" "github.com/stretchr/testify/assert" "tangled.org/core/api/tangled" ) @@ -838,3 +839,43 @@ func TestHasHashedCron(t *testing.T) { assert.False(t, HasHashedCron("0 3 * * THU")) assert.False(t, HasHashedCron("*/5 * * * *")) } + +func TestParseCronAllowsIntervalsBelowHostedMinimum(t *testing.T) { + _, err := ParseCron("*/4 * * * *", "UTC", "interval-seed") + assert.NoError(t, err) +} + +type stuckCronSchedule struct { + at time.Time +} + +func (s stuckCronSchedule) Next(time.Time) time.Time { + return s.at +} + +func TestValidateScheduleIntervalRejectsNonAdvancingSchedule(t *testing.T) { + after := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC) + schedule := stuckCronSchedule{at: after.Add(time.Minute)} + assert.Error(t, ValidateScheduleInterval([]cron.Schedule{schedule}, 5*time.Minute, after)) +} + +func TestValidateScheduleInterval(t *testing.T) { + after := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC) + closeFirst, err := ParseCron("0 * * * *", "UTC", "close-first") + assert.NoError(t, err) + closeSecond, err := ParseCron("2 * * * *", "UTC", "close-second") + assert.NoError(t, err) + assert.Error(t, ValidateScheduleInterval([]cron.Schedule{closeFirst, closeSecond}, 5*time.Minute, after)) + + boundary, err := ParseCron("5 * * * *", "UTC", "boundary") + assert.NoError(t, err) + assert.NoError(t, ValidateScheduleInterval([]cron.Schedule{closeFirst, boundary}, 5*time.Minute, after)) + + duplicate, err := ParseCron("0 9 * * *", "UTC", "duplicate") + assert.NoError(t, err) + assert.NoError(t, ValidateScheduleInterval([]cron.Schedule{duplicate, duplicate}, 5*time.Minute, after)) + + sparse, err := ParseCron("0 0 29 2 *", "UTC", "sparse") + assert.NoError(t, err) + assert.NoError(t, ValidateScheduleInterval([]cron.Schedule{sparse}, 5*time.Minute, after)) +} -- 2.51.2