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), minimumInterval: 5 * time.Minute, 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), minimumInterval: 5 * time.Minute, 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 TestPipelineSchedulerDispatchTimeoutReleasesClaim(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)), byRepo: make(map[string][]*scheduledWorkflow), minimumInterval: 5 * time.Minute, concurrency: 1, dispatchTimeout: 25 * time.Millisecond, dispatch: func(ctx context.Context, _ db.Repo, _ db.ScheduledRepo, _ []string, _ time.Time) (models.PipelineId, error) { <-ctx.Done() return models.PipelineId{}, ctx.Err() }, } scheduler.replaceCachedSchedules(entry.repo.RepoDid.String(), []scheduledWorkflow{entry}, at.Add(-time.Minute)) start := time.Now() scheduler.dispatchDue(context.Background(), at) if elapsed := time.Since(start); elapsed > 5*time.Second { t.Fatalf("dispatchDue returned after %s; want the dispatch timeout to unblock it", elapsed) } reclaimed, err := d.ClaimScheduleRun(context.Background(), entry.repo.RepoDid.String(), entry.name, at, scheduler.minimumInterval) if err != nil { t.Fatal(err) } if !reclaimed { t.Fatal("timed-out dispatch left its schedule claim behind") } } 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)), 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++ 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)), 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() 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, scheduler.minimumInterval) if err != nil { t.Fatal(err) } if !reclaimed { t.Fatal("canceled dispatch left its schedule claim") } } 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 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)), 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") }, } 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, scheduler.minimumInterval) 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 5 * * 1-5", Timezone: "America/New_York", }}, 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)), 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) } 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), minimumInterval: 5 * time.Minute, 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), minimumInterval: 5 * time.Minute, 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)), minimumInterval: 5 * time.Minute, concurrency: 1, now: time.Now, byRepo: make(map[string][]*scheduledWorkflow), load: func(_ context.Context, _ db.Repo, _ time.Time) ([]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)), minimumInterval: 5 * time.Minute, concurrency: 1, now: time.Now, byRepo: make(map[string][]*scheduledWorkflow), load: func(_ context.Context, _ db.Repo, _ time.Time) ([]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) } }