diff --git a/docs/DOCS.md b/docs/DOCS.md index 663c634ca..b6a6dc873 100644 --- a/docs/DOCS.md +++ b/docs/DOCS.md @@ -1174,8 +1174,8 @@ has the following fields: pull request is made or updated. - `manual`: The workflow can be triggered manually. - `schedule`: The workflow runs according to its `cron` - expressions on the repository's latest default-branch - commit. + expressions using the default-branch name and commit + saved when spindle last refreshed the repository's schedules. - `branch`: Defines which branches the workflow should run for. If used with the `push` event, commits to the branch(es) listed here will trigger the workflow. If used @@ -1234,9 +1234,10 @@ when: tag: ["v*", "stable"] ``` -Scheduled workflows run on the latest commit of the default -branch. For example, this runs at a stable minute during the -02:00 UTC hour every day and at 09:00 UTC on weekdays: +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: ```yaml when: @@ -1251,11 +1252,16 @@ repository DID and workflow filename. The minute stays stable across restarts; changing either identity changes the hash. Ranges and steps are supported, for example `H(0-29)/10` or `H/15`. Unbounded day-of-month hashes use days 1–28. -Exact expressions remain exact. +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_SCHEDULE_CONCURRENCY` controls parallel repository +dispatches and startup schedule refreshes (default `4`). +Dispatch finishes each minute's batch before starting the next. + To skip CI for a push, pass a Git push option: ```sh diff --git a/spindle/config/config.go b/spindle/config/config.go index eb0540126..f19fd1bd7 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -207,6 +207,10 @@ type Logging struct { Insecure bool `env:"INSECURE, default=false"` } +type Schedule struct { + Concurrency int `env:"CONCURRENCY, default=4"` +} + type Config struct { Role Role `env:"SPINDLE_ROLE, default=standalone"` Server Server `env:",prefix=SPINDLE_SERVER_"` @@ -220,6 +224,7 @@ type Config struct { Tracing Tracing `env:",prefix=SPINDLE_TRACING_"` Logging Logging `env:",prefix=SPINDLE_LOGGING_"` Cache Cache `env:",prefix=SPINDLE_CACHE_"` + Schedule Schedule `env:",prefix=SPINDLE_SCHEDULE_"` } func (c *Config) validate() error { @@ -244,6 +249,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.Concurrency < 1 { + return fmt.Errorf("SPINDLE_SCHEDULE_CONCURRENCY must be greater than zero") + } if c.NixCache.MaxStagedMiB <= 0 || c.NixCache.MaxStagedMiB > math.MaxInt64>>20 { return fmt.Errorf("SPINDLE_NIX_CACHE_MAX_STAGED_MIB must be greater than zero and fit in bytes") } diff --git a/spindle/config/config_test.go b/spindle/config/config_test.go index 300b8dcca..53df45710 100644 --- a/spindle/config/config_test.go +++ b/spindle/config/config_test.go @@ -68,6 +68,22 @@ func TestLoadRejectsNonPositiveDrainTimeout(t *testing.T) { } } +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 }{ + {"no workers", "SPINDLE_SCHEDULE_CONCURRENCY", "0"}, + {"negative workers", "SPINDLE_SCHEDULE_CONCURRENCY", "-1"}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Setenv(tc.key, tc.value) + if _, err := Load(context.Background()); err == nil { + t.Fatalf("Load accepted %s=%s", tc.key, tc.value) + } + }) + } +} + func TestLoadRequiresJumpHostKey(t *testing.T) { t.Setenv("SPINDLE_SERVER_HOSTNAME", "spindle.example.com") t.Setenv("SPINDLE_SERVER_OWNER", "did:web:spindle.example.com") diff --git a/spindle/db/db.go b/spindle/db/db.go index 29d29cd72..96fd913b8 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -155,6 +155,30 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { created_at_ns integer not null default 0 ); + create table if not exists schedule_runs ( + repo_did text not null, + workflow text not null, + scheduled_at integer not null, + pipeline_id text, + + primary key (repo_did, workflow, scheduled_at) + ); + + create table if not exists scheduled_repos ( + repo_did text primary key, + branch text not null default '', + sha text not null default '', + refreshed_at integer not null + ); + create table if not exists workflow_schedules ( + repo_did text not null, + workflow text not null, + expression text not null, + + primary key (repo_did, workflow, expression), + foreign key (repo_did) references scheduled_repos(repo_did) on delete cascade + ); + create table if not exists workflows ( id integer primary key autoincrement, pipeline_id text not null, @@ -1085,6 +1109,45 @@ func migratePipelineIdentityColumns(tx *sql.Tx) error { return err } } + if err := orm.RunMigration(conn, logger, "events-pipeline-schedule-index", func(tx *sql.Tx) error { + _, err := tx.Exec(` + drop index if exists idx_events_pipeline_lookup; + create index idx_events_pipeline_lookup on events( + coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), + json_extract(event, '$.triggerMetadata.repo.did')), + coalesce(json_extract(event, '$.triggerMetadata.push.newSha'), + json_extract(event, '$.triggerMetadata.pullRequest.sourceSha'), + json_extract(event, '$.triggerMetadata.manual.sha'), + json_extract(event, '$.triggerMetadata.schedule.sha')) + ) where nsid = 'sh.tangled.pipeline'; + `) + 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 + if err := tx.QueryRow( + `select count(*) from pragma_table_info('scheduled_repos') where name = ?`, + column, + ).Scan(&present); err != nil { + return err + } + if present != 0 { + continue + } + if _, err := tx.Exec( + `alter table scheduled_repos add column ` + column + ` text not null default ''`, + ); err != nil { + return err + } + } + return nil + }); err != nil { + return err + } return nil } diff --git a/spindle/db/schedule_runs.go b/spindle/db/schedule_runs.go new file mode 100644 index 000000000..442f0f310 --- /dev/null +++ b/spindle/db/schedule_runs.go @@ -0,0 +1,40 @@ +package db + +import ( + "context" + "time" +) + +func (d *DB) ClaimScheduleRun(ctx context.Context, repoDid, workflow string, scheduledAt time.Time) (bool, error) { + result, err := d.ExecContext(ctx, ` + insert or ignore into schedule_runs (repo_did, workflow, scheduled_at) + values (?, ?, ?) + `, repoDid, workflow, scheduledAt.Unix()) + if err != nil { + return false, err + } + rows, err := result.RowsAffected() + return rows == 1, err +} + +func (d *DB) CompleteScheduleRun(ctx context.Context, repoDid, workflow string, scheduledAt time.Time, pipelineID string) error { + _, err := d.ExecContext(ctx, ` + update schedule_runs + set pipeline_id = ? + where repo_did = ? and workflow = ? and scheduled_at = ? + `, pipelineID, repoDid, workflow, scheduledAt.Unix()) + return err +} + +func (d *DB) ReleaseScheduleRun(ctx context.Context, repoDid, workflow string, scheduledAt time.Time) error { + _, err := d.ExecContext(ctx, ` + delete from schedule_runs + where repo_did = ? and workflow = ? and scheduled_at = ? and pipeline_id is null + `, repoDid, workflow, scheduledAt.Unix()) + return err +} + +func (d *DB) PruneScheduleRuns(ctx context.Context, before time.Time) error { + _, err := d.ExecContext(ctx, `delete from schedule_runs where scheduled_at < ?`, before.Unix()) + return err +} diff --git a/spindle/db/schedules.go b/spindle/db/schedules.go new file mode 100644 index 000000000..94a5b9ecc --- /dev/null +++ b/spindle/db/schedules.go @@ -0,0 +1,117 @@ +package db + +import ( + "context" + "time" +) + +type ScheduledRepo struct { + RepoDid string + Branch string + SHA string +} + +type WorkflowSchedule struct { + RepoDid string + Workflow string + Expression string + Branch string + SHA string +} + +func (d *DB) ReplaceWorkflowSchedules(ctx context.Context, repo ScheduledRepo, schedules []WorkflowSchedule, refreshedAt time.Time) error { + tx, err := d.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + + if _, err := tx.ExecContext(ctx, ` + insert into scheduled_repos (repo_did, branch, sha, refreshed_at) + values (?, ?, ?, ?) + on conflict(repo_did) do update set + branch = excluded.branch, + sha = excluded.sha, + refreshed_at = excluded.refreshed_at + `, repo.RepoDid, repo.Branch, repo.SHA, refreshedAt.Unix()); err != nil { + return err + } + if _, err := tx.ExecContext(ctx, `delete from workflow_schedules where repo_did = ?`, repo.RepoDid); err != nil { + return err + } + for _, schedule := range schedules { + if _, err := tx.ExecContext(ctx, ` + insert or ignore into workflow_schedules (repo_did, workflow, expression) + values (?, ?, ?) + `, repo.RepoDid, schedule.Workflow, schedule.Expression); err != nil { + return err + } + } + return tx.Commit() +} + +func (d *DB) WorkflowSchedules(ctx context.Context) ([]WorkflowSchedule, error) { + rows, err := d.QueryContext(ctx, ` + select ws.repo_did, ws.workflow, ws.expression, + 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 + `) + if err != nil { + return nil, err + } + defer rows.Close() + + var schedules []WorkflowSchedule + for rows.Next() { + var schedule WorkflowSchedule + if err := rows.Scan( + &schedule.RepoDid, + &schedule.Workflow, + &schedule.Expression, + &schedule.Branch, + &schedule.SHA, + ); err != nil { + return nil, err + } + schedules = append(schedules, schedule) + } + return schedules, rows.Err() +} + +func (d *DB) UnindexedScheduleRepos(ctx context.Context) ([]Repo, error) { + rows, err := d.QueryContext(ctx, ` + select knot, owner, rkey, coalesce(repo_did, ''), coalesce(name, '') + from repos + where coalesce(repo_did, '') <> '' + and not exists ( + select 1 from scheduled_repos + where scheduled_repos.repo_did = repos.repo_did + and coalesce(scheduled_repos.branch, '') <> '' + and coalesce(scheduled_repos.sha, '') <> '' + ) + order by repos.id + `) + if err != nil { + return nil, err + } + defer rows.Close() + + var repos []Repo + for rows.Next() { + repo, err := scanRepo(rows) + if err != nil { + return nil, err + } + repos = append(repos, *repo) + } + return repos, rows.Err() +} + +func (d *DB) RemoveWorkflowSchedules(ctx context.Context, repoDid string) error { + _, err := d.ExecContext(ctx, `delete from scheduled_repos where repo_did = ?`, repoDid) + return err +} diff --git a/spindle/db/schedules_test.go b/spindle/db/schedules_test.go new file mode 100644 index 000000000..ba992e48b --- /dev/null +++ b/spindle/db/schedules_test.go @@ -0,0 +1,173 @@ +package db + +import ( + "context" + "database/sql" + "path/filepath" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +func TestWorkflowSchedulesTrackIndexedRepositoriesIncludingEmptyManifests(t *testing.T) { + ctx := context.Background() + d := newTestDB(t) + repo := Repo{ + Knot: "knot.test", + Owner: syntax.DID("did:plc:owner"), + Rkey: syntax.RecordKey("repo"), + RepoDid: syntax.DID("did:plc:repo"), + } + if err := d.AddRepo(repo); err != nil { + t.Fatalf("AddRepo: %v", err) + } + + unindexed, err := d.UnindexedScheduleRepos(ctx) + if err != nil { + t.Fatalf("UnindexedScheduleRepos: %v", err) + } + if len(unindexed) != 1 { + t.Fatalf("unindexed repositories = %d, want 1", len(unindexed)) + } + + 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: "nightly.yml", Expression: "0 3 * * *"}, + } + if err := d.ReplaceWorkflowSchedules(ctx, ScheduledRepo{ + RepoDid: repo.RepoDid.String(), + Branch: "main", + SHA: "sha-1", + }, schedules, time.Now()); err != nil { + t.Fatalf("ReplaceWorkflowSchedules: %v", err) + } + persisted, err := d.WorkflowSchedules(ctx) + if err != nil { + t.Fatalf("WorkflowSchedules: %v", err) + } + if len(persisted) != 2 { + t.Fatalf("persisted schedules = %d, want 2 deduplicated definitions", len(persisted)) + } + for i, schedule := range persisted { + if schedule.Branch != "main" || schedule.SHA != "sha-1" { + t.Fatalf("persisted snapshot at index %d = branch %q, sha %q; want main/sha-1", i, schedule.Branch, schedule.SHA) + } + } + + if err := d.ReplaceWorkflowSchedules(ctx, ScheduledRepo{ + RepoDid: repo.RepoDid.String(), + Branch: "main", + SHA: "sha-2", + }, nil, time.Now()); err != nil { + t.Fatalf("ReplaceWorkflowSchedules(empty): %v", err) + } + persisted, err = d.WorkflowSchedules(ctx) + if err != nil { + t.Fatalf("WorkflowSchedules after empty refresh: %v", err) + } + if len(persisted) != 0 { + t.Fatalf("persisted schedules after empty refresh = %d, want 0", len(persisted)) + } + var branch, sha string + if err := d.QueryRow( + `select branch, sha from scheduled_repos where repo_did = ?`, + repo.RepoDid.String(), + ).Scan(&branch, &sha); err != nil { + t.Fatalf("query snapshot after empty refresh: %v", err) + } + if branch != "main" || sha != "sha-2" { + t.Fatalf("snapshot after empty refresh = branch %q, sha %q; want main/sha-2", branch, sha) + } + unindexed, err = d.UnindexedScheduleRepos(ctx) + if err != nil { + t.Fatalf("UnindexedScheduleRepos after empty refresh: %v", err) + } + if len(unindexed) != 0 { + t.Fatalf("empty but indexed repository returned for backfill: %#v", unindexed) + } + + if err := d.RemoveWorkflowSchedules(ctx, repo.RepoDid.String()); err != nil { + t.Fatalf("RemoveWorkflowSchedules: %v", err) + } + unindexed, err = d.UnindexedScheduleRepos(ctx) + if err != nil { + t.Fatalf("UnindexedScheduleRepos after removal: %v", err) + } + if len(unindexed) != 1 { + t.Fatalf("unindexed repositories after removal = %d, want 1", len(unindexed)) + } +} + +func TestWorkflowSchedulesMigrationBackfillsIncompleteRevisions(t *testing.T) { + ctx := context.Background() + path := filepath.Join(t.TempDir(), "spindle.db") + legacy, err := sql.Open("sqlite3", path) + if err != nil { + t.Fatalf("open legacy database: %v", err) + } + if _, err := legacy.Exec(` + create table scheduled_repos ( + repo_did text primary key, + refreshed_at integer not null + ); + create table workflow_schedules ( + repo_did text not null, + workflow text not null, + expression text not null, + primary key (repo_did, workflow, expression) + ); + insert into scheduled_repos values ('did:plc:repo', 0); + insert into workflow_schedules values ('did:plc:repo', 'build.yml', '0 9 * * *'); + `); err != nil { + legacy.Close() + t.Fatalf("seed legacy database: %v", err) + } + if err := legacy.Close(); err != nil { + t.Fatalf("close legacy database: %v", err) + } + + d, err := Make(ctx, path) + if err != nil { + t.Fatalf("migrate database: %v", err) + } + defer d.Close() + + repo := Repo{ + Knot: "knot.test", + Owner: syntax.DID("did:plc:owner"), + Rkey: syntax.RecordKey("repo"), + RepoDid: syntax.DID("did:plc:repo"), + } + if err := d.AddRepo(repo); err != nil { + t.Fatalf("AddRepo: %v", err) + } + + var branch, sha string + if err := d.QueryRow( + `select branch, sha from scheduled_repos where repo_did = ?`, + repo.RepoDid.String(), + ).Scan(&branch, &sha); err != nil { + t.Fatalf("query migrated revision: %v", err) + } + if branch != "" || sha != "" { + t.Fatalf("migrated revision = branch %q, sha %q, want empty legacy snapshot", branch, sha) + } + + schedules, err := d.WorkflowSchedules(ctx) + if err != nil { + t.Fatalf("WorkflowSchedules: %v", err) + } + if len(schedules) != 0 { + t.Fatalf("migrated incomplete schedules = %#v, want none", schedules) + } + + unindexed, err := d.UnindexedScheduleRepos(ctx) + if err != nil { + t.Fatalf("UnindexedScheduleRepos: %v", err) + } + if len(unindexed) != 1 || unindexed[0].RepoDid != repo.RepoDid { + t.Fatalf("migrated incomplete repos = %#v, want %s", unindexed, repo.RepoDid) + } +} diff --git a/spindle/engine/cache_controller.go b/spindle/engine/cache_controller.go index 5094e0d3d..58e6b5b17 100644 --- a/spindle/engine/cache_controller.go +++ b/spindle/engine/cache_controller.go @@ -169,6 +169,9 @@ func pipelineRef(metadata *tangled.Pipeline_TriggerMetadata) string { if metadata.Manual != nil && metadata.Manual.Ref != nil { return *metadata.Manual.Ref } + if metadata.Schedule != nil { + return metadata.Schedule.Ref + } return "" } diff --git a/spindle/engine/cache_test.go b/spindle/engine/cache_test.go index ed0ea456c..710ef3a17 100644 --- a/spindle/engine/cache_test.go +++ b/spindle/engine/cache_test.go @@ -16,6 +16,7 @@ import ( "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage" @@ -194,6 +195,60 @@ func TestCachePlanPreservesRestoreAtEntryQuota(t *testing.T) { } } +func TestScheduledCacheRestoresPushCacheOnSameRef(t *testing.T) { + ctx := context.Background() + d := newCacheTestDB(t) + repoDid, _ := addCacheTestRepo(t, d) + base := &fakeStorage{objects: make(map[string][]byte)} + controller := NewLocalCacheController(d, base, "", 0, 0, discardLogger) + wf := &models.Workflow{ + Name: "build.yml", Engine: "microvm", + Caches: []models.CacheEntry{{Key: "deps", Paths: []string{"deps"}}}, + } + sha := strings.Repeat("a", 40) + push, err := controller.Plan(ctx, &models.Pipeline{ + RepoDid: repoDid, + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ + Kind: "push", + Push: &tangled.Pipeline_PushTriggerData{Ref: "refs/heads/main", NewSha: sha}, + }, + }, wf) + if err != nil { + t.Fatal(err) + } + if len(push) != 1 || push[0].SaveKey == "" { + t.Fatalf("push did not reserve a cache: %+v", push) + } + store := newTrackedCacheStore(base, controller, discardLogger, push) + if err := store.Put(ctx, push[0].SaveKey, strings.NewReader("push cache")); err != nil { + t.Fatal(err) + } + for _, ref := range []string{"refs/heads/main", "refs/heads/other"} { + t.Run(ref, func(t *testing.T) { + bindings, err := controller.Plan(ctx, &models.Pipeline{ + RepoDid: repoDid, + TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ + Kind: "schedule", + Schedule: &tangled.Pipeline_ScheduleTriggerData{Ref: ref, Sha: sha}, + }, + }, wf) + if err != nil { + t.Fatal(err) + } + if len(bindings) != 1 { + t.Fatalf("cache bindings = %+v", bindings) + } + want := "" + if ref == "refs/heads/main" { + want = push[0].SaveKey + } + if bindings[0].RestoreKey != want { + t.Fatalf("restore key = %q, want %q", bindings[0].RestoreKey, want) + } + }) + } +} + func TestCachePlanKeepsPendingEntryThroughWorkflowDeadline(t *testing.T) { d := newCacheTestDB(t) repoDid, _ := addCacheTestRepo(t, d) diff --git a/spindle/models/clone.go b/spindle/models/clone.go index 1f6be1cb1..4ba629e53 100644 --- a/spindle/models/clone.go +++ b/spindle/models/clone.go @@ -103,6 +103,12 @@ func ExtractCommitSHA(tr tangled.Pipeline_TriggerMetadata) (string, error) { } return tr.Manual.Sha, nil + case workflow.TriggerKindSchedule: + if tr.Schedule == nil { + return "", fmt.Errorf("schedule trigger metadata is nil") + } + return tr.Schedule.Sha, nil + default: return "", fmt.Errorf("unknown trigger kind: %s", tr.Kind) } diff --git a/spindle/models/clone_test.go b/spindle/models/clone_test.go index a5feb8c63..d5b996dc8 100644 --- a/spindle/models/clone_test.go +++ b/spindle/models/clone_test.go @@ -420,3 +420,24 @@ func TestBuildCloneStep_NilCloneOpts(t *testing.T) { t.Error("Commands should contain 'git init'") } } + +func TestBuildCloneStep_ScheduleTrigger(t *testing.T) { + tr := tangled.Pipeline_TriggerMetadata{ + Kind: string(workflow.TriggerKindSchedule), + Schedule: &tangled.Pipeline_ScheduleTriggerData{ + Sha: "scheduledsha123", + Ref: "refs/heads/main", + ScheduledAt: "2026-08-10T09:00:00Z", + }, + Repo: &tangled.Pipeline_TriggerRepo{ + Knot: "example.com", + Did: "did:plc:user123", + RepoDid: sp("did:plc:boltless"), + }, + } + + step := BuildCloneStep(tangled.Pipeline_Workflow{}, tr, false) + if commands := strings.Join(step.Commands(), " "); !strings.Contains(commands, "scheduledsha123") { + t.Fatalf("clone commands do not target scheduled SHA: %s", commands) + } +} diff --git a/spindle/models/pipeline_env.go b/spindle/models/pipeline_env.go index 9935f54d2..5bb7799fe 100644 --- a/spindle/models/pipeline_env.go +++ b/spindle/models/pipeline_env.go @@ -102,7 +102,17 @@ func PipelineEnvVarsForSource(tr *tangled.Pipeline_TriggerMetadata, pipelineId P env["TANGLED_INPUT_"+strings.ToUpper(pair.Key)] = pair.Value } } - } + case workflow.TriggerKindSchedule: + if tr.Schedule != nil { + refName := plumbing.ReferenceName(tr.Schedule.Ref) + env["TANGLED_REF"] = tr.Schedule.Ref + env["TANGLED_REF_NAME"] = refName.Short() + env["TANGLED_REF_TYPE"] = "branch" + env["TANGLED_SHA"] = tr.Schedule.Sha + env["TANGLED_COMMIT_SHA"] = tr.Schedule.Sha + env["TANGLED_SCHEDULED_AT"] = tr.Schedule.ScheduledAt + } + } return env } diff --git a/spindle/models/pipeline_env_test.go b/spindle/models/pipeline_env_test.go index fd57b69df..70c8fae9e 100644 --- a/spindle/models/pipeline_env_test.go +++ b/spindle/models/pipeline_env_test.go @@ -288,3 +288,35 @@ func TestPipelineEnvVars_NilPushData(t *testing.T) { t.Error("Should not have TANGLED_REF when push data is nil") } } + +func TestPipelineEnvVars_Schedule(t *testing.T) { + tr := &tangled.Pipeline_TriggerMetadata{ + Kind: string(workflow.TriggerKindSchedule), + Schedule: &tangled.Pipeline_ScheduleTriggerData{ + Sha: "scheduledsha123", + Ref: "refs/heads/main", + ScheduledAt: "2026-08-10T09:00:00Z", + }, + Repo: &tangled.Pipeline_TriggerRepo{ + Knot: "example.com", + Did: "did:plc:user123", + RepoDid: sp("did:plc:boltless"), + }, + } + + env := PipelineEnvVars(tr, PipelineId{Knot: "example.com", Rkey: "pipeline"}) + expected := map[string]string{ + "TANGLED_PIPELINE_KIND": "schedule", + "TANGLED_REF": "refs/heads/main", + "TANGLED_REF_NAME": "main", + "TANGLED_REF_TYPE": "branch", + "TANGLED_SHA": "scheduledsha123", + "TANGLED_COMMIT_SHA": "scheduledsha123", + "TANGLED_SCHEDULED_AT": "2026-08-10T09:00:00Z", + } + for key, value := range expected { + if env[key] != value { + t.Errorf("%s = %q, want %q", key, env[key], value) + } + } +} diff --git a/spindle/schedule.go b/spindle/schedule.go new file mode 100644 index 000000000..c9ce93ccb --- /dev/null +++ b/spindle/schedule.go @@ -0,0 +1,663 @@ +package spindle + +import ( + "container/heap" + "context" + "fmt" + "log/slog" + "slices" + "sync" + "time" + + indigoxrpc "github.com/bluesky-social/indigo/xrpc" + "github.com/go-git/go-git/v5/plumbing" + "github.com/robfig/cron/v3" + "tangled.org/core/api/tangled" + "tangled.org/core/log" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" + "tangled.org/core/workflow" +) + +const ( + scheduleRunPruneInterval = 6 * time.Hour + scheduleDispatchBacklog = 60 +) + +type scheduledWorkflow struct { + repo db.Repo + scheduledRepo db.ScheduledRepo + name string + expression string + schedule cron.Schedule + next time.Time + index int +} + +type scheduleHeap []*scheduledWorkflow + +func (h scheduleHeap) Len() int { return len(h) } +func (h scheduleHeap) Less(i, j int) bool { return h[i].next.Before(h[j].next) } +func (h scheduleHeap) Swap(i, j int) { + h[i], h[j] = h[j], h[i] + h[i].index = i + h[j].index = j +} +func (h *scheduleHeap) Push(value any) { + entry := value.(*scheduledWorkflow) + entry.index = len(*h) + *h = append(*h, entry) +} +func (h *scheduleHeap) Pop() any { + old := *h + last := len(old) - 1 + entry := old[last] + old[last] = nil + entry.index = -1 + *h = old[:last] + return entry +} + +type scheduleLoader func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) +type scheduleDispatcher func(context.Context, db.Repo, db.ScheduledRepo, []string, time.Time) (models.PipelineId, error) + +type pipelineScheduler struct { + db *db.DB + l *slog.Logger + load scheduleLoader + dispatch scheduleDispatcher + now func() time.Time + + concurrency int + + mu sync.Mutex + byRepo map[string][]*scheduledWorkflow + queue scheduleHeap + + refreshMu sync.Mutex + refreshLocks map[string]*sync.Mutex + dispatchMu sync.Mutex +} + +func newPipelineScheduler(s *Spindle) *pipelineScheduler { + 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), + } +} + +func (s *pipelineScheduler) Start(ctx context.Context) { + if err := s.restoreSchedules(ctx); err != nil { + s.l.Error("failed to restore schedules", "err", err) + } + go s.backfillSchedules(ctx) + go s.run(ctx) +} + +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 { + s.l.Error("failed to prune schedule claims", "err", err) + } + pruneTicker := time.NewTicker(scheduleRunPruneInterval) + defer pruneTicker.Stop() + + timer := time.NewTimer(time.Until(nextScheduleMinute(s.now()))) + defer timer.Stop() + + for { + select { + case <-ctx.Done(): + return + case firedAt := <-timer.C: + scheduledAt := firedAt.UTC().Truncate(time.Minute) + select { + case dispatches <- scheduledAt: + default: + s.l.Error("schedule dispatch backlog is full", "scheduledAt", scheduledAt) + } + 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 { + s.l.Error("failed to prune schedule claims", "err", err) + } + } + } +} + +func nextScheduleMinute(now time.Time) time.Time { + return now.UTC().Truncate(time.Minute).Add(time.Minute) +} + +func (s *pipelineScheduler) runDispatches(ctx context.Context, dispatches <-chan time.Time) { + for { + select { + case <-ctx.Done(): + return + case scheduledAt, ok := <-dispatches: + if !ok { + return + } + s.dispatchDue(ctx, scheduledAt) + } + } +} + +func (s *pipelineScheduler) restoreSchedules(ctx context.Context) error { + repos, err := s.db.AllRepos() + if err != nil { + return fmt.Errorf("listing repositories: %w", err) + } + repoByDid := make(map[string]db.Repo, len(repos)) + for _, repo := range repos { + repoByDid[repo.RepoDid.String()] = repo + } + + definitions, err := s.db.WorkflowSchedules(ctx) + if err != nil { + return fmt.Errorf("listing workflow schedules: %w", err) + } + grouped := make(map[string][]scheduledWorkflow) + seen := make(map[struct { + repo string + workflow string + expression string + }]struct{}) + for _, definition := range definitions { + repo, ok := repoByDid[definition.RepoDid] + if !ok { + continue + } + if _, ok := grouped[definition.RepoDid]; !ok { + grouped[definition.RepoDid] = nil + } + key := struct { + repo string + workflow string + expression string + }{definition.RepoDid, definition.Workflow, definition.Expression} + if _, ok := seen[key]; ok { + continue + } + seen[key] = struct{}{} + schedule, err := workflow.ParseCron( + definition.Expression, + 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) + continue + } + grouped[definition.RepoDid] = append(grouped[definition.RepoDid], scheduledWorkflow{ + repo: repo, + scheduledRepo: db.ScheduledRepo{ + RepoDid: definition.RepoDid, + Branch: definition.Branch, + SHA: definition.SHA, + }, + name: definition.Workflow, + expression: definition.Expression, + schedule: schedule, + }) + } + now := s.now() + for repoDid, entries := range grouped { + s.replaceCachedSchedules(repoDid, entries, now) + } + s.l.Info("restored persisted schedules", "repositories", len(grouped), "schedules", len(definitions)) + return nil +} + +func (s *pipelineScheduler) backfillSchedules(ctx context.Context) { + repos, err := s.db.UnindexedScheduleRepos(ctx) + if err != nil { + s.l.Error("failed to list repositories needing schedule backfill", "err", err) + return + } + if len(repos) == 0 { + return + } + workers := s.workerLimit(len(repos)) + jobs := make(chan db.Repo) + var wg sync.WaitGroup + wg.Add(workers) + for range workers { + go func() { + defer wg.Done() + for { + select { + case <-ctx.Done(): + return + case repo, ok := <-jobs: + if !ok || ctx.Err() != nil { + return + } + if err := s.RefreshRepo(ctx, repo); err != nil { + s.l.Warn("failed to backfill repository schedules", "repo", repo.RepoDid, "err", err) + } + } + } + }() + } + for _, repo := range repos { + select { + case <-ctx.Done(): + close(jobs) + wg.Wait() + return + case jobs <- repo: + } + } + close(jobs) + wg.Wait() +} + +func (s *pipelineScheduler) workerLimit(work int) int { + limit := s.concurrency + if limit < 1 { + limit = 1 + } + if work > 0 && limit > work { + return work + } + return limit +} + +func (s *pipelineScheduler) lockRefresh(repoDid string) func() { + s.refreshMu.Lock() + if s.refreshLocks == nil { + s.refreshLocks = make(map[string]*sync.Mutex) + } + lock, ok := s.refreshLocks[repoDid] + if !ok { + lock = &sync.Mutex{} + s.refreshLocks[repoDid] = lock + } + s.refreshMu.Unlock() + + lock.Lock() + return lock.Unlock +} + +func (s *pipelineScheduler) RefreshRepo(ctx context.Context, repo db.Repo) error { + unlock := s.lockRefresh(repo.RepoDid.String()) + defer unlock() + + current, err := s.db.GetRepoByDid(repo.RepoDid) + if err != nil { + s.removeCachedSchedules(repo.RepoDid.String()) + return fmt.Errorf("checking repository registration: %w", err) + } + repo = *current + entries, scheduledRepo, err := s.load(ctx, repo) + if err != nil { + s.removeCachedSchedules(repo.RepoDid.String()) + return err + } + if scheduledRepo.RepoDid == "" { + scheduledRepo.RepoDid = repo.RepoDid.String() + } + if scheduledRepo.RepoDid != repo.RepoDid.String() { + s.removeCachedSchedules(repo.RepoDid.String()) + return fmt.Errorf("loaded schedule snapshot for %s, want %s", scheduledRepo.RepoDid, repo.RepoDid) + } + if scheduledRepo.Branch == "" || scheduledRepo.SHA == "" { + s.removeCachedSchedules(repo.RepoDid.String()) + return fmt.Errorf("loaded schedule snapshot for %s is incomplete", repo.RepoDid) + } + definitions := make([]db.WorkflowSchedule, 0, len(entries)) + for i := range entries { + entries[i].repo = repo + entries[i].scheduledRepo = scheduledRepo + definitions = append(definitions, db.WorkflowSchedule{ + RepoDid: repo.RepoDid.String(), + Workflow: entries[i].name, + Expression: entries[i].expression, + Branch: scheduledRepo.Branch, + SHA: scheduledRepo.SHA, + }) + } + refreshedAt := s.now() + if err := s.db.ReplaceWorkflowSchedules(ctx, scheduledRepo, definitions, refreshedAt); err != nil { + return fmt.Errorf("persisting repository schedules: %w", err) + } + s.replaceCachedSchedules(repo.RepoDid.String(), entries, refreshedAt) + s.l.Debug("repository schedules refreshed", "repo", repo.RepoDid, "count", len(entries)) + return nil +} + +func (s *pipelineScheduler) RemoveRepo(ctx context.Context, repoDid string) error { + unlock := s.lockRefresh(repoDid) + defer unlock() + + err := s.db.RemoveWorkflowSchedules(ctx, repoDid) + s.removeCachedSchedules(repoDid) + return err +} + +type dueRepository struct { + repo db.Repo + scheduledRepo db.ScheduledRepo + workflows map[string]struct{} +} + +func (s *pipelineScheduler) replaceCachedSchedules(repoDid string, entries []scheduledWorkflow, after time.Time) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.byRepo == nil { + s.byRepo = make(map[string][]*scheduledWorkflow) + } + for _, entry := range s.byRepo[repoDid] { + if entry.index >= 0 { + heap.Remove(&s.queue, entry.index) + } + } + delete(s.byRepo, repoDid) + + after = after.UTC().Truncate(time.Minute) + cached := make([]*scheduledWorkflow, 0, len(entries)) + for i := range entries { + entry := new(scheduledWorkflow) + *entry = entries[i] + entry.next = entry.schedule.Next(after) + if entry.next.IsZero() || !entry.next.After(after) { + continue + } + heap.Push(&s.queue, entry) + cached = append(cached, entry) + } + if len(cached) > 0 { + s.byRepo[repoDid] = cached + } +} + +func (s *pipelineScheduler) removeCachedSchedules(repoDid string) { + s.mu.Lock() + defer s.mu.Unlock() + for _, entry := range s.byRepo[repoDid] { + if entry.index >= 0 { + heap.Remove(&s.queue, entry.index) + } + } + delete(s.byRepo, repoDid) +} + +func (s *pipelineScheduler) due(at time.Time) map[string]dueRepository { + at = at.UTC().Truncate(time.Minute) + s.mu.Lock() + defer s.mu.Unlock() + + due := make(map[string]dueRepository) + for s.queue.Len() > 0 { + entry := s.queue[0] + if entry.next.After(at) { + break + } + heap.Pop(&s.queue) + if entry.next.Equal(at) { + key := entry.repo.RepoDid.String() + group := due[key] + if group.workflows == nil { + group = dueRepository{ + repo: entry.repo, + scheduledRepo: entry.scheduledRepo, + workflows: make(map[string]struct{}), + } + } else if group.scheduledRepo.Branch == "" && group.scheduledRepo.SHA == "" { + group.scheduledRepo = entry.scheduledRepo + } + group.workflows[entry.name] = struct{}{} + due[key] = group + } + entry.next = entry.schedule.Next(at) + if entry.next.IsZero() || !entry.next.After(at) { + continue + } + heap.Push(&s.queue, entry) + } + return due +} + +func (s *pipelineScheduler) dispatchDue(ctx context.Context, at time.Time) { + at = at.UTC().Truncate(time.Minute) + s.dispatchMu.Lock() + defer s.dispatchMu.Unlock() + + due := s.due(at) + if len(due) == 0 { + return + } + + workers := s.workerLimit(len(due)) + jobs := make(chan dueRepository) + var wg sync.WaitGroup + wg.Add(workers) + for range workers { + go func() { + defer wg.Done() + for { + select { + case <-ctx.Done(): + return + case group, ok := <-jobs: + if !ok || ctx.Err() != nil { + return + } + s.dispatchRepository(ctx, group, at) + } + } + }() + } +send: + for _, group := range due { + select { + case <-ctx.Done(): + break send + case jobs <- group: + } + } + close(jobs) + wg.Wait() +} + +func (s *pipelineScheduler) dispatchRepository(ctx context.Context, group dueRepository, at time.Time) { + repoDid := group.repo.RepoDid.String() + names := make([]string, 0, len(group.workflows)) + for name := range group.workflows { + names = append(names, name) + } + slices.Sort(names) + + claimed := names[:0] + for _, name := range names { + ok, err := s.db.ClaimScheduleRun(ctx, repoDid, name, at) + if err != nil { + s.l.Error("failed to claim scheduled workflow", "repo", repoDid, "workflow", name, "scheduledAt", at, "err", err) + continue + } + if ok { + claimed = append(claimed, name) + } + } + if len(claimed) == 0 { + return + } + + pipelineID, err := s.dispatch(ctx, group.repo, group.scheduledRepo, claimed, at) + if err != nil { + s.l.Error("scheduled pipeline dispatch failed", "repo", repoDid, "scheduledAt", at, "err", err) + if pipelineID.Rkey == "" { + cleanupCtx := ctx + if ctx.Err() != nil { + cleanupCtx = context.WithoutCancel(ctx) + } + for _, name := range claimed { + if releaseErr := s.db.ReleaseScheduleRun(cleanupCtx, repoDid, name, at); releaseErr != nil { + s.l.Error("failed to release schedule claim", "repo", repoDid, "workflow", name, "scheduledAt", at, "err", releaseErr) + } + } + return + } + } + + pipelineAt := "" + if pipelineID.Rkey != "" { + pipelineAt = pipelineID.AtUri().String() + } + cleanupCtx := ctx + if ctx.Err() != nil { + cleanupCtx = context.WithoutCancel(ctx) + } + for _, name := range claimed { + if err := s.db.CompleteScheduleRun(cleanupCtx, repoDid, name, at, pipelineAt); err != nil { + s.l.Error("failed to complete schedule claim", "repo", repoDid, "workflow", name, "scheduledAt", at, "err", err) + } + } +} + +func (s *Spindle) loadRepoSchedules(ctx context.Context, repo db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { + branch, err := s.getRepoDefaultBranch(ctx, repo) + if err != nil { + return nil, db.ScheduledRepo{}, err + } + scheduledRepo := db.ScheduledRepo{ + RepoDid: repo.RepoDid.String(), + Branch: branch.Name, + SHA: branch.Hash, + } + raw, err := s.loadPipeline(ctx, s.newRepoCloneUrl(repo.Knot, repo.RepoDid), s.newRepoPath(repo.RepoDid), branch.Hash) + if err != nil { + return nil, scheduledRepo, err + } + + compiler := workflow.Compiler{} + parsed := compiler.Parse(raw) + for _, diagnostic := range compiler.Diagnostics.Errors { + s.l.Warn("scheduled workflow manifest is invalid", "repo", repo.RepoDid, "diagnostic", diagnostic.String()) + } + + var entries []scheduledWorkflow + seen := make(map[struct{ name, expression 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} + if _, ok := seen[key]; ok { + continue + } + seen[key] = struct{}{} + schedule, err := workflow.ParseCron( + expression, + 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) + invalid = true + continue + } + 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) + } + workflowEntries = append(workflowEntries, scheduledWorkflow{ + repo: repo, + scheduledRepo: scheduledRepo, + name: wf.Name, + expression: expression, + schedule: schedule, + }) + } + } + if invalid { + s.l.Warn("ignoring workflow with invalid schedule", "repo", repo.RepoDid, "workflow", wf.Name) + continue + } + entries = append(entries, workflowEntries...) + } + return entries, scheduledRepo, nil +} + +func (s *Spindle) getRepoDefaultBranch(ctx context.Context, repo db.Repo) (*tangled.RepoGetDefaultBranch_Output, error) { + scheme := "https" + if s.cfg.Server.Dev { + scheme = "http" + } + client := &indigoxrpc.Client{Host: fmt.Sprintf("%s://%s", scheme, repo.Knot)} + branch, err := tangled.RepoGetDefaultBranch(ctx, client, repo.RepoDid.String()) + if err != nil { + return nil, fmt.Errorf("resolving default branch for %s: %w", repo.RepoDid, err) + } + if branch.Name == "" { + return nil, fmt.Errorf("default branch for %s has no name", repo.RepoDid) + } + if branch.Hash == "" { + details, detailsErr := tangled.RepoBranch(ctx, client, branch.Name, repo.RepoDid.String()) + if detailsErr != nil { + return nil, fmt.Errorf("resolving default branch details for %s: %w", repo.RepoDid, detailsErr) + } + branch.Hash = details.Hash + branch.Message = details.Message + branch.ShortHash = details.ShortHash + branch.When = details.When + } + if branch.Hash == "" { + return nil, fmt.Errorf("default branch for %s has no commit hash", repo.RepoDid) + } + return branch, nil +} + +func (s *Spindle) dispatchScheduledPipeline(ctx context.Context, repo db.Repo, scheduledRepo db.ScheduledRepo, workflows []string, scheduledAt time.Time) (models.PipelineId, error) { + repoDid := repo.RepoDid.String() + if scheduledRepo.RepoDid == "" { + scheduledRepo.RepoDid = repoDid + } + if scheduledRepo.RepoDid != repoDid { + return models.PipelineId(""), fmt.Errorf("scheduled repository snapshot is for %s, want %s", scheduledRepo.RepoDid, repoDid) + } + if scheduledRepo.Branch == "" || scheduledRepo.SHA == "" { + return models.PipelineId(""), fmt.Errorf("scheduled repository snapshot for %s is incomplete", repoDid) + } + rkey := repo.Rkey.String() + ref := plumbing.NewBranchReferenceName(scheduledRepo.Branch).String() + triggerRepo := &tangled.Pipeline_TriggerRepo{ + Did: repo.Owner.String(), + Knot: repo.Knot, + Repo: &rkey, + RepoDid: &repoDid, + DefaultBranch: scheduledRepo.Branch, + } + trigger := tangled.Pipeline_TriggerMetadata{ + Kind: string(workflow.TriggerKindSchedule), + Repo: triggerRepo, + Schedule: &tangled.Pipeline_ScheduleTriggerData{ + Ref: ref, + ScheduledAt: scheduledAt.UTC().Truncate(time.Minute).Format(time.RFC3339), + Sha: scheduledRepo.SHA, + }, + } + return s.runPipeline( + ctx, + repo.RepoDid, + trigger, + nil, + s.newRepoCloneUrl(repo.Knot, repo.RepoDid), + s.newRepoPath(repo.RepoDid), + scheduledRepo.SHA, + workflows, + triggerRepo, + ) +} diff --git a/spindle/schedule_test.go b/spindle/schedule_test.go new file mode 100644 index 000000000..0a68a0167 --- /dev/null +++ b/spindle/schedule_test.go @@ -0,0 +1,556 @@ +package spindle + +import ( + "context" + "errors" + "io" + "log/slog" + "path/filepath" + "slices" + "sync" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" + "tangled.org/core/workflow" +) + +func newScheduleTestDB(t *testing.T) *db.DB { + t.Helper() + d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db")) + if err != nil { + t.Fatalf("db.Make: %v", err) + } + t.Cleanup(func() { _ = d.Close() }) + return d +} + +func scheduleTestRepo() db.Repo { + return db.Repo{ + Knot: "knot.test", + Owner: syntax.DID("did:plc:owner"), + Rkey: syntax.RecordKey("repo"), + RepoDid: syntax.DID("did:plc:repo"), + } +} + +func scheduleTestSnapshot() db.ScheduledRepo { + repo := scheduleTestRepo() + return db.ScheduledRepo{ + RepoDid: repo.RepoDid.String(), + Branch: "main", + SHA: "sha", + } +} + +func mustCron(t *testing.T, expression string) scheduledWorkflow { + t.Helper() + repo := scheduleTestRepo() + schedule, err := workflow.ParseCron(expression, repo.RepoDid.String()+"\x00build.yml") + if err != nil { + t.Fatalf("ParseCron(%q, UTC): %v", expression, err) + } + return scheduledWorkflow{ + repo: repo, + scheduledRepo: db.ScheduledRepo{ + RepoDid: repo.RepoDid.String(), + Branch: "main", + SHA: "sha", + }, + expression: expression, + schedule: schedule, + } +} + +func TestPipelineSchedulerDueUsesUTCMinuteAndDeduplicatesWorkflows(t *testing.T) { + weekdayMorning := mustCron(t, "0 9 * * 1-5") + weekdayMorning.name = "build.yml" + duplicate := mustCron(t, "0 9 * * *") + duplicate.name = "build.yml" + noon := mustCron(t, "0 12 * * *") + noon.name = "noon.yml" + + scheduler := &pipelineScheduler{ + byRepo: make(map[string][]*scheduledWorkflow), + concurrency: 1, + } + scheduler.replaceCachedSchedules( + scheduleTestRepo().RepoDid.String(), + []scheduledWorkflow{weekdayMorning, duplicate, noon}, + time.Date(2026, time.August, 10, 8, 59, 0, 0, time.UTC), + ) + at := time.Date(2026, time.August, 10, 11, 0, 47, 0, time.FixedZone("local", 2*60*60)) + due := scheduler.due(at) + group, ok := due[scheduleTestRepo().RepoDid.String()] + if !ok { + t.Fatal("expected repository to have due workflows") + } + var names []string + for name := range group.workflows { + names = append(names, name) + } + slices.Sort(names) + if !slices.Equal(names, []string{"build.yml"}) { + t.Fatalf("due workflows = %v, want [build.yml]", names) + } +} + +func TestPipelineSchedulerDispatchDueClaimsOccurrenceOnce(t *testing.T) { + d := newScheduleTestDB(t) + entry := mustCron(t, "0 9 * * *") + entry.name = "build.yml" + at := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) + + var calls int + dispatch := func(_ context.Context, repo db.Repo, _ db.ScheduledRepo, workflows []string, scheduledAt time.Time) (models.PipelineId, error) { + calls++ + if repo.RepoDid != entry.repo.RepoDid || !slices.Equal(workflows, []string{"build.yml"}) || !scheduledAt.Equal(at) { + t.Fatalf("unexpected dispatch: repo=%s workflows=%v scheduledAt=%s", repo.RepoDid, workflows, scheduledAt) + } + return models.PipelineId{Knot: repo.Knot, Rkey: "pipeline"}, nil + } + for range 2 { + scheduler := &pipelineScheduler{ + db: d, + l: slog.New(slog.NewTextHandler(io.Discard, nil)), + byRepo: make(map[string][]*scheduledWorkflow), + concurrency: 1, + dispatch: dispatch, + } + + scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) + scheduler.dispatchDue(context.Background(), at) + } + if calls != 1 { + t.Fatalf("dispatch calls = %d, want 1", calls) + } +} +func TestPipelineSchedulerBoundsConcurrentRepositoryDispatches(t *testing.T) { + at := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) + repoA := scheduleTestRepo() + repoA.RepoDid = syntax.DID("did:plc:repoa") + repoB := scheduleTestRepo() + repoB.RepoDid = syntax.DID("did:plc:repob") + entryA := mustCron(t, "0 9 * * *") + entryA.repo = repoA + entryA.scheduledRepo = db.ScheduledRepo{RepoDid: repoA.RepoDid.String(), Branch: "main", SHA: "sha-a"} + entryA.name = "build.yml" + entryB := mustCron(t, "0 9 * * *") + entryB.repo = repoB + entryB.scheduledRepo = db.ScheduledRepo{RepoDid: repoB.RepoDid.String(), Branch: "main", SHA: "sha-b"} + entryB.name = "build.yml" + + d := newScheduleTestDB(t) + started := make(chan struct{}, 2) + release := make(chan struct{}) + 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), + dispatch: func(_ context.Context, repo db.Repo, _ db.ScheduledRepo, _ []string, _ time.Time) (models.PipelineId, error) { + mu.Lock() + active++ + if active > maxActive { + maxActive = active + } + mu.Unlock() + started <- struct{}{} + <-release + mu.Lock() + active-- + mu.Unlock() + return models.PipelineId{Knot: repo.Knot, Rkey: "pipeline"}, nil + }, + } + scheduler.replaceCachedSchedules(repoA.RepoDid.String(), []scheduledWorkflow{entryA}, at.Add(-time.Minute)) + scheduler.replaceCachedSchedules(repoB.RepoDid.String(), []scheduledWorkflow{entryB}, at.Add(-time.Minute)) + + done := make(chan struct{}) + go func() { + scheduler.dispatchDue(context.Background(), at) + close(done) + }() + for range 2 { + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("repository dispatches did not start concurrently") + } + } + mu.Lock() + gotMaxActive := maxActive + mu.Unlock() + if gotMaxActive != 2 { + t.Fatalf("maximum concurrent dispatches = %d, want 2", gotMaxActive) + } + close(release) + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("dispatchDue did not wait for repository dispatches") + } +} + +func TestPipelineSchedulerReleasesClaimWhenDispatchIsCanceled(t *testing.T) { + at := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) + entry := mustCron(t, "0 9 * * *") + entry.name = "build.yml" + d := newScheduleTestDB(t) + started := make(chan struct{}) + 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), + dispatch: func(ctx context.Context, _ db.Repo, _ db.ScheduledRepo, _ []string, _ time.Time) (models.PipelineId, error) { + close(started) + <-ctx.Done() + return models.PipelineId{}, ctx.Err() + }, + } + scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) + done := make(chan struct{}) + go func() { + scheduler.dispatchDue(ctx, at) + close(done) + }() + <-started + cancel() + select { + case <-done: + 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) + if err != nil { + t.Fatal(err) + } + if !reclaimed { + t.Fatal("canceled dispatch left its schedule claim") + } +} + +func TestPipelineSchedulerDispatchFailureClaimLifecycle(t *testing.T) { + tests := []struct { + name string + pipelineID models.PipelineId + wantReleased bool + }{ + {name: "before pipeline creation", wantReleased: true}, + { + name: "after pipeline creation", + pipelineID: models.PipelineId{Knot: "knot.test", Rkey: "pipeline"}, + wantReleased: false, + }, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + d := newScheduleTestDB(t) + entry := mustCron(t, "0 9 * * *") + entry.name = "build.yml" + 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), + dispatch: func(context.Context, db.Repo, db.ScheduledRepo, []string, time.Time) (models.PipelineId, error) { + return test.pipelineID, errors.New("temporary failure") + }, + } + 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) + if err != nil { + t.Fatal(err) + } + if reclaimed != test.wantReleased { + t.Fatalf("schedule claim reclaimed = %v, want %v", reclaimed, test.wantReleased) + } + if test.wantReleased { + return + } + + var completed int + if err := d.QueryRow( + `select count(*) from schedule_runs where repo_did = ? and workflow = ? and scheduled_at = ? and pipeline_id = ?`, + entry.repo.RepoDid.String(), entry.name, at.Unix(), test.pipelineID.AtUri().String(), + ).Scan(&completed); err != nil { + t.Fatal(err) + } + if completed != 1 { + t.Fatal("pipeline creation error left its schedule claim incomplete") + } + }) + } +} + +func TestPipelineSchedulerRestoresPersistedSchedulesWithoutLoadingRepository(t *testing.T) { + ctx := context.Background() + d := newScheduleTestDB(t) + repo := scheduleTestRepo() + if err := d.AddRepo(repo); err != nil { + t.Fatalf("AddRepo: %v", err) + } + if err := d.ReplaceWorkflowSchedules(ctx, db.ScheduledRepo{ + RepoDid: repo.RepoDid.String(), + Branch: "main", + SHA: "persisted-sha", + }, []db.WorkflowSchedule{{ + RepoDid: repo.RepoDid.String(), + Workflow: "build.yml", + Expression: "30 9 * * 1-5", + }}, time.Now()); err != nil { + t.Fatalf("ReplaceWorkflowSchedules: %v", err) + } + + 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), + } + if err := scheduler.restoreSchedules(ctx); err != nil { + t.Fatalf("restoreSchedules: %v", err) + } + due := scheduler.due(at) + group, ok := due[repo.RepoDid.String()] + if !ok { + t.Fatal("persisted schedule was not restored") + } + if group.scheduledRepo.Branch != "main" || group.scheduledRepo.SHA != "persisted-sha" { + t.Fatalf("restored repository snapshot = %#v, want main/persisted-sha", group.scheduledRepo) + } +} + +type countingSchedule struct { + next time.Time + calls int +} + +func (s *countingSchedule) Next(time.Time) time.Time { + s.calls++ + return s.next +} + +func TestPipelineSchedulerDoesNotEvaluateNonDueSchedulesEachMinute(t *testing.T) { + now := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) + entries := make([]scheduledWorkflow, 10_000) + schedules := make([]*countingSchedule, len(entries)) + for i := range entries { + schedules[i] = &countingSchedule{next: now.Add(24 * time.Hour)} + entries[i] = scheduledWorkflow{ + repo: scheduleTestRepo(), + scheduledRepo: scheduleTestSnapshot(), + name: "build.yml", + schedule: schedules[i], + } + } + + scheduler := &pipelineScheduler{ + byRepo: make(map[string][]*scheduledWorkflow), + concurrency: 1, + } + scheduler.replaceCachedSchedules(scheduleTestRepo().RepoDid.String(), entries, now) + for _, schedule := range schedules { + schedule.calls = 0 + } + if due := scheduler.due(now.Add(time.Minute)); len(due) != 0 { + t.Fatalf("due repositories = %d, want 0", len(due)) + } + for i, schedule := range schedules { + if schedule.calls != 0 { + t.Fatalf("schedule %d evaluated %d times while not due", i, schedule.calls) + } + } +} + +type sequenceSchedule struct { + next []time.Time +} + +func (s *sequenceSchedule) Next(time.Time) time.Time { + if len(s.next) == 0 { + return time.Time{} + } + next := s.next[0] + s.next = s.next[1:] + return next +} + +func TestPipelineSchedulerDropsSchedulesThatDoNotAdvance(t *testing.T) { + at := time.Date(2026, time.August, 10, 9, 0, 0, 0, time.UTC) + tests := []struct { + name string + next []time.Time + due bool + }{ + {name: "zero initial occurrence"}, + {name: "zero after due occurrence", next: []time.Time{at, {}}, due: true}, + {name: "same occurrence repeats", next: []time.Time{at, at}, due: true}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + entry := scheduledWorkflow{ + repo: scheduleTestRepo(), + scheduledRepo: scheduleTestSnapshot(), + name: "build.yml", + schedule: &sequenceSchedule{next: test.next}, + } + scheduler := &pipelineScheduler{ + byRepo: make(map[string][]*scheduledWorkflow), + concurrency: 1, + } + scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) + due := scheduler.due(at) + if got := len(due) > 0; got != test.due { + t.Fatalf("due = %v, want %v", got, test.due) + } + if scheduler.queue.Len() != 0 { + t.Fatalf("queue length = %d, want 0 after non-advancing occurrence", scheduler.queue.Len()) + } + }) + } +} + +func TestPipelineSchedulerSerializesRepositoryRefreshes(t *testing.T) { + ctx := context.Background() + d := newScheduleTestDB(t) + repo := scheduleTestRepo() + if err := d.AddRepo(repo); err != nil { + t.Fatalf("AddRepo: %v", err) + } + oldEntry := mustCron(t, "0 8 * * *") + oldEntry.name = "old.yml" + newEntry := mustCron(t, "0 9 * * *") + newEntry.name = "new.yml" + + firstStarted := make(chan struct{}) + releaseFirst := make(chan struct{}) + secondLoaded := make(chan struct{}) + var releaseOnce sync.Once + release := func() { releaseOnce.Do(func() { close(releaseFirst) }) } + defer release() + 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), + load: func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { + callMu.Lock() + calls++ + call := calls + callMu.Unlock() + if call == 1 { + close(firstStarted) + <-releaseFirst + return []scheduledWorkflow{oldEntry}, scheduleTestSnapshot(), nil + } + close(secondLoaded) + return []scheduledWorkflow{newEntry}, scheduleTestSnapshot(), nil + }, + } + + firstDone := make(chan error, 1) + go func() { firstDone <- scheduler.RefreshRepo(ctx, repo) }() + <-firstStarted + secondAttempted := make(chan struct{}) + secondDone := make(chan error, 1) + go func() { + close(secondAttempted) + secondDone <- scheduler.RefreshRepo(ctx, repo) + }() + <-secondAttempted + select { + case <-secondLoaded: + t.Error("second repository load started before the first refresh completed") + case <-time.After(time.Second): + } + release() + if err := <-firstDone; err != nil { + t.Fatalf("first RefreshRepo: %v", err) + } + if err := <-secondDone; err != nil { + t.Fatalf("second RefreshRepo: %v", err) + } + + definitions, err := d.WorkflowSchedules(ctx) + if err != nil { + t.Fatalf("WorkflowSchedules: %v", err) + } + if len(definitions) != 1 || definitions[0].Workflow != "new.yml" { + t.Fatalf("persisted schedules = %#v, want only newest refresh", definitions) + } +} + +func TestPipelineSchedulerRemovalWaitsForInflightRefresh(t *testing.T) { + ctx := context.Background() + d := newScheduleTestDB(t) + repo := scheduleTestRepo() + if err := d.AddRepo(repo); err != nil { + t.Fatalf("AddRepo: %v", err) + } + entry := mustCron(t, "0 9 * * *") + entry.name = "build.yml" + loadStarted := make(chan struct{}) + releaseLoad := make(chan struct{}) + var releaseOnce sync.Once + 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), + load: func(context.Context, db.Repo) ([]scheduledWorkflow, db.ScheduledRepo, error) { + close(loadStarted) + <-releaseLoad + return []scheduledWorkflow{entry}, scheduleTestSnapshot(), nil + }, + } + + refreshDone := make(chan error, 1) + go func() { refreshDone <- scheduler.RefreshRepo(ctx, repo) }() + <-loadStarted + removeAttempted := make(chan struct{}) + removeDone := make(chan error, 1) + go func() { + close(removeAttempted) + removeDone <- scheduler.RemoveRepo(ctx, repo.RepoDid.String()) + }() + <-removeAttempted + select { + case <-removeDone: + t.Error("repository removal completed while its schedule refresh was still in flight") + case <-time.After(time.Second): + } + release() + if err := <-refreshDone; err != nil { + t.Fatalf("RefreshRepo: %v", err) + } + if err := <-removeDone; err != nil { + t.Fatalf("RemoveRepo: %v", err) + } + definitions, err := d.WorkflowSchedules(ctx) + if err != nil { + t.Fatalf("WorkflowSchedules: %v", err) + } + if len(definitions) != 0 { + t.Fatalf("persisted schedules after removal = %#v, want none", definitions) + } +} diff --git a/spindle/server.go b/spindle/server.go index 60ec0180e..f71f2002b 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -19,6 +19,7 @@ import ( "github.com/bluesky-social/indigo/atproto/syntax" indigoxrpc "github.com/bluesky-social/indigo/xrpc" "github.com/go-chi/chi/v5" + "github.com/go-git/go-git/v5/plumbing" "github.com/go-git/go-git/v5/plumbing/object" "github.com/hashicorp/go-version" "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp" @@ -89,6 +90,7 @@ type Spindle struct { engs map[string]models.Engine jobWake chan struct{} jobWorkers sync.WaitGroup + scheduler *pipelineScheduler cfg *config.Config feed *feed.Feed res *idresolver.Resolver @@ -287,6 +289,7 @@ func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]m spindle.res = idresolver.DefaultResolver(cfg.Server.PlcUrl) spindle.verify = repoverify.New(spindle.res, cfg.Server.Dev) + spindle.scheduler = newPipelineScheduler(spindle) err = e.AddSpindle(rbacDomain) if err != nil { @@ -382,6 +385,9 @@ func (s *Spindle) Start(ctx context.Context) error { if err != nil { return fmt.Errorf("starting metrics listener: %w", err) } + if s.scheduler != nil { + s.scheduler.Start(runCtx) + } // an executor dials out to its mill and takes work from it var execDone chan struct{} diff --git a/spindle/server_test.go b/spindle/server_test.go index 5e4be236d..6acbb4458 100644 --- a/spindle/server_test.go +++ b/spindle/server_test.go @@ -86,7 +86,7 @@ func TestExecutorRoleBuildsMinimalSpindle(t *testing.T) { t.Fatalf("New() error = %v", err) } - if s.jc != nil || s.tap != nil || s.e != nil || s.feed != nil || s.res != nil || s.vault != nil { + if s.jc != nil || s.tap != nil || s.e != nil || s.feed != nil || s.res != nil || s.vault != nil || s.scheduler != nil { t.Fatal("executor role built coordinator-only spindle dependencies") } diff --git a/spindle/tapclient.go b/spindle/tapclient.go index c28ca0adf..770b49376 100644 --- a/spindle/tapclient.go +++ b/spindle/tapclient.go @@ -133,7 +133,7 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error if record.Spindle == nil || *record.Spindle != hostname { if knownRepo { l.Info("tearing down repo reassigned from this spindle", "newSpindle", record.Spindle) - return t.teardownRepo(l, prior, ownerDid, rkey) + return t.teardownRepo(ctx, l, prior, ownerDid, rkey) } return nil } @@ -260,6 +260,11 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error t.spindle.jc.AddDid(ownerDid.String()) t.drainPendingCollabs(ctx, repoDid) + if t.spindle.scheduler != nil { + if err := t.spindle.scheduler.RefreshRepo(ctx, repo); err != nil { + l.Warn("failed to refresh repository schedules", "err", err) + } + } case tapc.RecordDeleteAction: repo, err := t.spindle.db.GetRepoByOwnerRkey(ownerDid, rkey) @@ -267,12 +272,12 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error l.Info("skipping delete for unknown repo") return nil } - return t.teardownRepo(l, repo, ownerDid, rkey) + return t.teardownRepo(ctx, l, repo, ownerDid, rkey) } return nil } -func (t *Tap) teardownRepo(l *slog.Logger, repo *db.Repo, ownerDid syntax.DID, rkey syntax.RecordKey) error { +func (t *Tap) teardownRepo(ctx context.Context, l *slog.Logger, repo *db.Repo, ownerDid syntax.DID, rkey syntax.RecordKey) error { if repo.RepoDid != "" { if err := t.spindle.WipeRepo(context.Background(), repo.RepoDid, "repo record removed"); err != nil { return err diff --git a/spindle/wipe.go b/spindle/wipe.go index 8c25766c3..009494527 100644 --- a/spindle/wipe.go +++ b/spindle/wipe.go @@ -63,6 +63,9 @@ func (s *Spindle) WipeRepo(ctx context.Context, repoDid syntax.DID, reason strin fail("delete mill leases", err) fail("delete queued jobs", s.db.DeleteJobsByRepo(ctx, repoDid.String())) + if s.scheduler != nil { + fail("delete schedules", s.scheduler.RemoveRepo(ctx, repoDid.String())) + } fail("delete quota state", s.db.DeleteQuotaStateForRepo(ctx, repoDid.String())) fail("delete secrets", s.vault.RemoveAllSecrets(ctx, secrets.RepoIdentifier(repoDid.String()))) fail("delete webhooks", s.db.DeleteWebhooksByRepo(repoDid))