package scheduling import ( "context" "encoding/json" "errors" "fmt" "log/slog" "time" "crussell/auth" "crussell/clock" "crussell/db" "github.com/jackc/pgx/v5" ) // NotifyUnpaidOneWeek inserts admin_notifications for bookings that ended // 7+ days ago with no completed payment. Runs daily. // TODO: Notify affected user via email/SMS when SMTP is wired (E5). func NotifyUnpaidOneWeek(ctx context.Context) (int, error) { tx, err := db.Conn.Begin(ctx) if err != nil { return 0, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if err := tx.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { slog.Error("failed to rollback transaction", "err", err) } }() rows, err := tx.Query(ctx, ` SELECT b.id, b.user_id FROM bookings b WHERE b.status NOT IN ('client_cancelled', 'we_cancelled') AND NOT EXISTS ( SELECT 1 FROM payments p WHERE p.booking_id = b.id AND p.status = 'completed' ) AND b.end_time >= NOW() - INTERVAL '30 days' AND b.end_time < NOW() - INTERVAL '7 days' AND NOT EXISTS ( SELECT 1 FROM admin_notifications an WHERE an.booking_id = b.id AND an.reason = '1_week_no_pay' ) `) if err != nil { return 0, fmt.Errorf("failed to query unpaid 1-week bookings: %w", err) } defer rows.Close() var ids, userIDs []string for rows.Next() { var id, userID string if err := rows.Scan(&id, &userID); err != nil { return 0, fmt.Errorf("failed to scan row: %w", err) } ids = append(ids, id) userIDs = append(userIDs, userID) } if err := rows.Err(); err != nil { return 0, fmt.Errorf("rows iteration error: %w", err) } if len(ids) == 0 { if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return 0, nil } _, err = tx.Exec(ctx, ` INSERT INTO admin_notifications (reason, booking_id, user_id) SELECT '1_week_no_pay', unnest($1::text[]), unnest($2::text[]) `, ids, userIDs) if err != nil { return 0, fmt.Errorf("failed to insert 1_week_no_pay notifications: %w", err) } if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return len(ids), nil } // NotifyUnpaidOneMonth inserts admin_notifications for bookings that ended // 30+ days ago with no completed payment. Runs daily. // TODO: Notify affected user via email/SMS when SMTP is wired (E5). func NotifyUnpaidOneMonth(ctx context.Context) (int, error) { tx, err := db.Conn.Begin(ctx) if err != nil { return 0, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if err := tx.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { slog.Error("failed to rollback transaction", "err", err) } }() rows, err := tx.Query(ctx, ` SELECT b.id, b.user_id FROM bookings b WHERE b.status NOT IN ('client_cancelled', 'we_cancelled') AND NOT EXISTS ( SELECT 1 FROM payments p WHERE p.booking_id = b.id AND p.status = 'completed' ) AND b.end_time < NOW() - INTERVAL '30 days' AND NOT EXISTS ( SELECT 1 FROM admin_notifications an WHERE an.booking_id = b.id AND an.reason = '1_month_no_pay' ) `) if err != nil { return 0, fmt.Errorf("failed to query unpaid 1-month bookings: %w", err) } defer rows.Close() var ids, userIDs []string for rows.Next() { var id, userID string if err := rows.Scan(&id, &userID); err != nil { return 0, fmt.Errorf("failed to scan row: %w", err) } ids = append(ids, id) userIDs = append(userIDs, userID) } if err := rows.Err(); err != nil { return 0, fmt.Errorf("rows iteration error: %w", err) } if len(ids) == 0 { if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return 0, nil } _, err = tx.Exec(ctx, ` INSERT INTO admin_notifications (reason, booking_id, user_id) SELECT '1_month_no_pay', unnest($1::text[]), unnest($2::text[]) `, ids, userIDs) if err != nil { return 0, fmt.Errorf("failed to insert 1_month_no_pay notifications: %w", err) } if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return len(ids), nil } // TransitionDiscountCampaigns auto-transitions campaign statuses based on // dates and redemption limits: // - draft → active when start_date <= NOW() // - active → completed when end_date < NOW() or max_redemptions reached func TransitionDiscountCampaigns(ctx context.Context) (int, error) { tx, err := db.Conn.Begin(ctx) if err != nil { return 0, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if err := tx.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { slog.Error("failed to rollback transaction", "err", err) } }() result, err := tx.Exec(ctx, ` UPDATE discount_campaigns SET status = 'active' WHERE status = 'draft' AND campaign_type = 'time_based' AND start_date <= NOW() `) if err != nil { return 0, fmt.Errorf("failed to activate campaigns: %w", err) } activated := result.RowsAffected() result, err = tx.Exec(ctx, ` UPDATE discount_campaigns SET status = 'completed' WHERE status = 'active' AND ( end_date < NOW() OR (max_redemptions IS NOT NULL AND times_redeemed >= max_redemptions) ) `) if err != nil { return 0, fmt.Errorf("failed to complete campaigns: %w", err) } completed := result.RowsAffected() if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return int(activated + completed), nil } // CleanupExpiredVerificationCodes deletes expired verification codes and // used codes older than 30 days. func CleanupExpiredVerificationCodes(ctx context.Context) (int, error) { tx, err := db.Conn.Begin(ctx) if err != nil { return 0, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if err := tx.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { slog.Error("failed to rollback transaction", "err", err) } }() result, err := tx.Exec(ctx, ` DELETE FROM verification_codes WHERE (used_at IS NOT NULL AND used_at < NOW() - INTERVAL '30 days') OR (expires_at < NOW() AND used_at IS NULL) `) if err != nil { return 0, fmt.Errorf("failed to cleanup verification codes: %w", err) } if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return int(result.RowsAffected()), nil } // CleanupExpiredRefreshTokens deletes expired refresh tokens and // revoked tokens older than the shared RefreshTokenLifetime window (auth). func CleanupExpiredRefreshTokens(ctx context.Context) (int, error) { tx, err := db.Conn.Begin(ctx) if err != nil { return 0, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if err := tx.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { slog.Error("failed to rollback transaction", "err", err) } }() result, err := tx.Exec(ctx, ` DELETE FROM refresh_tokens WHERE expires_at < NOW() OR (revoked = TRUE AND created_at < NOW() - make_interval(days => $1)) `, int64(auth.RefreshTokenLifetime/(24*time.Hour))) if err != nil { return 0, fmt.Errorf("failed to cleanup refresh tokens: %w", err) } if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } return int(result.RowsAffected()), nil } // LondonDateString returns the Europe/London calendar date (YYYY-MM-DD) for // the given instant. ApplyScheduledDefaultHours uses this so staged // default-hours changes roll over at London midnight rather than UTC midnight: // during BST a London date begins at 23:00 UTC the previous day, so using the // UTC date would misdate the effective day inside the 00:00-01:00 BST window. func LondonDateString(t time.Time) string { return t.In(clock.London).Format("2006-01-02") } // ApplyScheduledDefaultHours applies any pending default hours changes that // have reached their effective_date. Runs daily at 00:05 to catch midnight // roll-overs even if the cron is slightly delayed. func ApplyScheduledDefaultHours(ctx context.Context) (int, error) { todayStr := LondonDateString(clock.Now()) tx, err := db.Conn.Begin(ctx) if err != nil { return 0, fmt.Errorf("failed to begin transaction: %w", err) } defer func() { if err := tx.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { slog.Error("failed to rollback transaction", "err", err) } }() // Find any changes where effective_date <= today (London) and not yet applied/cancelled var rowID int var hoursJSON string err = tx.QueryRow(ctx, ` SELECT id, hours::text FROM default_hours_scheduled_changes WHERE effective_date <= $1::date AND applied_at IS NULL AND cancelled_at IS NULL LIMIT 1 `, todayStr).Scan(&rowID, &hoursJSON) if err != nil { // No matching change — nothing to do return 0, tx.Commit(ctx) } // Parse the staged hours var hours []struct { Weekday int `json:"weekday"` StartTime string `json:"startTime"` EndTime string `json:"endTime"` IsOpen bool `json:"isOpen"` } if err := json.Unmarshal([]byte(hoursJSON), &hours); err != nil { return 0, fmt.Errorf("failed to parse hours JSON for change %d: %w", rowID, err) } // Bulk-update working_hours with staged values weekdays := make([]int, len(hours)) startTimes := make([]string, len(hours)) endTimes := make([]string, len(hours)) isOpenFlags := make([]bool, len(hours)) for i, h := range hours { weekdays[i] = h.Weekday startTimes[i] = h.StartTime endTimes[i] = h.EndTime isOpenFlags[i] = h.IsOpen } if _, err := tx.Exec(ctx, ` UPDATE working_hours AS wh SET start_time = v.start_time, end_time = v.end_time, is_open = v.is_open FROM ( SELECT unnest($1::smallint[]) AS weekday, unnest($2::time[]) AS start_time, unnest($3::time[]) AS end_time, unnest($4::boolean[]) AS is_open ) v WHERE wh.weekday = v.weekday `, weekdays, startTimes, endTimes, isOpenFlags); err != nil { return 0, fmt.Errorf("failed to update working_hours for change %d: %w", rowID, err) } // Mark change as applied if _, err := tx.Exec(ctx, ` UPDATE default_hours_scheduled_changes SET applied_at = NOW() WHERE id = $1 `, rowID); err != nil { return 0, fmt.Errorf("failed to mark change %d as applied: %w", rowID, err) } // Insert admin notification if _, err := tx.Exec(ctx, ` INSERT INTO admin_notifications (reason, created_at) VALUES ('default_hours_changed', NOW()) `); err != nil { return 0, fmt.Errorf("failed to insert notification for change %d: %w", rowID, err) } if err := tx.Commit(ctx); err != nil { return 0, fmt.Errorf("failed to commit: %w", err) } slog.Info("Applied scheduled default hours change", "change_id", rowID, "effective_date", todayStr) return 1, nil }