From 41e322c53d7b4b796eb377d0df9c29ecd10ba431 Mon Sep 17 00:00:00 2001 From: Chia Date: Thu, 6 Aug 2026 09:29:41 +1200 Subject: feat: complete commercial control plane, billing, auth, and model catalog - add PostgreSQL control-plane persistence with Redis-degraded hot reload - implement prepaid balance, usage ledger, Stripe top-up and reconciliation - add registration, email verification, password reset, invitations and RBAC - support TOTP, Passkey MFA, device sessions, quotas and rate limits - add tenant billing profiles, audit logs and operational readiness checks - build authenticated admin console, Quickstart, Playground and usage analytics - add public model catalog with pricing, filtering and cost estimation - support OpenAI Responses providers and provider health failover - validate real upstream usage reporting and balance settlement --- internal/billing/auto_topup.go | 541 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 541 insertions(+) create mode 100644 internal/billing/auto_topup.go (limited to 'internal/billing/auto_topup.go') 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, ¤cy, &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, ¤cy, &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 +} -- cgit v1.2.3