Files
Crussell/backend/internal/jobs/cleanup.go
T
popertotsandSisyphus 62dca184df fix: account security + GDPR erasure — current-password lockout budgets (atomic 15/30/60 escalation, uniform 401), DAV + S3 deletion durable in-tx, batch erasure outbox
- ChangePassword/DeleteAccount: failed-attempt lockout matching login escalation, atomic check-increment (no burst), uniform 401 with distinct bodies, passwordless accounts require 2FA unconditionally to delete, NULL-password change-password clear error
- erasure: CardDAV dav_cards rows deleted inside the erasure transaction (was fire-and-forget goroutine); S3 profile-pic deletion via pending_s3_deletions outbox + retry job; stale-guest/idle-account batch paths write the outbox in-tx and skip the guessed-bucket fallback
- S3_PROFILE_PICS_BUCKET unset -> fail-closed warning (once per process)
- scheduler test: 27 jobs (retry-s3-deletions)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)
2026-08-22 00:34:51 +01:00

669 lines
24 KiB
Go

package jobs
import (
"context"
"database/sql"
"errors"
"fmt"
"log"
"log/slog"
"time"
"crussell/auth"
"crussell/db"
authHandlers "crussell/handlers/auth"
"crussell/handlers/payments"
"crussell/handlers/scheduling"
"crussell/handlers/user"
"crussell/internal/adminnotify"
"crussell/internal/s3"
"crussell/internal/square"
"crussell/mw"
"github.com/jackc/pgx/v5"
)
// RegisterAll registers every background maintenance job on the scheduler.
// Call once during server startup, before s.Start().
func RegisterAll(s *Scheduler) {
// === HIGH FREQUENCY — every 5 minutes ===
s.Register(Job{
Name: "cleanup-reservations",
Schedule: "*/5 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupOldReservations,
})
s.Register(Job{
Name: "cleanup-expired-deposits",
Schedule: "*/5 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupExpiredDeposits,
})
s.Register(Job{
Name: "cleanup-rate-limiters",
Schedule: "*/5 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: mw.CleanupAllRateLimiters,
})
s.Register(Job{
Name: "cleanup-gdpr-export-cache",
Schedule: "*/5 * * * *",
Timeout: 10 * time.Second,
Concurrency: 1,
Handler: user.CleanupGDPRExportCache,
})
s.Register(Job{
Name: "sweep-pending-square-refunds",
Schedule: "*/5 * * * *",
Timeout: 60 * time.Second,
Concurrency: 1,
Handler: payments.SweepPendingSquareRefunds,
})
// Offset from the refund sweep (which also writes payments rows) by one
// minute to avoid the two sweeps contending on the same table.
s.Register(Job{
Name: "sweep-stale-pending-payments",
Schedule: "1,6,11,16,21,26,31,36,41,46,51,56 * * * *",
Timeout: 60 * time.Second,
Concurrency: 1,
Handler: payments.SweepStalePendingPayments,
})
// Cancels terminal (card-machine) checkouts still pending at Square after
// an hour — a never-polled checkout would otherwise sit live indefinitely
// and complete into an invisible, untracked charge.
s.Register(Job{
Name: "sweep-stale-terminal-checkouts",
Schedule: "*/15 * * * *",
Timeout: 60 * time.Second,
Concurrency: 1,
Handler: payments.SweepStaleTerminalCheckouts,
})
// === MID FREQUENCY — every minute (progressive rate limiter was on 30s) ===
s.Register(Job{
Name: "cleanup-progressive-rate-limiter",
Schedule: "* * * * *",
Timeout: 10 * time.Second,
Concurrency: 1,
Handler: mw.CleanupProgressiveRateLimiter,
})
// === MID FREQUENCY — hourly ===
s.Register(Job{
Name: "cleanup-expired-loyalty-redemptions",
Schedule: "0 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupExpiredLoyaltyRedemptions,
})
s.Register(Job{
Name: "cleanup-old-idempotency-keys",
Schedule: "0 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupOldIdempotencyKeys,
})
s.Register(Job{
Name: "cleanup-revoked-jtis",
Schedule: "0 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: auth.CleanupRevokedJTIs,
})
s.Register(Job{
Name: "cleanup-stale-login-entries",
Schedule: "0 * * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: authHandlers.CleanupStaleLoginEntries,
})
// === LOW FREQUENCY — daily, off-peak (staggered to avoid DB contention) ===
s.Register(Job{
Name: "anonymize-stale-guest-accounts",
Schedule: "0 3 * * *",
Timeout: 5 * time.Minute,
Concurrency: 1,
Handler: scheduling.AnonymizeStaleGuestAccounts,
})
s.Register(Job{
Name: "cleanup-expired-financial-records",
Schedule: "0 4 * * *",
Timeout: 10 * time.Minute,
Concurrency: 1,
Handler: scheduling.CleanupExpiredFinancialRecords,
})
s.Register(Job{
Name: "cleanup-idle-accounts",
Schedule: "30 3 * * *",
Timeout: 5 * time.Minute,
Concurrency: 1,
Handler: scheduling.CleanupIdleAccounts,
})
s.Register(Job{
Name: "cleanup-expired-gift-cards",
Schedule: "0 5 * * *",
Timeout: 5 * time.Minute,
Concurrency: 1,
Handler: scheduling.CleanupExpiredGiftCards,
})
s.Register(Job{
Name: "cleanup-old-name-history",
Schedule: "30 4 * * *",
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupOldNameHistory,
})
// === STAGED HOURS CHANGE ===
s.Register(Job{
Name: "apply-default-hours",
Schedule: "5 0 * * *", // Daily at 00:05 — after midnight to avoid race
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.ApplyScheduledDefaultHours,
})
// === BUSINESS LOGIC JOBS ===
s.Register(Job{
Name: "notify-unpaid-1-week",
Schedule: "0 7 * * *", // Daily at 7am — end of business day + 7 days
Timeout: 2 * time.Minute,
Concurrency: 1,
Handler: scheduling.NotifyUnpaidOneWeek,
})
s.Register(Job{
Name: "notify-unpaid-1-month",
Schedule: "30 7 * * *", // Daily at 7:30am (staggered from notify-unpaid-1-week)
Timeout: 2 * time.Minute,
Concurrency: 1,
Handler: scheduling.NotifyUnpaidOneMonth,
})
s.Register(Job{
Name: "transition-discount-campaigns",
Schedule: "0 * * * *", // Hourly
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.TransitionDiscountCampaigns,
})
s.Register(Job{
Name: "cleanup-verification-codes",
Schedule: "0 2 * * *", // Daily at 2am
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupExpiredVerificationCodes,
})
s.Register(Job{
Name: "cleanup-refresh-tokens",
Schedule: "0 2 * * *", // Daily at 2am
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: scheduling.CleanupExpiredRefreshTokens,
})
// Staggered from cleanup-verification-codes / cleanup-refresh-tokens
// (both at 0 2 * * *) to avoid DB contention.
s.Register(Job{
Name: "sweep-square-webhook-events",
Schedule: "30 2 * * *", // Daily at 2:30am
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: SweepSquareWebhookEvents,
})
// Surfaces unresolved money events (stale pending payments/till sales,
// refunds at the retry cap) as admin notifications so the owner sees the
// "money may have moved at Square but the DB couldn't record it" situations
// that would otherwise live only in un-watched CRITICAL log lines (Gap
// Backlog T14). Retire this job when a proper log/alert pipeline (T14
// Sentry) lands — the notifications page is the stopgap while the app has
// none. Staggered at 2:45am between the 2am and 3am daily batches.
s.Register(Job{
Name: "scan-critical-payment-logs",
Schedule: "45 2 * * *", // Daily at 2:45am
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: ScanCriticalPaymentLogs,
})
// Durable safety net for the GDPR account-deletion Square outbox (Fault
// A1): retries the Square card/customer deletions that the
// DeleteAccountHandler async cleanup could not finish (process crash or
// exhausted retries), using the outbox rows the handler persisted inside
// its anonymization tx before it committed. Hourly — the async goroutine
// handles the common case within seconds, so this only catches stragglers.
// The schedule is deliberately unshared so a long run cannot contend with
// the payment sweeps.
s.Register(Job{
Name: "retry-square-erasures",
Schedule: "17 * * * *", // Hourly at :17
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: RetryPendingSquareErasures,
})
// Durable safety net for the account-deletion S3/R2 profile-picture outbox
// (Fault A2): retries the object deletions that DeleteAccountHandler's
// async cleanup goroutine could not finish (process crash or an S3 outage),
// using the pending_s3_deletions rows the handler persisted inside its
// anonymization tx before it committed. Hourly, offset from
// retry-square-erasures (:37 vs :17) so the two erasure jobs cannot
// contend.
s.Register(Job{
Name: "retry-s3-deletions",
Schedule: "37 * * * *", // Hourly at :37
Timeout: 30 * time.Second,
Concurrency: 1,
Handler: RetryPendingS3Deletions,
})
}
// SweepSquareWebhookEvents deletes square_webhook_events rows older than the
// shared auth.RefreshTokenLifetime window (90 days). Every accepted webhook
// event_id is stored permanently for restart-safe dedup (payment.updated fires
// on every payment update), so without retention the table would grow without
// bound. 90 days comfortably exceeds Square's webhook replay window while
// keeping the table bounded. The 90-day window deliberately reuses
// auth.RefreshTokenLifetime so the two coinciding retention windows cannot
// drift apart.
func SweepSquareWebhookEvents(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)
}
}()
tag, err := tx.Exec(ctx, `
DELETE FROM square_webhook_events
WHERE received_at < NOW() - make_interval(days => $1)
`, int64(auth.RefreshTokenLifetime/(24*time.Hour)))
if err != nil {
return 0, fmt.Errorf("failed to sweep square webhook events: %w", err)
}
return int(tag.RowsAffected()), tx.Commit(ctx)
}
// criticalPaymentStaleAge is how long a pending payment / till sale may sit
// unresolved before the scan surfaces it to the admin notification centre. A
// pending row with no update in this window is exactly the "money may have
// moved at Square but the DB couldn't record it" situation the CRITICAL payment
// logs describe (handlers.go, giftcards.go, refunds.go, till.go). Deliberately
// shorter than the 24h stale-pending sweep (handlers/payments/sweep.go) so the
// owner hears about it while the row could still be rescued.
const criticalPaymentStaleAge = "2 hours"
// ScanCriticalPaymentLogs surfaces unresolved critical payment states as admin
// notifications (reason='critical_payment_log'), giving the owner an in-app
// view of money events that would otherwise be visible only in un-watched
// CRITICAL log lines (the app has no log/alert pipeline — Gap Backlog T14).
//
// Candidates:
// - payments / till_sales rows still 'pending' with no update for >2h: a
// lost-response charge the DB never recorded.
// - refunds rows still 'pending' at the 3-attempt retry cap: a refund the
// sweep could not resolve (money may have moved at Square).
//
// Each candidate inserts one admin_notifications row keyed on (reason,
// booking_id) — but ONLY if no unacknowledged notification for that key
// already exists (NULL-safe via IS NOT DISTINCT FROM, matching the
// insertRefundFailedNotifications dedup in handlers/payments/refunds.go), so
// the bell never floods. Acknowledging the notification re-arms the scan:
// while the row stays unresolved it is surfaced again on the next run.
//
// Flood cap (Round 2 Loop B finding 1): the unacknowledged
// 'critical_payment_log' queue is globally capped at
// adminnotify.MaxUnacknowledgedCriticalLogs — the SAME shared cap as every
// other insert site (webhooks, account erasure, time blockers). The cap is
// folded into the INSERT's WHERE clause (count-then-insert is atomic, closing
// the TOCTOU where two concurrent scans could both read a below-cap count and
// overshoot together), and the pre-check logs the suppression so a capped-out
// scan stays visible to the operator. A hostile flood of unresolved money rows
// (attacker-registered accounts, a runaway reconciliation loop) must not be
// able to bury the single-operator notification centre.
//
// This is the "grep CRITICAL" the audit asked for, but DB-backed since the app
// has no log pipeline. When a proper log/alert pipeline (T14 Sentry) lands,
// this job can be retired.
func ScanCriticalPaymentLogs(ctx context.Context) (int, error) {
// Pre-check: at the cap, skip the scan entirely (the fold inside the INSERT
// enforces the cap atomically even if a concurrent insert races this check;
// the pre-check only decides whether to log the suppression).
if adminnotify.CriticalLogsCapExceeded(ctx, db.Conn, "critical_payment_log") {
log.Printf("[SCAN] critical_payment_log admin notification suppressed — %d unacknowledged rows at the cap; acknowledge outstanding notifications to re-arm", adminnotify.MaxUnacknowledgedCriticalLogs)
return 0, nil
}
tag, err := db.Conn.Exec(ctx, `
INSERT INTO admin_notifications (reason, booking_id, created_at)
SELECT DISTINCT 'critical_payment_log'::admin_notification_reason, src.booking_id, NOW()
FROM (
-- Lost-response payments: still pending with no update past the threshold.
SELECT id, booking_id FROM payments
WHERE status = 'pending' AND updated_at < NOW() - INTERVAL '`+criticalPaymentStaleAge+`'
UNION ALL
-- Lost-response till sales (no booking linkage — booking_id NULL).
SELECT id, NULL::CHAR(12) FROM till_sales
WHERE status = 'pending' AND updated_at < NOW() - INTERVAL '`+criticalPaymentStaleAge+`'
UNION ALL
-- Refunds stuck at the 3-attempt retry cap.
SELECT id, booking_id FROM refunds
WHERE status = 'pending' AND refund_attempts >= 3
) src
WHERE NOT EXISTS (
SELECT 1 FROM admin_notifications an
WHERE an.reason = 'critical_payment_log'
AND an.booking_id IS NOT DISTINCT FROM src.booking_id
AND an.acknowledged_at IS NULL
)
-- admin_notifications.booking_id is a hard FK to bookings(id); a
-- candidate whose booking was hard-deleted (7-year retention) would
-- otherwise fail the whole INSERT and silently kill ALL critical
-- alerts. Skip orphans: keep NULL booking_ids (till sales) and only
-- payments/refunds whose booking still exists.
AND (src.booking_id IS NULL OR EXISTS (
SELECT 1 FROM bookings b WHERE b.id = src.booking_id
))
-- Flood cap (Round 2 Loop B finding 1): fold the global
-- MaxUnacknowledgedCriticalLogs bound INTO the insert so the
-- count-then-insert is atomic — two concurrent scans cannot both read a
-- below-cap count and overshoot together.
AND (SELECT COUNT(*) FROM admin_notifications _an
WHERE _an.reason = 'critical_payment_log'
AND _an.acknowledged_at IS NULL) < $1
`, adminnotify.MaxUnacknowledgedCriticalLogs)
if err != nil {
return 0, fmt.Errorf("failed to scan critical payment logs: %w", err)
}
n := int(tag.RowsAffected())
if n > 0 {
log.Printf("[SCAN] Inserted %d admin_notification(s) for unresolved critical payment states (pending payments/till sales > %s, refunds at retry cap)", n, criticalPaymentStaleAge)
}
return n, nil
}
// erasureNotificationKey derives the notification-dedup key for a pending
// outbox row: the user id when the row still carries one (registered users
// whose deletion tx has not yet run), otherwise a stable row-scoped key (guest
// rows are unlinked by delete_guest_user and the account-deletion outbox
// NULLs user_id on the rows it writes).
func erasureNotificationKey(userID sql.NullString, rowID string) string {
if userID.Valid && userID.String != "" {
return userID.String
}
return "row:" + rowID
}
// raiseErasureNotification raises the deduped critical notification once per
// key. Repeated failures across job runs collapse to a single alert (the
// deterministic admin_notifications id in
// user.InsertSquareErasureCriticalNotification does the ON CONFLICT dedup).
func raiseErasureNotification(ctx context.Context, key string, notified map[string]bool) {
if notified[key] {
return
}
notified[key] = true
user.InsertSquareErasureCriticalNotification(ctx, key)
}
// RetryPendingSquareErasures is the durable safety net for the account-deletion
// Square outbox (Fault A1). DeleteAccountHandler persists the Square
// card/customer erasure targets on the scrubbed, soft-deleted user_saved_cards
// rows (last_4 = 'XXXX') inside its anonymization transaction, before it
// commits. If the process crashes between that commit and the async cleanup
// goroutine finishing, those rows are the only remaining record of the
// Square-side PII — this job finds them and retries the Square deletion so the
// card/customer is never permanently orphaned at Square. On success (or a
// Square NOT_FOUND — the data is already gone) it drains the outbox columns; on
// final failure it raises a critical payment notification, deduped per affected
// user/row. Returns the number of outbox rows drained.
func RetryPendingSquareErasures(ctx context.Context) (int, error) {
if payments.SquareClient == nil {
// Square not configured: no external erasure is possible, and the
// handler only writes outbox entries when a client was configured.
return 0, nil
}
client := payments.SquareClient
rows, err := db.Conn.Query(ctx, `
SELECT id, user_id, square_card_id, square_customer_id
FROM user_saved_cards
WHERE deleted_at IS NOT NULL
AND last_4 = 'XXXX'
AND (square_card_id IS NOT NULL OR square_customer_id IS NOT NULL)
ORDER BY id
`)
if err != nil {
return 0, fmt.Errorf("failed to query pending Square erasures: %w", err)
}
defer rows.Close()
type pendingErasure struct {
rowID string
userID sql.NullString
cardID sql.NullString
customerID sql.NullString
}
var pending []pendingErasure
for rows.Next() {
var p pendingErasure
if err := rows.Scan(&p.rowID, &p.userID, &p.cardID, &p.customerID); err != nil {
return 0, fmt.Errorf("failed to scan pending Square erasure: %w", err)
}
pending = append(pending, p)
}
if err := rows.Err(); err != nil {
return 0, fmt.Errorf("failed to iterate pending Square erasures: %w", err)
}
if len(pending) == 0 {
return 0, nil
}
// Group the outbox rows: each card id maps to exactly one row (per-user
// UNIQUE), each customer id may map to several rows (shared across the
// user's saved cards) and may also appear on other deleted accounts' rows.
cardRows := map[string][]pendingErasure{}
customerRows := map[string][]pendingErasure{}
for _, p := range pending {
if p.cardID.Valid && p.cardID.String != "" {
cardRows[p.cardID.String] = append(cardRows[p.cardID.String], p)
}
if p.customerID.Valid && p.customerID.String != "" {
customerRows[p.customerID.String] = append(customerRows[p.customerID.String], p)
}
}
notified := map[string]bool{}
drained := map[string]bool{}
// Cards: each ccof: token is erased once. NOT_FOUND means Square no longer
// has the card — the erasure is complete, so the outbox is drained rather
// than alerted on.
for cardID, cardPend := range cardRows {
rowID := cardPend[0].rowID
err := user.RetrySquareDeletion(ctx, func(actx context.Context) error {
return client.DeleteCardOnFile(actx, cardID)
})
if err != nil && !square.IsNotFound(err) {
raiseErasureNotification(ctx, erasureNotificationKey(cardPend[0].userID, rowID), notified)
log.Printf("Error: retry-square-erasures failed to delete Square card %s (outbox row %s): %v", square.TokenPrefix(cardID), rowID, err)
slog.Error("square card erasure retry failed after attempts", "row", rowID, "card", square.TokenPrefix(cardID), "error", err)
continue
}
if _, err := db.Conn.Exec(ctx, `
UPDATE user_saved_cards SET square_card_id = NULL
WHERE id = $1 AND deleted_at IS NOT NULL
`, rowID); err != nil {
return 0, fmt.Errorf("failed to clear card erasure outbox row %s: %w", rowID, err)
}
drained[rowID] = true
}
// Customers: one DeleteCustomer per distinct id, guarded by the
// still-referenced-by-another-account check (a shared Square customer must
// survive while any active card of another account references it).
for customerID, custPend := range customerRows {
stillReferenced := false
for _, p := range custPend {
var ref bool
if err := db.Conn.QueryRow(ctx, `
SELECT EXISTS(
SELECT 1 FROM user_saved_cards
WHERE square_customer_id = $1 AND deleted_at IS NULL
AND user_id IS DISTINCT FROM $2
)
`, customerID, p.userID).Scan(&ref); err != nil {
return 0, fmt.Errorf("failed to check Square customer %s references before deletion: %w", square.TokenPrefix(customerID), err)
}
if ref {
stillReferenced = true
break
}
}
if stillReferenced {
// Deliberately kept (shared customer) — not a pending erasure.
// Drain the outbox rows so the job stops retrying a deletion that
// must not happen; the customer is erased when the last referencing
// account is itself erased.
for _, p := range custPend {
if _, err := db.Conn.Exec(ctx, `
UPDATE user_saved_cards SET square_customer_id = NULL
WHERE id = $1 AND deleted_at IS NOT NULL
`, p.rowID); err != nil {
return 0, fmt.Errorf("failed to clear skipped customer erasure outbox row %s: %w", p.rowID, err)
}
drained[p.rowID] = true
}
continue
}
err := user.RetrySquareDeletion(ctx, func(actx context.Context) error {
return client.DeleteCustomer(actx, customerID)
})
if err != nil && !square.IsNotFound(err) {
for _, p := range custPend {
raiseErasureNotification(ctx, erasureNotificationKey(p.userID, p.rowID), notified)
}
log.Printf("Error: retry-square-erasures failed to delete Square customer %s (%d outbox rows): %v", square.TokenPrefix(customerID), len(custPend), err)
slog.Error("square customer erasure retry failed after attempts", "customer", square.TokenPrefix(customerID), "rows", len(custPend), "error", err)
continue
}
for _, p := range custPend {
if _, err := db.Conn.Exec(ctx, `
UPDATE user_saved_cards SET square_customer_id = NULL
WHERE id = $1 AND deleted_at IS NOT NULL
`, p.rowID); err != nil {
return 0, fmt.Errorf("failed to clear customer erasure outbox row %s: %w", p.rowID, err)
}
drained[p.rowID] = true
}
}
if n := len(drained); n > 0 {
log.Printf("[ERASURE] retry-square-erasures drained %d pending Square erasure outbox row(s)", n)
}
return len(drained), nil
}
// RetryPendingS3Deletions is the durable safety net for the account-deletion
// S3/R2 profile-picture outbox (Fault A2). DeleteAccountHandler persists the
// deletion target (bucket + object key) inside its anonymization transaction,
// before it commits. If the process crashes between that commit and the async
// cleanup goroutine finishing, the photo's PII would otherwise be retained in
// the object store indefinitely — this job finds the pending rows and retries
// the deletion so the erasure is eventually complete. On success it deletes
// the outbox row; on failure it leaves the row in place (bumping attempts and
// recording the error) so the next run retries again. Returns the number of
// outbox rows drained.
func RetryPendingS3Deletions(ctx context.Context) (int, error) {
if s3.Client == nil {
// No object store configured: no deletion is possible, and the handler
// only writes outbox entries when a client was configured.
return 0, nil
}
client := s3.Client
rows, err := db.Conn.Query(ctx, `
SELECT id, bucket, object_key
FROM pending_s3_deletions
ORDER BY created_at
`)
if err != nil {
return 0, fmt.Errorf("failed to query pending S3 deletions: %w", err)
}
defer rows.Close()
type pendingDelete struct {
id string
bucket string
key string
}
var pending []pendingDelete
for rows.Next() {
var p pendingDelete
if err := rows.Scan(&p.id, &p.bucket, &p.key); err != nil {
return 0, fmt.Errorf("failed to scan pending S3 deletion: %w", err)
}
pending = append(pending, p)
}
if err := rows.Err(); err != nil {
return 0, fmt.Errorf("failed to iterate pending S3 deletions: %w", err)
}
if len(pending) == 0 {
return 0, nil
}
drained := 0
for _, p := range pending {
actx, cancel := context.WithTimeout(ctx, 30*time.Second)
err := client.Delete(actx, p.bucket, p.key)
cancel()
if err != nil {
if _, uerr := db.Conn.Exec(ctx, `
UPDATE pending_s3_deletions SET attempts = attempts + 1, last_error = $2
WHERE id = $1
`, p.id, err.Error()); uerr != nil {
return drained, fmt.Errorf("failed to record S3 deletion retry failure %s: %w", p.id, uerr)
}
log.Printf("Warning: retry-s3-deletions failed to delete profile picture %s (outbox %s): %v", p.key, p.id, err)
continue
}
if _, err := db.Conn.Exec(ctx, `DELETE FROM pending_s3_deletions WHERE id = $1`, p.id); err != nil {
return drained, fmt.Errorf("failed to drain S3 deletion outbox row %s: %w", p.id, err)
}
drained++
}
if drained > 0 {
log.Printf("[ERASURE] retry-s3-deletions drained %d pending S3 deletion outbox row(s)", drained)
}
return drained, nil
}