package tasks import ( "context" "fmt" "time" "atlas9.dev/c/core/dbi" ) // SqliteScheduleStore implements ScheduleStore over a SQLite table. Every // schedule table shares the same coordination columns: an id primary key, plus // schedule, next_run_at, and last_run_at. type SqliteScheduleStore struct { db dbi.DBI table string } func NewSqliteScheduleStore(db dbi.DBI, table string) *SqliteScheduleStore { return &SqliteScheduleStore{db: db, table: table} } func (s *SqliteScheduleStore) Due(ctx context.Context) ([]Due, error) { rows, err := s.db.Query(ctx, "SELECT id, schedule FROM "+s.table+" WHERE next_run_at <= datetime('now')") if err != nil { return nil, fmt.Errorf("querying due schedules: %w", err) } defer rows.Close() var dues []Due for rows.Next() { var d Due if err := rows.Scan(&d.ID, &d.Schedule); err != nil { return nil, fmt.Errorf("scanning due schedule: %w", err) } dues = append(dues, d) } return dues, rows.Err() } func (s *SqliteScheduleStore) Claim(ctx context.Context, id string, next time.Time) (bool, error) { res, err := s.db.Exec(ctx, "UPDATE "+s.table+" SET last_run_at = datetime('now'), next_run_at = $1"+ " WHERE id = $2 AND next_run_at <= datetime('now')", next.UTC().Format(sqliteTimeFormat), id) if err != nil { return false, fmt.Errorf("claiming schedule: %w", err) } n, err := res.RowsAffected() if err != nil { return false, fmt.Errorf("checking rows affected: %w", err) } return n > 0, nil }