From cf97fff650a4bafa8ac03b075f0811049b0de452 Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 18 Sep 2026 13:55:31 +0300 Subject: [PATCH] spindle/schedule: bound pipeline dispatches with a timeout Signed-off-by: dawn --- spindle/schedule.go | 16 +++++++++++++++- spindle/schedule_test.go | 35 +++++++++++++++++++++++++++++++++++ 2 files changed, 50 insertions(+), 1 deletion(-) diff --git a/spindle/schedule.go b/spindle/schedule.go index 033c1e50f..c40e966bd 100644 --- a/spindle/schedule.go +++ b/spindle/schedule.go @@ -22,6 +22,9 @@ import ( const ( scheduleRunPruneInterval = 6 * time.Hour scheduleDispatchBacklog = 60 + // sized so only a genuinely stalled fetch can hit it; a hung dispatch + // would otherwise wedge the single-minute dispatch loop + scheduleDispatchTimeout = 10 * time.Minute ) type scheduledWorkflow struct { @@ -71,6 +74,7 @@ type pipelineScheduler struct { minimumInterval time.Duration concurrency int + dispatchTimeout time.Duration mu sync.Mutex byRepo map[string][]*scheduledWorkflow @@ -95,6 +99,7 @@ func newPipelineScheduler(s *Spindle) *pipelineScheduler { now: time.Now, minimumInterval: minimumInterval, concurrency: concurrency, + dispatchTimeout: scheduleDispatchTimeout, byRepo: make(map[string][]*scheduledWorkflow), } } @@ -140,6 +145,13 @@ func (s *pipelineScheduler) run(ctx context.Context) { } } +func (s *pipelineScheduler) effectiveDispatchTimeout() time.Duration { + if s.dispatchTimeout > 0 { + return s.dispatchTimeout + } + return scheduleDispatchTimeout +} + func (s *pipelineScheduler) schedulePruneBefore(now time.Time) time.Time { retention := 24 * time.Hour if s.minimumInterval > retention { @@ -538,7 +550,9 @@ func (s *pipelineScheduler) dispatchRepository(ctx context.Context, group dueRep return } - pipelineID, err := s.dispatch(ctx, group.repo, group.scheduledRepo, claimed, at) + dispatchCtx, cancel := context.WithTimeout(ctx, s.effectiveDispatchTimeout()) + defer cancel() + pipelineID, err := s.dispatch(dispatchCtx, 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 == "" { diff --git a/spindle/schedule_test.go b/spindle/schedule_test.go index 97cc984a7..8e21b47b5 100644 --- a/spindle/schedule_test.go +++ b/spindle/schedule_test.go @@ -129,6 +129,41 @@ func TestPipelineSchedulerDispatchDueClaimsOccurrenceOnce(t *testing.T) { 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() -- 2.51.2