package billing import ( "context" "errors" "fmt" "math" "strings" "github.com/jackc/pgx/v5" ) func (s *Service) ListAccounts(ctx context.Context, tenantID string) ([]Account, error) { query := ` SELECT t.id::text, t.name, COALESCE(w.currency, $1), COALESCE(w.balance_micros, 0), COALESCE(w.reserved_micros, 0), COALESCE(w.balance_micros - w.reserved_micros, 0), COALESCE(w.updated_at, t.created_at) FROM tenants t LEFT JOIN tenant_wallets w ON w.tenant_id = t.id` args := []any{s.currency} if strings.TrimSpace(tenantID) != "" { query += ` WHERE t.id=$2` args = append(args, tenantID) } query += ` ORDER BY t.name` rows, err := s.db.Query(ctx, query, args...) if err != nil { return nil, fmt.Errorf("query billing accounts: %w", err) } defer rows.Close() result := make([]Account, 0) for rows.Next() { var item Account if err := rows.Scan(&item.TenantID, &item.TenantName, &item.Currency, &item.BalanceMicros, &item.ReservedMicros, &item.AvailableMicros, &item.UpdatedAt); err != nil { return nil, fmt.Errorf("scan billing account: %w", err) } result = append(result, item) } return result, rows.Err() } func (s *Service) ListLedger(ctx context.Context, tenantID string, limit int) ([]LedgerEntry, error) { if limit < 1 || limit > 500 { limit = 200 } query := ` SELECT id::text, tenant_id::text, COALESCE(project_id::text, ''), currency, amount_micros, balance_after_micros, kind, source_type, source_id, description, created_at FROM billing_ledger` args := []any{} if strings.TrimSpace(tenantID) != "" { query += ` WHERE tenant_id = $1 ORDER BY created_at DESC LIMIT $2` args = append(args, tenantID, limit) } else { query += ` ORDER BY created_at DESC LIMIT $1` args = append(args, limit) } rows, err := s.db.Query(ctx, query, args...) if err != nil { return nil, fmt.Errorf("query billing ledger: %w", err) } defer rows.Close() result := make([]LedgerEntry, 0) for rows.Next() { var item LedgerEntry if err := rows.Scan(&item.ID, &item.TenantID, &item.ProjectID, &item.Currency, &item.AmountMicros, &item.BalanceAfterMicros, &item.Kind, &item.SourceType, &item.SourceID, &item.Description, &item.CreatedAt); err != nil { return nil, fmt.Errorf("scan billing ledger: %w", err) } result = append(result, item) } return result, rows.Err() } func (s *Service) ListTopUpOrders(ctx context.Context, tenantID string, limit int) ([]TopUpOrder, error) { if limit < 1 || limit > 200 { limit = 50 } query := `SELECT id::text,tenant_id::text,amount_minor,amount_micros,currency,status,trigger_type, COALESCE(stripe_session_id,''),COALESCE(checkout_url,''),created_at,paid_at, COALESCE(stripe_customer_id,''),COALESCE(stripe_payment_intent_id,''),COALESCE(stripe_charge_id,''), COALESCE(stripe_invoice_id,''),COALESCE(invoice_url,''),COALESCE(invoice_pdf_url,''),COALESCE(receipt_url,''), refunded_micros,disputed_micros,reconciliation_status,reconciled_at,reconciliation_error FROM topup_orders` args := []any{} if strings.TrimSpace(tenantID) != "" { query += ` WHERE tenant_id=$1 ORDER BY created_at DESC LIMIT $2` args = []any{tenantID, limit} } else { query += ` ORDER BY created_at DESC LIMIT $1` args = []any{limit} } rows, err := s.db.Query(ctx, query, args...) if err != nil { return nil, fmt.Errorf("query top-up orders: %w", err) } defer rows.Close() result := make([]TopUpOrder, 0) for rows.Next() { var item TopUpOrder if err := rows.Scan(&item.ID, &item.TenantID, &item.AmountMinor, &item.AmountMicros, &item.Currency, &item.Status, &item.TriggerType, &item.StripeSessionID, &item.CheckoutURL, &item.CreatedAt, &item.PaidAt, &item.StripeCustomerID, &item.StripePaymentIntentID, &item.StripeChargeID, &item.StripeInvoiceID, &item.InvoiceURL, &item.InvoicePDFURL, &item.ReceiptURL, &item.RefundedMicros, &item.DisputedMicros, &item.ReconciliationStatus, &item.ReconciledAt, &item.ReconciliationError); err != nil { return nil, err } result = append(result, item) } return result, rows.Err() } func (s *Service) GetTopUpOrder(ctx context.Context, tenantID, orderID string) (TopUpOrder, error) { var result TopUpOrder query := `SELECT id::text,tenant_id::text,amount_minor,amount_micros,currency,status,trigger_type, COALESCE(stripe_session_id,''),COALESCE(checkout_url,''),created_at,paid_at, COALESCE(stripe_customer_id,''),COALESCE(stripe_payment_intent_id,''),COALESCE(stripe_charge_id,''), COALESCE(stripe_invoice_id,''),COALESCE(invoice_url,''),COALESCE(invoice_pdf_url,''),COALESCE(receipt_url,''), refunded_micros,disputed_micros,reconciliation_status,reconciled_at,reconciliation_error FROM topup_orders WHERE id=$1` args := []any{orderID} if strings.TrimSpace(tenantID) != "" { query += ` AND tenant_id=$2` args = append(args, tenantID) } err := s.db.QueryRow(ctx, query, args...).Scan(&result.ID, &result.TenantID, &result.AmountMinor, &result.AmountMicros, &result.Currency, &result.Status, &result.TriggerType, &result.StripeSessionID, &result.CheckoutURL, &result.CreatedAt, &result.PaidAt, &result.StripeCustomerID, &result.StripePaymentIntentID, &result.StripeChargeID, &result.StripeInvoiceID, &result.InvoiceURL, &result.InvoicePDFURL, &result.ReceiptURL, &result.RefundedMicros, &result.DisputedMicros, &result.ReconciliationStatus, &result.ReconciledAt, &result.ReconciliationError) if errors.Is(err, pgx.ErrNoRows) { return TopUpOrder{}, ErrTopUpOrderNotFound } if err != nil { return TopUpOrder{}, fmt.Errorf("get top-up order: %w", err) } return result, nil } func (s *Service) AdjustBalance(ctx context.Context, input AdjustmentInput) (LedgerEntry, error) { input.TenantID = strings.TrimSpace(input.TenantID) input.Description = normalizeDescription(input.Description) if input.TenantID == "" || input.AmountMicros == 0 { return LedgerEntry{}, ErrInvalidAmount } tx, err := s.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted}) if err != nil { return LedgerEntry{}, fmt.Errorf("begin balance adjustment: %w", err) } defer tx.Rollback(ctx) if _, err := tx.Exec(ctx, ` INSERT INTO tenant_wallets (tenant_id, currency) VALUES ($1, $2) ON CONFLICT (tenant_id) DO NOTHING`, input.TenantID, s.currency); err != nil { return LedgerEntry{}, fmt.Errorf("ensure adjustment wallet: %w", err) } var currency string var balance, held int64 if err := tx.QueryRow(ctx, `SELECT currency, balance_micros, reserved_micros FROM tenant_wallets WHERE tenant_id = $1 FOR UPDATE`, input.TenantID).Scan(¤cy, &balance, &held); err != nil { return LedgerEntry{}, fmt.Errorf("lock adjustment wallet: %w", err) } if currency != s.currency { return LedgerEntry{}, fmt.Errorf("tenant wallet currency %s does not match %s", currency, s.currency) } if input.AmountMicros > 0 && balance > math.MaxInt64-input.AmountMicros { return LedgerEntry{}, ErrInvalidAmount } newBalance := balance + input.AmountMicros if newBalance < held || newBalance < 0 { return LedgerEntry{}, ErrInsufficientBalance } if _, err := tx.Exec(ctx, `UPDATE tenant_wallets SET balance_micros = $2, updated_at = now() WHERE tenant_id = $1`, input.TenantID, newBalance); err != nil { return LedgerEntry{}, fmt.Errorf("apply balance adjustment: %w", err) } sourceID := "adj_" + randomHex(16) var result LedgerEntry 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,'adjustment','admin',$5,$6) RETURNING id::text, tenant_id::text, '', currency, amount_micros, balance_after_micros, kind, source_type, source_id, description, created_at`, input.TenantID, s.currency, input.AmountMicros, newBalance, sourceID, input.Description, ).Scan(&result.ID, &result.TenantID, &result.ProjectID, &result.Currency, &result.AmountMicros, &result.BalanceAfterMicros, &result.Kind, &result.SourceType, &result.SourceID, &result.Description, &result.CreatedAt) if err != nil { return LedgerEntry{}, fmt.Errorf("write adjustment ledger entry: %w", err) } if err := tx.Commit(ctx); err != nil { return LedgerEntry{}, fmt.Errorf("commit balance adjustment: %w", err) } return result, nil } // ReleaseUnmeteredReservation is an audited operational escape hatch for a // fail-closed success that cannot be reconciled. It never invents usage or // changes wallet balance; it only returns the existing hold to availability. func (s *Service) ReleaseUnmeteredReservation(ctx context.Context, tenantID, requestID string, input ReleaseReservationInput, actor ResolutionActor) (ReservationRelease, error) { tenantID = strings.TrimSpace(tenantID) requestID = strings.TrimSpace(requestID) reason := normalizeDescription(input.Reason) if tenantID == "" || requestID == "" || reason == "" || actor.ID == "" || actor.Type == "" { return ReservationRelease{}, fmt.Errorf("%w: tenant, request, reason, and actor are required", ErrReservationNotReleasable) } tx, err := s.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted}) if err != nil { return ReservationRelease{}, fmt.Errorf("begin reservation release: %w", err) } defer tx.Rollback(ctx) var result ReservationRelease var projectID, currency, status string if err := tx.QueryRow(ctx, `SELECT request_id,tenant_id::text,project_id::text,currency,reserved_micros,status FROM billing_reservations WHERE request_id=$1 AND tenant_id=$2 FOR UPDATE`, requestID, tenantID).Scan( &result.RequestID, &result.TenantID, &projectID, ¤cy, &result.ReservedMicros, &status); errors.Is(err, pgx.ErrNoRows) { return ReservationRelease{}, ErrReservationNotReleasable } else if err != nil { return ReservationRelease{}, fmt.Errorf("lock reservation for release: %w", err) } if status == "released" { var evidenceExists bool if err := tx.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM billing_ledger WHERE source_type='unmetered_reservation' AND source_id=$1)`, requestID).Scan(&evidenceExists); err != nil { return ReservationRelease{}, fmt.Errorf("read reservation release evidence: %w", err) } if !evidenceExists { return ReservationRelease{}, fmt.Errorf("%w: reservation was released by normal settlement", ErrReservationNotReleasable) } if err := tx.QueryRow(ctx, `SELECT COALESCE(settled_at,created_at) FROM billing_reservations WHERE request_id=$1`, requestID).Scan(&result.ReleasedAt); err != nil { return ReservationRelease{}, err } result.Status = status return result, tx.Commit(ctx) } if status != "metering_failed" { return ReservationRelease{}, fmt.Errorf("%w: reservation status is %s", ErrReservationNotReleasable, status) } var balance, held int64 if err := tx.QueryRow(ctx, `SELECT balance_micros,reserved_micros FROM tenant_wallets WHERE tenant_id=$1 FOR UPDATE`, tenantID).Scan(&balance, &held); err != nil { return ReservationRelease{}, fmt.Errorf("lock wallet for reservation release: %w", err) } if result.ReservedMicros > held { return ReservationRelease{}, errors.New("wallet reservation invariant violated during release") } if _, err := tx.Exec(ctx, `UPDATE tenant_wallets SET reserved_micros=reserved_micros-$2,updated_at=now() WHERE tenant_id=$1`, tenantID, result.ReservedMicros); err != nil { return ReservationRelease{}, fmt.Errorf("release wallet hold: %w", err) } if err := tx.QueryRow(ctx, `UPDATE billing_reservations SET status='released',settled_at=now() WHERE request_id=$1 RETURNING status,settled_at`, requestID).Scan(&result.Status, &result.ReleasedAt); err != nil { return ReservationRelease{}, fmt.Errorf("mark reservation released: %w", err) } command, err := tx.Exec(ctx, `UPDATE usage_events SET metering_status='released_unmetered' WHERE request_id=$1 AND tenant_id=$2 AND metering_status='missing' AND usage_reported=FALSE`, requestID, tenantID) if err != nil { return ReservationRelease{}, fmt.Errorf("mark unmetered usage resolved: %w", err) } if command.RowsAffected() != 1 { return ReservationRelease{}, errors.New("unmetered usage invariant violated during release") } description := fmt.Sprintf("Unmetered reservation released by %s %s: %s", actor.Type, actor.ID, reason) if _, err := tx.Exec(ctx, `INSERT INTO billing_ledger (tenant_id,project_id,currency,amount_micros,balance_after_micros,kind,source_type,source_id,description) VALUES ($1,$2,$3,0,$4,'release','unmetered_reservation',$5,$6) ON CONFLICT (source_type,source_id) DO NOTHING`, tenantID, projectID, currency, balance, requestID, description); err != nil { return ReservationRelease{}, fmt.Errorf("write reservation release evidence: %w", err) } if err := tx.Commit(ctx); err != nil { return ReservationRelease{}, fmt.Errorf("commit reservation release: %w", err) } return result, nil } func (s *Service) createTopUpOrder(ctx context.Context, input CheckoutInput) (string, int64, error) { if !s.stripeEnabled { return "", 0, ErrStripeDisabled } if strings.TrimSpace(input.TenantID) == "" || input.AmountMinor < s.minTopUpMinor || input.AmountMinor > s.maxTopUpMinor { return "", 0, ErrInvalidAmount } amountMicros, err := minorToMicros(s.currency, input.AmountMinor) if err != nil { return "", 0, err } var orderID string err = s.db.QueryRow(ctx, ` INSERT INTO topup_orders (tenant_id, amount_minor, amount_micros, currency) VALUES ($1,$2,$3,$4) RETURNING id::text`, input.TenantID, input.AmountMinor, amountMicros, s.currency, ).Scan(&orderID) if err != nil { return "", 0, fmt.Errorf("create top-up order: %w", err) } return orderID, amountMicros, nil } func minorToMicros(currency string, amount int64) (int64, error) { if amount <= 0 { return 0, ErrInvalidAmount } factor := int64(10_000) switch strings.ToLower(currency) { case "bif", "clp", "djf", "gnf", "jpy", "kmf", "krw", "mga", "pyg", "rwf", "ugx", "vnd", "vuv", "xaf", "xof", "xpf": factor = 1_000_000 case "bhd", "jod", "kwd", "omr", "tnd": factor = 1_000 } if amount > math.MaxInt64/factor { return 0, ErrInvalidAmount } return amount * factor, nil } func isNotFound(err error) bool { return errors.Is(err, pgx.ErrNoRows) }