Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636package 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) }}