diff --git a/docs/DOCS.md b/docs/DOCS.md index f2df96440..5b664d8de 100644 --- a/docs/DOCS.md +++ b/docs/DOCS.md @@ -1196,10 +1196,11 @@ has the following fields: field is required for `schedule` events and invalid for other events. Each entry requires a `cron` field containing a five-field expression (`minute hour day-of-month month - day-of-week`). Expressions are evaluated in UTC. Hashed `H` - fields are supported; seconds, nicknames such as `@daily`, - and timezone prefixes are not. A condition containing - `schedule` cannot also specify `branch`, `tag`, `types`, or `paths`. + day-of-week`) and accepts an optional IANA `timezone`, + which defaults to `UTC`. Hashed `H` fields are supported; + seconds, nicknames such as `@daily`, and timezone prefixes + are not. A condition containing `schedule` cannot also + specify `branch`, `tag`, `types`, or `paths`. For example, if you'd like to define a workflow that runs when commits are pushed to the `main` and `develop` @@ -1235,16 +1236,19 @@ when: ``` Scheduled workflows use the saved default-branch revision, -including after a spindle restart. For example, this runs at a -stable minute during the 02:00 UTC hour every day and at 09:00 UTC -on weekdays: +including after a spindle restart. A condition can contain +multiple schedule entries, each with its own timezone. +For example, this runs at a stable minute during the 02:00 +hour in New York and at 09:00 on weekdays in London: ```yaml when: - event: schedule schedule: - cron: "H 2 * * *" + timezone: America/New_York - cron: "0 9 * * 1-5" + timezone: Europe/London ``` `H` spreads scheduled work using a deterministic hash of the @@ -1255,12 +1259,12 @@ 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 evaluates schedules once per minute. Occurrences -while the spindle is offline are not replayed. +Spindle checks schedules once per minute. `SPINDLE_SCHEDULE_CONCURRENCY` controls parallel repository dispatches and startup schedule refreshes (default `4`). Dispatch finishes each minute's batch before starting the next. +Occurrences while spindle is offline are not replayed. To skip CI for a push, pass a Git push option: diff --git a/spindle/db/db.go b/spindle/db/db.go index 6cf07e169..9d4512d33 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -159,8 +159,9 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { repo_did text not null, workflow text not null, expression text not null, + timezone text not null default 'UTC', - primary key (repo_did, workflow, expression), + primary key (repo_did, workflow, expression, timezone), foreign key (repo_did) references scheduled_repos(repo_did) on delete cascade ); @@ -942,6 +943,35 @@ func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error return err } + if err := orm.RunMigration(conn, logger, "workflow-schedules-timezone", func(tx *sql.Tx) error { + var hasTimezone int + if err := tx.QueryRow( + `select count(*) from pragma_table_info('workflow_schedules') where name = 'timezone'`, + ).Scan(&hasTimezone); err != nil { + return err + } + if hasTimezone > 0 { + return nil + } + _, err := tx.Exec(` + create table workflow_schedules_new ( + repo_did text not null, + workflow text not null, + expression text not null, + timezone text not null default 'UTC', + + primary key (repo_did, workflow, expression, timezone), + foreign key (repo_did) references scheduled_repos(repo_did) on delete cascade + ); + insert into workflow_schedules_new (repo_did, workflow, expression, timezone) + select repo_did, workflow, expression, 'UTC' from workflow_schedules; + drop table workflow_schedules; + alter table workflow_schedules_new rename to workflow_schedules; + `) + return err + }); err != nil { + return err + } if err := orm.RunMigration(conn, logger, "scheduled-repos-revision", func(tx *sql.Tx) error { for _, column := range []string{"branch", "sha"} { var present int diff --git a/spindle/db/schedules.go b/spindle/db/schedules.go index 94a5b9ecc..05ce77aac 100644 --- a/spindle/db/schedules.go +++ b/spindle/db/schedules.go @@ -15,6 +15,7 @@ type WorkflowSchedule struct { RepoDid string Workflow string Expression string + Timezone string Branch string SHA string } @@ -40,10 +41,14 @@ func (d *DB) ReplaceWorkflowSchedules(ctx context.Context, repo ScheduledRepo, s return err } for _, schedule := range schedules { + timezone := schedule.Timezone + if timezone == "" { + timezone = "UTC" + } if _, err := tx.ExecContext(ctx, ` - insert or ignore into workflow_schedules (repo_did, workflow, expression) - values (?, ?, ?) - `, repo.RepoDid, schedule.Workflow, schedule.Expression); err != nil { + insert or ignore into workflow_schedules (repo_did, workflow, expression, timezone) + values (?, ?, ?, ?) + `, repo.RepoDid, schedule.Workflow, schedule.Expression, timezone); err != nil { return err } } @@ -52,13 +57,13 @@ func (d *DB) ReplaceWorkflowSchedules(ctx context.Context, repo ScheduledRepo, s func (d *DB) WorkflowSchedules(ctx context.Context) ([]WorkflowSchedule, error) { rows, err := d.QueryContext(ctx, ` - select ws.repo_did, ws.workflow, ws.expression, + select ws.repo_did, ws.workflow, ws.expression, ws.timezone, coalesce(sr.branch, ''), coalesce(sr.sha, '') from workflow_schedules ws join scheduled_repos sr on sr.repo_did = ws.repo_did where coalesce(sr.branch, '') <> '' and coalesce(sr.sha, '') <> '' - order by ws.repo_did, ws.workflow, ws.expression + order by ws.repo_did, ws.workflow, ws.expression, ws.timezone `) if err != nil { return nil, err @@ -72,6 +77,7 @@ func (d *DB) WorkflowSchedules(ctx context.Context) ([]WorkflowSchedule, error) &schedule.RepoDid, &schedule.Workflow, &schedule.Expression, + &schedule.Timezone, &schedule.Branch, &schedule.SHA, ); err != nil { diff --git a/spindle/db/schedules_test.go b/spindle/db/schedules_test.go index ba992e48b..ed555646e 100644 --- a/spindle/db/schedules_test.go +++ b/spindle/db/schedules_test.go @@ -32,8 +32,9 @@ func TestWorkflowSchedulesTrackIndexedRepositoriesIncludingEmptyManifests(t *tes } schedules := []WorkflowSchedule{ - {RepoDid: repo.RepoDid.String(), Workflow: "build.yml", Expression: "0 9 * * *"}, - {RepoDid: repo.RepoDid.String(), Workflow: "build.yml", Expression: "0 9 * * *"}, + {RepoDid: repo.RepoDid.String(), Workflow: "build.yml", Expression: "0 9 * * *", Timezone: "UTC"}, + {RepoDid: repo.RepoDid.String(), Workflow: "build.yml", Expression: "0 9 * * *", Timezone: "UTC"}, + {RepoDid: repo.RepoDid.String(), Workflow: "build.yml", Expression: "0 9 * * *", Timezone: "America/New_York"}, {RepoDid: repo.RepoDid.String(), Workflow: "nightly.yml", Expression: "0 3 * * *"}, } if err := d.ReplaceWorkflowSchedules(ctx, ScheduledRepo{ @@ -47,8 +48,11 @@ func TestWorkflowSchedulesTrackIndexedRepositoriesIncludingEmptyManifests(t *tes if err != nil { t.Fatalf("WorkflowSchedules: %v", err) } - if len(persisted) != 2 { - t.Fatalf("persisted schedules = %d, want 2 deduplicated definitions", len(persisted)) + if len(persisted) != 3 { + t.Fatalf("persisted schedules = %d, want 3 deduplicated definitions", len(persisted)) + } + if persisted[0].Timezone != "America/New_York" || persisted[1].Timezone != "UTC" || persisted[2].Timezone != "UTC" { + t.Fatalf("persisted timezones = %q, %q, %q", persisted[0].Timezone, persisted[1].Timezone, persisted[2].Timezone) } for i, schedule := range persisted { if schedule.Branch != "main" || schedule.SHA != "sha-1" { @@ -100,7 +104,7 @@ func TestWorkflowSchedulesTrackIndexedRepositoriesIncludingEmptyManifests(t *tes } } -func TestWorkflowSchedulesMigrationBackfillsIncompleteRevisions(t *testing.T) { +func TestWorkflowSchedulesMigrationDefaultsTimezoneToUTC(t *testing.T) { ctx := context.Background() path := filepath.Join(t.TempDir(), "spindle.db") legacy, err := sql.Open("sqlite3", path) @@ -144,6 +148,17 @@ func TestWorkflowSchedulesMigrationBackfillsIncompleteRevisions(t *testing.T) { t.Fatalf("AddRepo: %v", err) } + var timezone string + if err := d.QueryRow( + `select timezone from workflow_schedules where repo_did = ?`, + repo.RepoDid.String(), + ).Scan(&timezone); err != nil { + t.Fatalf("query migrated timezone: %v", err) + } + if timezone != "UTC" { + t.Fatalf("migrated timezone = %q, want UTC", timezone) + } + var branch, sha string if err := d.QueryRow( `select branch, sha from scheduled_repos where repo_did = ?`, diff --git a/spindle/schedule.go b/spindle/schedule.go index 8a2f233a6..18f3bdab4 100644 --- a/spindle/schedule.go +++ b/spindle/schedule.go @@ -29,6 +29,7 @@ type scheduledWorkflow struct { scheduledRepo db.ScheduledRepo name string expression string + timezone string schedule cron.Schedule next time.Time index int @@ -173,6 +174,7 @@ func (s *pipelineScheduler) restoreSchedules(ctx context.Context) error { repo string workflow string expression string + timezone string }]struct{}) for _, definition := range definitions { repo, ok := repoByDid[definition.RepoDid] @@ -182,21 +184,27 @@ func (s *pipelineScheduler) restoreSchedules(ctx context.Context) error { if _, ok := grouped[definition.RepoDid]; !ok { grouped[definition.RepoDid] = nil } + timezone := definition.Timezone + if timezone == "" { + timezone = "UTC" + } key := struct { repo string workflow string expression string - }{definition.RepoDid, definition.Workflow, definition.Expression} + timezone string + }{definition.RepoDid, definition.Workflow, definition.Expression, timezone} if _, ok := seen[key]; ok { continue } seen[key] = struct{}{} schedule, err := workflow.ParseCron( definition.Expression, + timezone, definition.RepoDid+"\x00"+definition.Workflow, ) if err != nil { - s.l.Warn("ignoring invalid persisted schedule", "repo", definition.RepoDid, "workflow", definition.Workflow, "expression", definition.Expression, "err", err) + s.l.Warn("ignoring invalid persisted schedule", "repo", definition.RepoDid, "workflow", definition.Workflow, "expression", definition.Expression, "timezone", timezone, "err", err) continue } grouped[definition.RepoDid] = append(grouped[definition.RepoDid], scheduledWorkflow{ @@ -208,6 +216,7 @@ func (s *pipelineScheduler) restoreSchedules(ctx context.Context) error { }, name: definition.Workflow, expression: definition.Expression, + timezone: timezone, schedule: schedule, }) } @@ -324,6 +333,7 @@ func (s *pipelineScheduler) RefreshRepo(ctx context.Context, repo db.Repo) error RepoDid: repo.RepoDid.String(), Workflow: entries[i].name, Expression: entries[i].expression, + Timezone: entries[i].timezone, Branch: scheduledRepo.Branch, SHA: scheduledRepo.SHA, }) @@ -549,24 +559,29 @@ func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]schedu } var entries []scheduledWorkflow - seen := make(map[struct{ name, expression string }]struct{}) + seen := make(map[struct{ name, expression, timezone string }]struct{}) for _, wf := range parsed { var workflowEntries []scheduledWorkflow invalid := false for _, constraint := range wf.When { for _, definition := range constraint.Schedule { expression := definition.Cron - key := struct{ name, expression string }{wf.Name, expression} + timezone := definition.Timezone + if timezone == "" { + timezone = "UTC" + } + key := struct{ name, expression, timezone string }{wf.Name, expression, timezone} if _, ok := seen[key]; ok { continue } seen[key] = struct{}{} schedule, err := workflow.ParseCron( expression, + timezone, repo.RepoDid.String()+"\x00"+wf.Name, ) if err != nil { - s.l.Warn("ignoring invalid workflow schedule", "repo", repo.RepoDid, "workflow", wf.Name, "expression", expression, "err", err) + s.l.Warn("ignoring invalid workflow schedule", "repo", repo.RepoDid, "workflow", wf.Name, "expression", expression, "timezone", timezone, "err", err) invalid = true continue } @@ -578,6 +593,7 @@ func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]schedu scheduledRepo: scheduledRepo, name: wf.Name, expression: expression, + timezone: timezone, schedule: schedule, }) } diff --git a/spindle/schedule_test.go b/spindle/schedule_test.go index 0a68a0167..e2ca73be2 100644 --- a/spindle/schedule_test.go +++ b/spindle/schedule_test.go @@ -48,7 +48,7 @@ func scheduleTestSnapshot() db.ScheduledRepo { func mustCron(t *testing.T, expression string) scheduledWorkflow { t.Helper() repo := scheduleTestRepo() - schedule, err := workflow.ParseCron(expression, repo.RepoDid.String()+"\x00build.yml") + schedule, err := workflow.ParseCron(expression, "", repo.RepoDid.String()+"\x00build.yml") if err != nil { t.Fatalf("ParseCron(%q, UTC): %v", expression, err) } @@ -308,7 +308,8 @@ func TestPipelineSchedulerRestoresPersistedSchedulesWithoutLoadingRepository(t * }, []db.WorkflowSchedule{{ RepoDid: repo.RepoDid.String(), Workflow: "build.yml", - Expression: "30 9 * * 1-5", + Expression: "30 5 * * 1-5", + Timezone: "America/New_York", }}, time.Now()); err != nil { t.Fatalf("ReplaceWorkflowSchedules: %v", err) } diff --git a/workflow/def.go b/workflow/def.go index c93fb46bc..0bdf735f5 100644 --- a/workflow/def.go +++ b/workflow/def.go @@ -9,6 +9,7 @@ import ( "strconv" "strings" "time" + _ "time/tzdata" "tangled.org/core/api/tangled" @@ -49,7 +50,8 @@ type ( } Schedule struct { - Cron string `yaml:"cron"` + Cron string `yaml:"cron"` + Timezone string `yaml:"timezone"` // optional; IANA timezone, defaults to UTC } CloneOpts struct { @@ -111,13 +113,13 @@ func HasHashedCron(expression string) bool { return false } -// ParseCron parses a five-field UTC cron expression using seed to +// ParseCron parses a five-field cron expression in timezone using seed to // deterministically expand Jenkins-style hashed fields. -func ParseCron(expression, seed string) (cron.Schedule, error) { - return parseCron(expression, seed, true) +func ParseCron(expression, timezone, seed string) (cron.Schedule, error) { + return parseCron(expression, timezone, seed, true) } -func parseCron(expression, seed string, validateOccurrence bool) (cron.Schedule, error) { +func parseCron(expression, timezone, seed string, validateOccurrence bool) (cron.Schedule, error) { if len(strings.Fields(expression)) != 5 { return nil, fmt.Errorf("expected exactly 5 fields") } @@ -134,7 +136,14 @@ func parseCron(expression, seed string, validateOccurrence bool) (cron.Schedule, expression = strings.Join(fields, " ") } - schedule, err := cronParser.Parse("CRON_TZ=UTC " + expression) + if timezone == "" { + timezone = "UTC" + } + location, err := time.LoadLocation(timezone) + if err != nil { + return nil, fmt.Errorf("invalid timezone %q: %w", timezone, err) + } + schedule, err := cronParser.Parse("CRON_TZ=" + location.String() + " " + expression) if err != nil { return nil, err } @@ -344,7 +353,7 @@ func (w Workflow) validateSchedules() error { if entry.Cron == "" { return fmt.Errorf("when[%d].schedule[%d]: cron is required", i, j) } - if _, err := parseCron(entry.Cron, w.Name, !HasHashedCron(entry.Cron)); err != nil { + if _, err := parseCron(entry.Cron, entry.Timezone, w.Name, !HasHashedCron(entry.Cron)); err != nil { return fmt.Errorf("when[%d].schedule[%d]: invalid cron expression %q: %w", i, j, entry.Cron, err) } } diff --git a/workflow/def_test.go b/workflow/def_test.go index 04e613cec..62646f260 100644 --- a/workflow/def_test.go +++ b/workflow/def_test.go @@ -638,12 +638,14 @@ when: - event: schedule schedule: - cron: "30 2 * * *" + timezone: America/New_York - cron: "0 9 * * 1-5" + timezone: Europe/London `)) assert.NoError(t, err) assert.Equal(t, []Schedule{ - {Cron: "30 2 * * *"}, - {Cron: "0 9 * * 1-5"}, + {Cron: "30 2 * * *", Timezone: "America/New_York"}, + {Cron: "0 9 * * 1-5", Timezone: "Europe/London"}, }, wf.When[0].Schedule) } @@ -668,8 +670,8 @@ func TestUnmarshalWorkflowRejectsInvalidCronSchedule(t *testing.T) { yaml: "when:\n - event: schedule\n", }, { - name: "schedule entry without event", - yaml: "when:\n - event: push\n branch: main\n schedule:\n - cron: \"0 9 * * *\"\n", + name: "schedule entry with timezone without event", + yaml: "when:\n - event: push\n branch: main\n schedule:\n - cron: \"0 9 * * *\"\n timezone: America/New_York\n", }, { name: "schedule entry without cron", @@ -687,6 +689,10 @@ func TestUnmarshalWorkflowRejectsInvalidCronSchedule(t *testing.T) { name: "impossible date", yaml: "when:\n - event: schedule\n schedule:\n - cron: \"0 0 30 2 *\"\n", }, + { + name: "invalid timezone", + yaml: "when:\n - event: schedule\n schedule:\n - cron: \"0 9 * * *\"\n timezone: Mars/Olympus_Mons\n", + }, } for _, test := range tests { @@ -732,31 +738,38 @@ func TestUnmarshalWorkflowRejectsScheduleFilters(t *testing.T) { } } +func TestParseCronUsesTimezone(t *testing.T) { + after := time.Date(2026, time.August, 10, 9, 29, 0, 0, time.UTC) + schedule, err := ParseCron("30 5 * * 1-5", "America/New_York", "timezone-test") + assert.NoError(t, err) + assert.Equal(t, time.Date(2026, time.August, 10, 9, 30, 0, 0, time.UTC), schedule.Next(after)) +} + func TestParseCronKeepsExactExpressionsIndependentOfSeed(t *testing.T) { after := time.Date(2026, time.August, 10, 9, 29, 0, 0, time.UTC) - first, err := ParseCron("30 5 * * 1-5", "first-seed") + first, err := ParseCron("30 5 * * 1-5", "UTC", "first-seed") assert.NoError(t, err) - second, err := ParseCron("30 5 * * 1-5", "second-seed") + second, err := ParseCron("30 5 * * 1-5", "UTC", "second-seed") assert.NoError(t, err) assert.Equal(t, first.Next(after), second.Next(after)) } func TestParseCronPreservesNamedWeekdaysAlongsideHashedTerms(t *testing.T) { after := time.Date(2026, time.January, 1, 2, 59, 0, 0, time.UTC) - named, err := ParseCron("0 3 * * THU", "named-weekday-seed") + named, err := ParseCron("0 3 * * THU", "UTC", "named-weekday-seed") assert.NoError(t, err) assert.Equal(t, time.Thursday, named.Next(after).Weekday()) - mixed, err := ParseCron("0 3 * * THU,H", "named-weekday-seed") + mixed, err := ParseCron("0 3 * * THU,H", "UTC", "named-weekday-seed") assert.NoError(t, err) assert.NotZero(t, mixed.Next(after)) } func TestParseCronHashedFieldsAreStableAndBounded(t *testing.T) { after := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC) - first, err := ParseCron("H H H H *", "stable-seed") + first, err := ParseCron("H H H H *", "UTC", "stable-seed") assert.NoError(t, err) - second, err := ParseCron("H H H H *", "stable-seed") + second, err := ParseCron("H H H H *", "UTC", "stable-seed") assert.NoError(t, err) firstOccurrence := first.Next(after) @@ -777,7 +790,7 @@ func TestParseCronHashedFieldsAreStableAndBounded(t *testing.T) { func TestParseCronHashedRangeAndStep(t *testing.T) { after := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC) - schedule, err := ParseCron("H(10-20)/5 * * * *", "range-step-seed") + schedule, err := ParseCron("H(10-20)/5 * * * *", "UTC", "range-step-seed") assert.NoError(t, err) @@ -794,7 +807,7 @@ func TestParseCronHashedRangeAndStep(t *testing.T) { func TestParseCronHashedStepOneKeepsJenkinsSingleValueSemantics(t *testing.T) { after := time.Date(2026, time.January, 1, 0, 0, 0, 0, time.UTC) - schedule, err := ParseCron("H/1 * * * *", "step-one-seed") + schedule, err := ParseCron("H/1 * * * *", "UTC", "step-one-seed") assert.NoError(t, err) first := schedule.Next(after) @@ -812,7 +825,7 @@ func TestParseCronRejectsMalformedHashedFields(t *testing.T) { "5H * * * *", } { t.Run(expression, func(t *testing.T) { - _, err := ParseCron(expression, "malformed-seed") + _, err := ParseCron(expression, "UTC", "malformed-seed") assert.Error(t, err) }) }