package tasks import ( "context" "database/sql" "log/slog" "time" "atlas9.dev/c/core/dbi" ) // Runner fires due schedules from one table. It lists due rows, then for each // atomically claims it (a conditional UPDATE) and enqueues its task in one // transaction, so multiple instances coordinate without a leader and never // double-enqueue a single occurrence. // // The schedule table is authoritative: its rows, their schedule strings, and // their run times all live in the DB — seeded by migration for system // schedules, written by users for dynamic ones. The Runner holds no schedule // definitions of its own; it only reads them through Store and enqueues each // due row's task through Dispatch. One Runner drives one table, the way one // Worker drains one queue. type Runner struct { DB *sql.DB // Store opens the schedule store on a transaction. Store func(dbi.DBI) ScheduleStore // Enqueue writes the task record for a due schedule, identified by its id. // It runs inside the claim transaction so the claim and the enqueue commit // atomically. Do nothing here but enqueue — it holds an open write // transaction, so any real work belongs in the task, not here. For system // schedules the id is the schedule type, and Enqueue switches on it. Enqueue func(ctx context.Context, tx dbi.DBI, id string) error // Interval is how long the Runner sleeps between ticks. Schedules are // minute-granular (ParseSchedule floors at 1m), so polling any faster buys // nothing; a fixed tick is enough and keeps this in step with Worker. // Defaults to 1m. Interval time.Duration } func (r *Runner) Run(ctx context.Context) { for { if ctx.Err() != nil { return } r.tick(ctx) select { case <-time.After(r.interval()): case <-ctx.Done(): return } } } func (r *Runner) interval() time.Duration { if r.Interval > 0 { return r.Interval } return time.Minute } // tick fires every schedule whose next run has passed. Each is fired // independently so one bad schedule doesn't wedge the others. func (r *Runner) tick(ctx context.Context) { var dues []Due err := dbi.ReadOnly(ctx, r.DB, func(tx dbi.DBI) error { var err error dues, err = r.Store(tx).Due(ctx) return err }) if err != nil { slog.ErrorContext(ctx, "listing due schedules", "err", err) return } for _, d := range dues { r.fire(ctx, d) } } // fire claims one schedule and, if it wins the claim, enqueues its task in the // same transaction. The claim re-checks next_run_at against DB time, so only // one caller across the fleet can win a given occurrence. func (r *Runner) fire(ctx context.Context, d Due) { schedule, err := ParseSchedule(d.Schedule) if err != nil { // Bad schedule, e.g. an operator typo in the DB. Skip so others run. slog.ErrorContext(ctx, "invalid schedule", "schedule", d.ID, "value", d.Schedule, "err", err) return } // Claim and enqueue in one transaction: the conditional UPDATE is the // fleet-wide coordination point (exactly one caller wins), and enqueuing in // the same tx makes advancing next_run_at and writing the task atomic. If // the enqueue fails the claim rolls back, so the next tick retries. next := schedule.Next(time.Now()) err = dbi.ReadWrite(ctx, r.DB, func(tx dbi.DBI) error { won, err := r.Store(tx).Claim(ctx, d.ID, next) if err != nil { return err } if !won { // Lost the claim to another tick or instance. return nil } return r.Enqueue(ctx, tx, d.ID) }) if err != nil { slog.ErrorContext(ctx, "firing schedule", "schedule", d.ID, "err", err) } }