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 // 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 { repo db.Repo scheduledRepo db.ScheduledRepo name string expression string timezone 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, time.Time) ([]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 minimumInterval time.Duration concurrency int dispatchTimeout time.Duration 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 { minimumInterval := s.cfg.Schedule.MinimumInterval 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, minimumInterval: minimumInterval, concurrency: concurrency, dispatchTimeout: scheduleDispatchTimeout, 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.schedulePruneBefore(s.now())); 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.schedulePruneBefore(s.now())); err != nil { s.l.Error("failed to prune schedule claims", "err", err) } } } } 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 { retention = s.minimumInterval } return now.UTC().Add(-retention) } 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) groupedSchedules := make(map[string]map[string][]cron.Schedule) seen := make(map[struct { repo string workflow string expression string timezone string }]struct{}) for _, definition := range definitions { repo, ok := repoByDid[definition.RepoDid] if !ok { continue } if _, ok := grouped[definition.RepoDid]; !ok { grouped[definition.RepoDid] = nil } timezone := definition.Timezone if timezone == "" { timezone = "UTC" } key := struct { repo string workflow string expression string timezone string }{definition.RepoDid, definition.Workflow, definition.Expression, timezone} if _, ok := seen[key]; ok { continue } seen[key] = struct{}{} schedule, err := workflow.ParseCron( definition.Expression, timezone, definition.RepoDid+"\x00"+definition.Workflow, ) if err != nil { s.l.Warn("ignoring invalid persisted schedule", "repo", definition.RepoDid, "workflow", definition.Workflow, "expression", definition.Expression, "timezone", timezone, "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, timezone: timezone, schedule: schedule, }) if groupedSchedules[definition.RepoDid] == nil { groupedSchedules[definition.RepoDid] = make(map[string][]cron.Schedule) } groupedSchedules[definition.RepoDid][definition.Workflow] = append( groupedSchedules[definition.RepoDid][definition.Workflow], schedule, ) } now := s.now() for repoDid, entries := range grouped { invalidWorkflows := make(map[string]struct{}) for workflowName, schedules := range groupedSchedules[repoDid] { if err := workflow.ValidateScheduleInterval(schedules, s.minimumInterval, now); err != nil { s.l.Warn("ignoring invalid persisted workflow schedules", "repo", repoDid, "workflow", workflowName, "minimumInterval", s.minimumInterval, "err", err) invalidWorkflows[workflowName] = struct{}{} } } valid := make([]scheduledWorkflow, 0, len(entries)) for _, entry := range entries { if _, invalid := invalidWorkflows[entry.name]; invalid { continue } valid = append(valid, entry) } s.replaceCachedSchedules(repoDid, valid, 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, s.now()) 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, Timezone: entries[i].timezone, 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, s.minimumInterval) 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 } 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 == "" { 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, after time.Time) ([]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()) } minimumInterval := s.cfg.Schedule.MinimumInterval var entries []scheduledWorkflow seen := make(map[struct{ name, expression, timezone string }]struct{}) for _, wf := range parsed { var workflowSchedules []cron.Schedule var workflowEntries []scheduledWorkflow invalid := false for _, constraint := range wf.When { for _, definition := range constraint.Schedule { expression := definition.Cron timezone := definition.Timezone if timezone == "" { timezone = "UTC" } key := struct{ name, expression, timezone string }{wf.Name, expression, timezone} if _, ok := seen[key]; ok { continue } seen[key] = struct{}{} schedule, err := workflow.ParseCron( expression, timezone, repo.RepoDid.String()+"\x00"+wf.Name, ) if err != nil { s.l.Warn("ignoring invalid workflow schedule", "repo", repo.RepoDid, "workflow", wf.Name, "expression", expression, "timezone", timezone, "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) } workflowSchedules = append(workflowSchedules, schedule) workflowEntries = append(workflowEntries, scheduledWorkflow{ repo: repo, scheduledRepo: scheduledRepo, name: wf.Name, expression: expression, timezone: timezone, schedule: schedule, }) } } if invalid { s.l.Warn("ignoring workflow with invalid schedule", "repo", repo.RepoDid, "workflow", wf.Name) continue } if err := workflow.ValidateScheduleInterval(workflowSchedules, minimumInterval, after); err != nil { s.l.Warn("ignoring workflow with schedules that are too close", "repo", repo.RepoDid, "workflow", wf.Name, "minimumInterval", minimumInterval, "err", err) 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, ) }