summaryrefslogtreecommitdiff
path: root/internal/billing/auto_topup.go
diff options
context:
space:
mode:
Diffstat (limited to 'internal/billing/auto_topup.go')
-rw-r--r--internal/billing/auto_topup.go541
1 files changed, 541 insertions, 0 deletions
diff --git a/internal/billing/auto_topup.go b/internal/billing/auto_topup.go
new file mode 100644
index 0000000..a90405b
--- /dev/null
+++ b/internal/billing/auto_topup.go
@@ -0,0 +1,541 @@
+package billing
+
+import (
+ "context"
+ "errors"
+ "fmt"
+ "net/url"
+ "strings"
+ "time"
+
+ "github.com/jackc/pgx/v5"
+ "github.com/stripe/stripe-go/v86"
+)
+
+type stripeSetupIntentRetriever func(context.Context, string, *stripe.SetupIntentRetrieveParams) (*stripe.SetupIntent, error)
+type stripePaymentIntentCreator func(context.Context, *stripe.PaymentIntentCreateParams) (*stripe.PaymentIntent, error)
+type stripePaymentIntentRetriever func(context.Context, string, *stripe.PaymentIntentRetrieveParams) (*stripe.PaymentIntent, error)
+
+const autoTopUpAction = "auto_topup_setup"
+
+func (s *Service) defaultAutoTopUpAmountMinor() int64 {
+ amount := int64(2000)
+ if amount < s.minTopUpMinor {
+ amount = s.minTopUpMinor
+ }
+ if amount > s.maxTopUpMinor {
+ amount = s.maxTopUpMinor
+ }
+ return amount
+}
+
+func (s *Service) defaultAutoTopUpThresholdMicros() int64 {
+ amount, err := minorToMicros(s.currency, s.defaultAutoTopUpAmountMinor())
+ if err != nil || amount <= 0 {
+ return 0
+ }
+ if amount/4 > 5*microsPerUnit {
+ return 5 * microsPerUnit
+ }
+ return amount / 4
+}
+
+func (s *Service) GetAutoTopUpSettings(ctx context.Context, tenantID string) (AutoTopUpSettings, error) {
+ tenantID = strings.TrimSpace(tenantID)
+ if tenantID == "" {
+ return AutoTopUpSettings{}, ErrBillingAccountNotFound
+ }
+ var result AutoTopUpSettings
+ var paymentMethodID string
+ err := s.db.QueryRow(ctx, `
+ SELECT t.id::text, COALESCE(w.currency,$2), $3::boolean,
+ COALESCE(a.enabled,FALSE), COALESCE(a.threshold_micros,$4),
+ COALESCE(a.topup_amount_minor,$5), COALESCE(a.stripe_payment_method_id,''),
+ COALESCE(a.payment_method_type,''), COALESCE(a.payment_method_brand,''),
+ COALESCE(a.payment_method_last4,''), COALESCE(a.payment_method_exp_month,0),
+ COALESCE(a.payment_method_exp_year,0), COALESCE(a.status,'not_configured'),
+ COALESCE(a.last_error,''), a.last_attempt_at, a.last_succeeded_at,
+ a.next_attempt_at, COALESCE(a.updated_at,t.created_at)
+ FROM tenants t
+ LEFT JOIN tenant_wallets w ON w.tenant_id=t.id
+ LEFT JOIN tenant_auto_topup_settings a ON a.tenant_id=t.id
+ WHERE t.id=$1`, tenantID, s.currency, s.stripeEnabled, s.defaultAutoTopUpThresholdMicros(), s.defaultAutoTopUpAmountMinor()).
+ Scan(&result.TenantID, &result.Currency, &result.StripeEnabled, &result.Enabled,
+ &result.ThresholdMicros, &result.TopUpAmountMinor, &paymentMethodID,
+ &result.PaymentMethodType, &result.PaymentMethodBrand, &result.PaymentMethodLast4,
+ &result.PaymentMethodExpMonth, &result.PaymentMethodExpYear, &result.Status,
+ &result.LastError, &result.LastAttemptAt, &result.LastSucceededAt,
+ &result.NextAttemptAt, &result.UpdatedAt)
+ if errors.Is(err, pgx.ErrNoRows) {
+ return AutoTopUpSettings{}, ErrBillingAccountNotFound
+ }
+ if err != nil {
+ return AutoTopUpSettings{}, fmt.Errorf("query automatic top-up settings: %w", err)
+ }
+ result.PaymentMethodConfigured = paymentMethodID != ""
+ return result, nil
+}
+
+func (s *Service) UpdateAutoTopUp(ctx context.Context, input UpdateAutoTopUpInput) (AutoTopUpSettings, error) {
+ input.TenantID = strings.TrimSpace(input.TenantID)
+ if input.TenantID == "" || input.ThresholdMicros < 0 || input.TopUpAmountMinor < s.minTopUpMinor || input.TopUpAmountMinor > s.maxTopUpMinor {
+ return AutoTopUpSettings{}, ErrInvalidAmount
+ }
+ topUpMicros, err := minorToMicros(s.currency, input.TopUpAmountMinor)
+ if err != nil || topUpMicros <= input.ThresholdMicros {
+ return AutoTopUpSettings{}, fmt.Errorf("%w: automatic top-up amount must exceed the balance threshold", ErrInvalidAmount)
+ }
+ if input.Enabled && !s.stripeEnabled {
+ return AutoTopUpSettings{}, ErrStripeDisabled
+ }
+ tx, err := s.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted})
+ if err != nil {
+ return AutoTopUpSettings{}, err
+ }
+ defer tx.Rollback(ctx)
+ if _, err := tx.Exec(ctx, `INSERT INTO tenant_auto_topup_settings (tenant_id,threshold_micros,topup_amount_minor)
+ VALUES ($1,$2,$3) ON CONFLICT (tenant_id) DO NOTHING`, input.TenantID, input.ThresholdMicros, input.TopUpAmountMinor); err != nil {
+ return AutoTopUpSettings{}, fmt.Errorf("initialize automatic top-up settings: %w", err)
+ }
+ var paymentMethodID, status string
+ if err := tx.QueryRow(ctx, `SELECT COALESCE(stripe_payment_method_id,''),status FROM tenant_auto_topup_settings WHERE tenant_id=$1 FOR UPDATE`, input.TenantID).Scan(&paymentMethodID, &status); err != nil {
+ if errors.Is(err, pgx.ErrNoRows) {
+ return AutoTopUpSettings{}, ErrBillingAccountNotFound
+ }
+ return AutoTopUpSettings{}, err
+ }
+ if input.Enabled && paymentMethodID == "" {
+ return AutoTopUpSettings{}, ErrPaymentMethodRequired
+ }
+ if input.Enabled && status == "action_required" {
+ return AutoTopUpSettings{}, ErrAutoTopUpNeedsAttention
+ }
+ next := any(nil)
+ newStatus := "not_configured"
+ if paymentMethodID != "" {
+ newStatus = "ready"
+ }
+ if input.Enabled {
+ newStatus = "ready"
+ next = time.Now().UTC()
+ }
+ if _, err := tx.Exec(ctx, `UPDATE tenant_auto_topup_settings SET enabled=$2,threshold_micros=$3,topup_amount_minor=$4,
+ status=$5,last_error=CASE WHEN $2 THEN '' ELSE last_error END,
+ next_attempt_at=$6,updated_at=now() WHERE tenant_id=$1`, input.TenantID, input.Enabled,
+ input.ThresholdMicros, input.TopUpAmountMinor, newStatus, next); err != nil {
+ return AutoTopUpSettings{}, fmt.Errorf("update automatic top-up settings: %w", err)
+ }
+ if err := tx.Commit(ctx); err != nil {
+ return AutoTopUpSettings{}, err
+ }
+ return s.GetAutoTopUpSettings(ctx, input.TenantID)
+}
+
+func (s *Service) DisableAutoTopUp(ctx context.Context, tenantID string) (AutoTopUpSettings, error) {
+ tenantID = strings.TrimSpace(tenantID)
+ if tenantID == "" {
+ return AutoTopUpSettings{}, ErrBillingAccountNotFound
+ }
+ if _, err := s.db.Exec(ctx, `UPDATE tenant_auto_topup_settings SET enabled=FALSE,status=CASE WHEN stripe_payment_method_id IS NULL THEN 'not_configured' ELSE 'ready' END,next_attempt_at=NULL,updated_at=now() WHERE tenant_id=$1`, tenantID); err != nil {
+ return AutoTopUpSettings{}, err
+ }
+ return s.GetAutoTopUpSettings(ctx, tenantID)
+}
+
+// CreateAutoTopUpSetupSession opens a Stripe-hosted SetupIntent flow. Stripe
+// owns card collection; this service only receives a PaymentMethod ID after a
+// signed webhook confirms that the setup succeeded.
+func (s *Service) CreateAutoTopUpSetupSession(ctx context.Context, input AutoTopUpSetupInput) (AutoTopUpSetupResult, error) {
+ if !s.stripeEnabled || s.createStripeCheckout == nil {
+ return AutoTopUpSetupResult{}, ErrStripeDisabled
+ }
+ input.TenantID = strings.TrimSpace(input.TenantID)
+ if input.TenantID == "" {
+ return AutoTopUpSetupResult{}, ErrBillingAccountNotFound
+ }
+ if _, err := s.GetAutoTopUpSettings(ctx, input.TenantID); err != nil {
+ return AutoTopUpSetupResult{}, err
+ }
+ customerID, err := s.ensureStripeCustomer(ctx, input.TenantID)
+ if err != nil {
+ return AutoTopUpSetupResult{}, err
+ }
+ params := &stripe.CheckoutSessionCreateParams{
+ Mode: stripe.String(string(stripe.CheckoutSessionModeSetup)),
+ Currency: stripe.String(s.currency),
+ ClientReferenceID: stripe.String(input.TenantID),
+ IntegrationIdentifier: stripe.String(s.integrationIdentifier),
+ SuccessURL: stripe.String(autoTopUpReturnURL(s.stripeSuccessURL, true)),
+ CancelURL: stripe.String(autoTopUpReturnURL(s.stripeCancelURL, false)),
+ Metadata: map[string]string{
+ "aigw_action": autoTopUpAction,
+ "aigw_tenant_id": input.TenantID,
+ },
+ }
+ if customerID != "" {
+ params.Customer = stripe.String(customerID)
+ } else {
+ params.CustomerCreation = stripe.String(string(stripe.CheckoutSessionCustomerCreationAlways))
+ if strings.TrimSpace(input.CustomerEmail) != "" {
+ params.CustomerEmail = stripe.String(strings.TrimSpace(input.CustomerEmail))
+ }
+ }
+ params.SetIdempotencyKey("aigw_autotopup_setup_" + randomHex(16))
+ session, err := s.createStripeCheckout(ctx, params)
+ if err != nil {
+ return AutoTopUpSetupResult{}, fmt.Errorf("create automatic top-up setup session: %w", err)
+ }
+ if session.ID == "" || session.URL == "" {
+ return AutoTopUpSetupResult{}, errors.New("Stripe returned an incomplete setup session")
+ }
+ if _, err := s.db.Exec(ctx, `INSERT INTO tenant_auto_topup_settings (tenant_id,threshold_micros,topup_amount_minor,stripe_setup_session_id)
+ VALUES ($1,$2,$3,$4) ON CONFLICT (tenant_id) DO UPDATE SET stripe_setup_session_id=EXCLUDED.stripe_setup_session_id,updated_at=now()`,
+ input.TenantID, s.defaultAutoTopUpThresholdMicros(), s.defaultAutoTopUpAmountMinor(), session.ID); err != nil {
+ return AutoTopUpSetupResult{}, fmt.Errorf("persist automatic top-up setup session: %w", err)
+ }
+ return AutoTopUpSetupResult{SessionID: session.ID, URL: session.URL}, nil
+}
+
+func autoTopUpReturnURL(raw string, success bool) string {
+ parsed, err := url.Parse(raw)
+ if err != nil {
+ return raw
+ }
+ query := parsed.Query()
+ query.Set("autotopup", "setup")
+ if success {
+ query.Set("session_id", "{CHECKOUT_SESSION_ID}")
+ } else {
+ query.Set("autotopup", "cancel")
+ query.Del("session_id")
+ }
+ parsed.RawQuery = strings.ReplaceAll(query.Encode(), url.QueryEscape("{CHECKOUT_SESSION_ID}"), "{CHECKOUT_SESSION_ID}")
+ return parsed.String()
+}
+
+func (s *Service) processAutoTopUpSetupEvent(ctx context.Context, event stripe.Event, session *stripe.CheckoutSession) error {
+ if s.retrieveStripeSetupIntent == nil || session == nil || event.ID == "" || session.ID == "" || session.ClientReferenceID == "" {
+ return ErrInvalidAmount
+ }
+ if event.Type != stripe.EventTypeCheckoutSessionCompleted {
+ return nil
+ }
+ setupIntentID := ""
+ if session.SetupIntent != nil {
+ setupIntentID = session.SetupIntent.ID
+ }
+ if setupIntentID == "" {
+ return ErrPaymentMethodRequired
+ }
+ intent, err := s.retrieveStripeSetupIntent(ctx, setupIntentID, &stripe.SetupIntentRetrieveParams{})
+ if err != nil {
+ return fmt.Errorf("retrieve automatic top-up setup intent: %w", err)
+ }
+ if intent == nil || intent.Status != stripe.SetupIntentStatusSucceeded || intent.PaymentMethod == nil || intent.PaymentMethod.ID == "" {
+ return ErrPaymentMethodRequired
+ }
+ tx, err := s.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted})
+ if err != nil {
+ return err
+ }
+ defer tx.Rollback(ctx)
+ if _, err := tx.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtextextended($1,0))`, event.ID); err != nil {
+ return err
+ }
+ if _, err := tx.Exec(ctx, `INSERT INTO stripe_webhook_events (event_id,event_type) VALUES ($1,$2) ON CONFLICT DO NOTHING`, event.ID, string(event.Type)); err != nil {
+ return err
+ }
+ var storedSetupSession string
+ if err := tx.QueryRow(ctx, `SELECT COALESCE(stripe_setup_session_id,'') FROM tenant_auto_topup_settings WHERE tenant_id=$1 FOR UPDATE`, session.ClientReferenceID).Scan(&storedSetupSession); err != nil || storedSetupSession != session.ID {
+ return ErrInvalidAmount
+ }
+ customerID := ""
+ if session.Customer != nil {
+ customerID = session.Customer.ID
+ }
+ if customerID == "" && intent.Customer != nil {
+ customerID = intent.Customer.ID
+ }
+ if customerID == "" {
+ return ErrPaymentMethodRequired
+ }
+ email := session.CustomerEmail
+ if session.CustomerDetails != nil && session.CustomerDetails.Email != "" {
+ email = session.CustomerDetails.Email
+ }
+ if _, err := tx.Exec(ctx, `INSERT INTO stripe_customers (tenant_id,stripe_customer_id,email) VALUES ($1,$2,$3)
+ ON CONFLICT (tenant_id) DO UPDATE SET stripe_customer_id=EXCLUDED.stripe_customer_id,
+ email=CASE WHEN EXCLUDED.email='' THEN stripe_customers.email ELSE EXCLUDED.email END,updated_at=now()`, session.ClientReferenceID, customerID, email); err != nil {
+ return err
+ }
+ methodType, brand, last4 := intent.PaymentMethod.Type, "", ""
+ var expMonth, expYear int64
+ if intent.PaymentMethod.Card != nil {
+ brand, last4 = string(intent.PaymentMethod.Card.Brand), intent.PaymentMethod.Card.Last4
+ expMonth, expYear = intent.PaymentMethod.Card.ExpMonth, intent.PaymentMethod.Card.ExpYear
+ }
+ if _, err := tx.Exec(ctx, `INSERT INTO tenant_auto_topup_settings
+ (tenant_id,threshold_micros,topup_amount_minor,stripe_payment_method_id,payment_method_type,payment_method_brand,payment_method_last4,payment_method_exp_month,payment_method_exp_year,stripe_setup_session_id,status,last_error,failure_count,updated_at)
+ VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,'ready','',0,now())
+ ON CONFLICT (tenant_id) DO UPDATE SET stripe_payment_method_id=EXCLUDED.stripe_payment_method_id,
+ payment_method_type=EXCLUDED.payment_method_type,payment_method_brand=EXCLUDED.payment_method_brand,
+ payment_method_last4=EXCLUDED.payment_method_last4,payment_method_exp_month=EXCLUDED.payment_method_exp_month,
+ payment_method_exp_year=EXCLUDED.payment_method_exp_year,stripe_setup_session_id=EXCLUDED.stripe_setup_session_id,
+ status='ready',last_error='',failure_count=0,next_attempt_at=NULL,updated_at=now()`,
+ session.ClientReferenceID, s.defaultAutoTopUpThresholdMicros(), s.defaultAutoTopUpAmountMinor(), intent.PaymentMethod.ID,
+ methodType, brand, last4, expMonth, expYear, session.ID); err != nil {
+ return err
+ }
+ if _, err := tx.Exec(ctx, `UPDATE stripe_webhook_events SET processed_at=now(),processing_error='' WHERE event_id=$1`, event.ID); err != nil {
+ return err
+ }
+ return tx.Commit(ctx)
+}
+
+// processAutoTopUpOnce claims one eligible tenant before making a Stripe call.
+// The row lock and unique pending-order index make this safe across gateways.
+func (s *Service) processAutoTopUpOnce(ctx context.Context) (bool, error) {
+ if !s.stripeEnabled || s.createStripePaymentIntent == nil {
+ return false, nil
+ }
+ tx, err := s.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted})
+ if err != nil {
+ return false, err
+ }
+ defer tx.Rollback(ctx)
+ var tenantID, customerID, customerEmail, paymentMethodID, currency, orderID string
+ var amountMinor, threshold, balance, reserved int64
+ row := tx.QueryRow(ctx, `
+ SELECT a.tenant_id::text,c.stripe_customer_id,c.email,a.stripe_payment_method_id,w.currency,
+ a.topup_amount_minor,a.threshold_micros,w.balance_micros,w.reserved_micros
+ ,COALESCE(o.id::text,'')
+ FROM tenant_auto_topup_settings a
+ JOIN stripe_customers c ON c.tenant_id=a.tenant_id
+ JOIN tenant_wallets w ON w.tenant_id=a.tenant_id
+ LEFT JOIN LATERAL (SELECT id FROM topup_orders WHERE tenant_id=a.tenant_id AND trigger_type='auto' AND status='pending' ORDER BY created_at DESC LIMIT 1) o ON TRUE
+ WHERE a.enabled AND a.stripe_payment_method_id IS NOT NULL
+ AND a.status IN ('ready','failed','charging')
+ AND (a.next_attempt_at IS NULL OR a.next_attempt_at <= now())
+ AND w.balance_micros-w.reserved_micros <= a.threshold_micros
+ ORDER BY w.balance_micros-w.reserved_micros,a.updated_at
+ FOR UPDATE OF a SKIP LOCKED LIMIT 1`)
+ if err := row.Scan(&tenantID, &customerID, &customerEmail, &paymentMethodID, &currency, &amountMinor, &threshold, &balance, &reserved, &orderID); err != nil {
+ if errors.Is(err, pgx.ErrNoRows) {
+ return false, tx.Commit(ctx)
+ }
+ return false, err
+ }
+ if balance-reserved > threshold {
+ return false, tx.Commit(ctx)
+ }
+ amountMicros, err := minorToMicros(currency, amountMinor)
+ if err != nil {
+ return false, err
+ }
+ if orderID == "" {
+ if err := tx.QueryRow(ctx, `INSERT INTO topup_orders (tenant_id,amount_minor,amount_micros,currency,trigger_type)
+ VALUES ($1,$2,$3,$4,'auto') RETURNING id::text`, tenantID, amountMinor, amountMicros, currency).Scan(&orderID); err != nil {
+ return false, err
+ }
+ }
+ if _, err := tx.Exec(ctx, `UPDATE tenant_auto_topup_settings SET status='charging',last_attempt_at=now(),next_attempt_at=now()+interval '15 minutes',updated_at=now() WHERE tenant_id=$1`, tenantID); err != nil {
+ return false, err
+ }
+ if err := tx.Commit(ctx); err != nil {
+ return false, err
+ }
+ params := &stripe.PaymentIntentCreateParams{
+ Amount: stripe.Int64(amountMinor), Currency: stripe.String(currency), Customer: stripe.String(customerID),
+ PaymentMethod: stripe.String(paymentMethodID), Confirm: stripe.Bool(true), OffSession: stripe.Bool(true),
+ ErrorOnRequiresAction: stripe.Bool(true), Description: stripe.String("AIGW automatic prepaid balance top-up"),
+ Metadata: map[string]string{"aigw_action": "auto_topup", "aigw_tenant_id": tenantID, "aigw_topup_order_id": orderID},
+ }
+ if strings.TrimSpace(customerEmail) != "" {
+ params.ReceiptEmail = stripe.String(strings.TrimSpace(customerEmail))
+ }
+ params.SetIdempotencyKey("aigw_autotopup_" + orderID)
+ intent, callErr := s.createStripePaymentIntent(ctx, params)
+ if callErr != nil {
+ var stripeErr *stripe.Error
+ if errors.As(callErr, &stripeErr) && stripeErr.PaymentIntent != nil {
+ intent = stripeErr.PaymentIntent
+ if intent.ID != "" {
+ _, _ = s.db.Exec(ctx, `UPDATE topup_orders SET stripe_payment_intent_id=$2,stripe_customer_id=$3 WHERE id=$1`, orderID, intent.ID, customerID)
+ }
+ if intent.Status == stripe.PaymentIntentStatusSucceeded {
+ return true, s.creditAutoTopUpPaymentIntent(ctx, intent)
+ }
+ return true, s.applyAutoTopUpPaymentIntentFailure(ctx, intent)
+ }
+ return true, s.scheduleAutoTopUpRetry(ctx, tenantID, orderID, callErr)
+ }
+ if intent == nil || intent.ID == "" {
+ return true, s.failAutoTopUp(ctx, tenantID, orderID, errors.New("Stripe returned an incomplete automatic top-up PaymentIntent"))
+ }
+ if _, err := s.db.Exec(ctx, `UPDATE topup_orders SET stripe_payment_intent_id=$2,stripe_customer_id=$3 WHERE id=$1`, orderID, intent.ID, customerID); err != nil {
+ return true, err
+ }
+ if intent.Status == stripe.PaymentIntentStatusSucceeded {
+ return true, s.creditAutoTopUpPaymentIntent(ctx, intent)
+ }
+ if intent.Status == stripe.PaymentIntentStatusProcessing {
+ return true, nil
+ }
+ if intent.Status == stripe.PaymentIntentStatusRequiresAction || intent.Status == stripe.PaymentIntentStatusRequiresPaymentMethod || intent.Status == stripe.PaymentIntentStatusCanceled {
+ return true, s.markAutoTopUpAttention(ctx, tenantID, orderID, ErrAutoTopUpNeedsAttention)
+ }
+ return true, s.failAutoTopUp(ctx, tenantID, orderID, fmt.Errorf("automatic top-up PaymentIntent ended in status %s", intent.Status))
+}
+
+func (s *Service) scheduleAutoTopUpRetry(ctx context.Context, tenantID, orderID string, cause error) error {
+ message := truncateError(cause)
+ _, err := s.db.Exec(ctx, `UPDATE topup_orders SET reconciliation_error=$2 WHERE id=$1 AND status='pending'`, orderID, message)
+ if err != nil {
+ return errors.Join(cause, err)
+ }
+ _, settingsErr := s.db.Exec(ctx, `UPDATE tenant_auto_topup_settings SET status='failed',failure_count=failure_count+1,last_error=$2,
+ next_attempt_at=now()+interval '1 hour',updated_at=now() WHERE tenant_id=$1`, tenantID, message)
+ if settingsErr != nil {
+ return errors.Join(cause, settingsErr)
+ }
+ return cause
+}
+
+func (s *Service) failAutoTopUp(ctx context.Context, tenantID, orderID string, cause error) error {
+ message := truncateError(cause)
+ _, err := s.db.Exec(ctx, `UPDATE topup_orders SET status='failed',reconciliation_error=$2 WHERE id=$1`, orderID, message)
+ if err != nil {
+ return errors.Join(cause, err)
+ }
+ _, settingsErr := s.db.Exec(ctx, `UPDATE tenant_auto_topup_settings SET status='failed',failure_count=failure_count+1,last_error=$2,
+ next_attempt_at=now()+interval '1 hour',updated_at=now() WHERE tenant_id=$1`, tenantID, message)
+ if settingsErr != nil {
+ return errors.Join(cause, settingsErr)
+ }
+ return cause
+}
+
+func (s *Service) markAutoTopUpAttention(ctx context.Context, tenantID, orderID string, cause error) error {
+ message := truncateError(cause)
+ _, err := s.db.Exec(ctx, `UPDATE topup_orders SET status='failed',reconciliation_error=$2 WHERE id=$1`, orderID, message)
+ if err != nil {
+ return err
+ }
+ _, err = s.db.Exec(ctx, `UPDATE tenant_auto_topup_settings SET enabled=FALSE,status='action_required',last_error=$2,next_attempt_at=NULL,updated_at=now() WHERE tenant_id=$1`, tenantID, message)
+ return err
+}
+
+func (s *Service) creditAutoTopUpPaymentIntent(ctx context.Context, intent *stripe.PaymentIntent) error {
+ if intent == nil || intent.ID == "" {
+ return ErrInvalidAmount
+ }
+ tx, err := s.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted})
+ if err != nil {
+ return err
+ }
+ defer tx.Rollback(ctx)
+ if _, err := tx.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtextextended($1,0))`, intent.ID); err != nil {
+ return err
+ }
+ if err := s.applyAutoTopUpPaymentIntentTx(ctx, tx, intent); err != nil {
+ return err
+ }
+ return tx.Commit(ctx)
+}
+
+func (s *Service) applyAutoTopUpPaymentIntentTx(ctx context.Context, tx pgx.Tx, intent *stripe.PaymentIntent) error {
+ if intent.Status != stripe.PaymentIntentStatusSucceeded || intent.Metadata["aigw_action"] != "auto_topup" {
+ return nil
+ }
+ orderID, tenantID := intent.Metadata["aigw_topup_order_id"], intent.Metadata["aigw_tenant_id"]
+ if orderID == "" || tenantID == "" {
+ return ErrInvalidAmount
+ }
+ var amountMinor, amountMicros int64
+ var currency, status, storedPI string
+ if err := tx.QueryRow(ctx, `SELECT amount_minor,amount_micros,currency,status,COALESCE(stripe_payment_intent_id,'') FROM topup_orders WHERE id=$1 AND tenant_id=$2 FOR UPDATE`, orderID, tenantID).Scan(&amountMinor, &amountMicros, &currency, &status, &storedPI); err != nil {
+ return err
+ }
+ if intent.Amount != amountMinor || string(intent.Currency) != currency || (storedPI != "" && storedPI != intent.ID) {
+ return ErrInvalidAmount
+ }
+ if intent.AmountReceived != 0 && intent.AmountReceived != amountMinor {
+ return ErrInvalidAmount
+ }
+ if intent.Customer != nil && intent.Customer.ID != "" {
+ if _, err := tx.Exec(ctx, `UPDATE topup_orders SET stripe_customer_id=$2 WHERE id=$1`, orderID, intent.Customer.ID); err != nil {
+ return err
+ }
+ }
+ if status != "paid" {
+ if _, err := tx.Exec(ctx, `INSERT INTO tenant_wallets (tenant_id,currency) VALUES ($1,$2) ON CONFLICT DO NOTHING`, tenantID, currency); err != nil {
+ return err
+ }
+ var balance, reserved int64
+ if err := tx.QueryRow(ctx, `SELECT balance_micros,reserved_micros FROM tenant_wallets WHERE tenant_id=$1 FOR UPDATE`, tenantID).Scan(&balance, &reserved); err != nil {
+ return err
+ }
+ if balance > int64(^uint64(0)>>1)-amountMicros {
+ return ErrInvalidAmount
+ }
+ newBalance := balance + amountMicros
+ if newBalance < reserved {
+ return ErrInsufficientBalance
+ }
+ var inserted int64
+ if err := tx.QueryRow(ctx, `INSERT INTO billing_ledger (tenant_id,currency,amount_micros,balance_after_micros,kind,source_type,source_id,description)
+ VALUES ($1,$2,$3,$4,'topup','stripe_payment_intent',$5,'Automatic prepaid balance top-up')
+ ON CONFLICT (source_type,source_id) DO NOTHING RETURNING amount_micros`, tenantID, currency, amountMicros, newBalance, intent.ID).Scan(&inserted); err != nil && !errors.Is(err, pgx.ErrNoRows) {
+ return err
+ }
+ if inserted != 0 {
+ if _, err := tx.Exec(ctx, `UPDATE tenant_wallets SET balance_micros=$2,updated_at=now() WHERE tenant_id=$1`, tenantID, newBalance); err != nil {
+ return err
+ }
+ }
+ }
+ if _, err := tx.Exec(ctx, `UPDATE topup_orders SET status='paid',paid_at=COALESCE(paid_at,now()),stripe_payment_intent_id=$2,reconciliation_status='ok',reconciled_at=now(),reconciliation_error='' WHERE id=$1`, orderID, intent.ID); err != nil {
+ return err
+ }
+ _, err := tx.Exec(ctx, `UPDATE tenant_auto_topup_settings SET status='ready',last_error='',failure_count=0,last_succeeded_at=now(),next_attempt_at=NULL,updated_at=now() WHERE tenant_id=$1`, tenantID)
+ return err
+}
+
+func (s *Service) applyAutoTopUpPaymentIntentFailure(ctx context.Context, intent *stripe.PaymentIntent) error {
+ if intent == nil || intent.Metadata["aigw_action"] != "auto_topup" {
+ return nil
+ }
+ tenantID, orderID := intent.Metadata["aigw_tenant_id"], intent.Metadata["aigw_topup_order_id"]
+ message := "automatic top-up payment failed"
+ if intent.LastPaymentError != nil && intent.LastPaymentError.Msg != "" {
+ message = intent.LastPaymentError.Msg
+ }
+ if intent.Status == stripe.PaymentIntentStatusRequiresAction || intent.Status == stripe.PaymentIntentStatusRequiresPaymentMethod || intent.Status == stripe.PaymentIntentStatusCanceled {
+ return s.markAutoTopUpAttention(ctx, tenantID, orderID, errors.New(message))
+ }
+ return s.failAutoTopUp(ctx, tenantID, orderID, errors.New(message))
+}
+
+func (s *Service) applyAutoTopUpPaymentIntentFailureTx(ctx context.Context, tx pgx.Tx, intent *stripe.PaymentIntent) error {
+ if intent == nil || intent.Metadata["aigw_action"] != "auto_topup" {
+ return nil
+ }
+ tenantID, orderID := intent.Metadata["aigw_tenant_id"], intent.Metadata["aigw_topup_order_id"]
+ if tenantID == "" || orderID == "" {
+ return ErrInvalidAmount
+ }
+ message := "automatic top-up payment failed"
+ if intent.LastPaymentError != nil && intent.LastPaymentError.Msg != "" {
+ message = truncateError(errors.New(intent.LastPaymentError.Msg))
+ }
+ if _, err := tx.Exec(ctx, `UPDATE topup_orders SET status='failed',reconciliation_error=$2,stripe_payment_intent_id=COALESCE(NULLIF($3,''),stripe_payment_intent_id) WHERE id=$1 AND status='pending'`, orderID, message, intent.ID); err != nil {
+ return err
+ }
+ status, enabled := "failed", true
+ if intent.Status == stripe.PaymentIntentStatusRequiresAction || intent.Status == stripe.PaymentIntentStatusRequiresPaymentMethod || intent.Status == stripe.PaymentIntentStatusCanceled {
+ status, enabled = "action_required", false
+ }
+ _, err := tx.Exec(ctx, `UPDATE tenant_auto_topup_settings SET enabled=$2,status=$3,last_error=$4,failure_count=failure_count+1,
+ next_attempt_at=CASE WHEN $2 THEN now()+interval '1 hour' ELSE NULL END,updated_at=now() WHERE tenant_id=$1`, tenantID, enabled, status, message)
+ return err
+}