package repository import ( "database/sql" "fmt" "time" "vidarchive/internal/models" ) type SubscriptionRepository struct { db *sql.DB } func NewSubscriptionRepository(db *sql.DB) *SubscriptionRepository { return &SubscriptionRepository{db: db} } const subscriptionColumns = `id, name, url, enabled, refresh_mode, schedule_kind, cron_expr, preset_id, format_override, custom_flags, output_dir, prune_removed, last_run_at, next_run_at, last_status, created_at` func (r *SubscriptionRepository) Create(s *models.Subscription) error { result, err := r.db.Exec( `INSERT INTO subscriptions (name, url, enabled, refresh_mode, schedule_kind, cron_expr, preset_id, format_override, custom_flags, output_dir, prune_removed, next_run_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, s.Name, s.URL, boolToInt(s.Enabled), s.RefreshMode, s.ScheduleKind, s.CronExpr, s.PresetID, s.FormatOverride, s.CustomFlags, s.OutputDir, boolToInt(s.PruneRemoved), s.NextRunAt, ) if err != nil { return err } id, err := result.LastInsertId() if err != nil { return fmt.Errorf("subscription last insert id: %w", err) } s.ID = id return nil } func (r *SubscriptionRepository) GetByID(id int64) (*models.Subscription, error) { row := r.db.QueryRow(`SELECT `+subscriptionColumns+` FROM subscriptions WHERE id = ?`, id) return scanSubscription(row) } func (r *SubscriptionRepository) GetAll() ([]*models.Subscription, error) { rows, err := r.db.Query(`SELECT ` + subscriptionColumns + ` FROM subscriptions ORDER BY created_at DESC`) if err != nil { return nil, err } defer rows.Close() return scanSubscriptions(rows) } // GetDue returns enabled subscriptions due at or before now, or never run. func (r *SubscriptionRepository) GetDue(now time.Time) ([]*models.Subscription, error) { rows, err := r.db.Query( `SELECT `+subscriptionColumns+` FROM subscriptions WHERE enabled = 1 AND (next_run_at IS NULL OR next_run_at <= ?) ORDER BY created_at ASC`, now, ) if err != nil { return nil, err } defer rows.Close() return scanSubscriptions(rows) } func (r *SubscriptionRepository) Update(s *models.Subscription) error { _, err := r.db.Exec( `UPDATE subscriptions SET name=?, url=?, enabled=?, refresh_mode=?, schedule_kind=?, cron_expr=?, preset_id=?, format_override=?, custom_flags=?, output_dir=?, prune_removed=?, next_run_at=? WHERE id=?`, s.Name, s.URL, boolToInt(s.Enabled), s.RefreshMode, s.ScheduleKind, s.CronExpr, s.PresetID, s.FormatOverride, s.CustomFlags, s.OutputDir, boolToInt(s.PruneRemoved), s.NextRunAt, s.ID, ) return err } func (r *SubscriptionRepository) SetEnabled(id int64, enabled bool) error { _, err := r.db.Exec(`UPDATE subscriptions SET enabled = ? WHERE id = ?`, boolToInt(enabled), id) return err } func (r *SubscriptionRepository) MarkRun(id int64, lastRunAt, nextRunAt time.Time, status string) error { _, err := r.db.Exec( `UPDATE subscriptions SET last_run_at = ?, next_run_at = ?, last_status = ? WHERE id = ?`, lastRunAt, nextRunAt, status, id, ) return err } // MarkManualRun records a run started by hand. Unlike MarkRun it leaves // next_run_at alone: running a subscription now must not move its schedule. func (r *SubscriptionRepository) MarkManualRun(id int64, lastRunAt time.Time, status string) error { _, err := r.db.Exec( `UPDATE subscriptions SET last_run_at = ?, last_status = ? WHERE id = ?`, lastRunAt, status, id, ) return err } // SetLastStatus leaves the run times alone. The worker calls it as the run // progresses, so the subscription does not stay on the status it was queued with. func (r *SubscriptionRepository) SetLastStatus(id int64, status string) error { _, err := r.db.Exec(`UPDATE subscriptions SET last_status = ? WHERE id = ?`, status, id) return err } func (r *SubscriptionRepository) Delete(id int64) error { _, err := r.db.Exec(`DELETE FROM subscriptions WHERE id = ?`, id) return err } func scanSubscriptions(rows *sql.Rows) ([]*models.Subscription, error) { var subs []*models.Subscription for rows.Next() { s, err := scanSubscription(rows) if err != nil { return nil, err } subs = append(subs, s) } return subs, rows.Err() } func scanSubscription(row interface{ Scan(...interface{}) error }) (*models.Subscription, error) { var s models.Subscription var enabled, pruneRemoved int var refreshMode, scheduleKind, cronExpr, formatOverride, customFlags, outputDir sql.NullString err := row.Scan( &s.ID, &s.Name, &s.URL, &enabled, &refreshMode, &scheduleKind, &cronExpr, &s.PresetID, &formatOverride, &customFlags, &outputDir, &pruneRemoved, &s.LastRunAt, &s.NextRunAt, &s.LastStatus, &s.CreatedAt, ) if err != nil { return nil, err } s.Enabled = enabled == 1 s.PruneRemoved = pruneRemoved == 1 s.RefreshMode = refreshMode.String s.ScheduleKind = scheduleKind.String s.CronExpr = cronExpr.String s.FormatOverride = formatOverride.String s.CustomFlags = customFlags.String s.OutputDir = outputDir.String return &s, nil }