- tip gate: CreateTipPayment saved-card 2FA gate now has scaTokenizedSavedCard skip matching every other charge surface (booking, terminal, gift-card); isSCATokenizeResultShape escape added to tip SAVE gate - webhook: align UPDATE clears VAT fields before re-apply (matches sweep rescue); 503 unknown-event tracking with 24h timeout notification via square_webhook_events table - cash-tip: cashChargeBasePence no longer restores campaign or subtracts loyalty — overcharge and tip shortfall fixed; 2FA dead code remnants removed from gift-card buy flow; TwoFactorCodeInput help text deconfused; refund pre-fill unit mismatch fixed (pounds vs pence); SCA buyer names split from full_name; passwordless delete UI accepts empty password - lockout: successful current-password clears shared failed_attempts/locked_until (victim can recover from login lockout via password change); passwordless delete condition changed to require 2FA only in enforced env - erasure: stale-guest batch erasure persists Square card/customer targets to durable outbox before NULLing them (crash-safe); S3 deletion retry capped at 10 attempts with admin notification; S3_PROFILE_PICS_BUCKET startup check added - env parsing: IsExplicitDevOrMockEnv and Square HTTP client base-URL switch now normalize (ToLower+TrimSpace) for consistency - auth: change-password/delete-account get per-user rate limiters (10/min); consume param dead code suppressed with TODO - frontend: 2FA/SCA dead code removed from gift-card buy flow, TwoFactorCodeInput help text fixed, refund pre-fill unit mismatch fixed, buyer names populated from full_name Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent)
691 lines
26 KiB
Go
691 lines
26 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
|
|
}
|
|
|
|
// maxS3DeletionRetries is the cap on retry attempts for a pending S3 deletion
|
|
// outbox row. After this many consecutive failures the row is marked as
|
|
// final-failed and a critical admin notification is raised — the operator must
|
|
// investigate and delete the object manually.
|
|
const maxS3DeletionRetries = 10
|
|
|
|
// 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. After maxS3DeletionRetries
|
|
// consecutive failures the row is marked as final-failed and a critical admin
|
|
// notification is raised — the operator must investigate and delete the object
|
|
// manually. 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, attempts
|
|
FROM pending_s3_deletions
|
|
WHERE attempts < $1
|
|
ORDER BY created_at
|
|
`, maxS3DeletionRetries)
|
|
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
|
|
attempts int
|
|
}
|
|
var pending []pendingDelete
|
|
for rows.Next() {
|
|
var p pendingDelete
|
|
if err := rows.Scan(&p.id, &p.bucket, &p.key, &p.attempts); 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 {
|
|
// Truncate the error to 200 chars to prevent unbounded growth.
|
|
errMsg := err.Error()
|
|
if len(errMsg) > 200 {
|
|
errMsg = errMsg[:200]
|
|
}
|
|
newAttempts := p.attempts + 1
|
|
if _, uerr := db.Conn.Exec(ctx, `
|
|
UPDATE pending_s3_deletions SET attempts = $2, last_error = $3
|
|
WHERE id = $1
|
|
`, p.id, newAttempts, errMsg); 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)
|
|
// After maxS3DeletionRetries consecutive failures, raise a critical
|
|
// admin notification and stop retrying — the operator must investigate.
|
|
if newAttempts >= maxS3DeletionRetries {
|
|
user.InsertSquareErasureCriticalNotification(ctx, "s3:"+p.id)
|
|
log.Printf("CRITICAL: retry-s3-deletions exhausted %d attempts for outbox %s (key %s) — raising admin notification; operator must delete the object manually", maxS3DeletionRetries, p.id, p.key)
|
|
}
|
|
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
|
|
}
|