package fx

import (
	"context"
	"database/sql"
	"errors"
	"fmt"
	"math"
	"strings"
	"time"

	"github.com/jackc/pgx/v5"
	"github.com/jackc/pgx/v5/pgconn"
	"github.com/jackc/pgx/v5/pgxpool"

	"github.com/niels/banking-app/backend/internal/domain"
	"github.com/niels/banking-app/backend/internal/ledger"
	"github.com/niels/banking-app/backend/internal/risk"
)

const (
	rateScale      int64 = 1_000_000
	basisPointBase int64 = 10_000
	maxRateMicros  int64 = 1_000_000_000_000
	quoteTTL             = time.Minute
)

type Repository struct {
	db     *pgxpool.Pool
	ledger *ledger.Repository
	risk   *risk.Repository
}

type ExchangeRate struct {
	ID              string     `json:"id"`
	BaseCurrency    string     `json:"base_currency"`
	QuoteCurrency   string     `json:"quote_currency"`
	RateMicros      int64      `json:"rate_micros"`
	SpreadBps       int        `json:"spread_bps"`
	Source          string     `json:"source"`
	Provider        string     `json:"provider"`
	SourceReference string     `json:"source_reference,omitempty"`
	SourceTimestamp *time.Time `json:"source_timestamp,omitempty"`
	ReceivedAt      time.Time  `json:"received_at"`
	ValidFrom       time.Time  `json:"valid_from"`
	StaleAfter      *time.Time `json:"stale_after,omitempty"`
	IsFallback      bool       `json:"is_fallback"`
	CreatedAt       time.Time  `json:"created_at"`
}

type Quote struct {
	ID                  string           `json:"id"`
	UserID              string           `json:"user_id"`
	FromAccountID       string           `json:"from_account_id"`
	ToAccountID         string           `json:"to_account_id"`
	FromCurrency        string           `json:"from_currency"`
	ToCurrency          string           `json:"to_currency"`
	FromAmountCents     int64            `json:"from_amount_cents"`
	ToAmountCents       int64            `json:"to_amount_cents"`
	RateMicros          int64            `json:"rate_micros"`
	MarketRateMicros    int64            `json:"market_rate_micros,omitempty"`
	SpreadBps           int              `json:"spread_bps"`
	FeeCents            int64            `json:"fee_cents"`
	ExchangeRateID      string           `json:"exchange_rate_id,omitempty"`
	RateSource          string           `json:"rate_source,omitempty"`
	RateProvider        string           `json:"rate_provider,omitempty"`
	RateSourceTimestamp *time.Time       `json:"rate_source_timestamp,omitempty"`
	RateStaleAfter      *time.Time       `json:"rate_stale_after,omitempty"`
	Disclosure          *QuoteDisclosure `json:"disclosure,omitempty"`
	Status              string           `json:"status"`
	IdempotencyKey      string           `json:"idempotency_key,omitempty"`
	ExpiresAt           time.Time        `json:"expires_at"`
	AcceptedAt          *time.Time       `json:"accepted_at,omitempty"`
	CreatedAt           time.Time        `json:"created_at"`
}

type Conversion struct {
	ID               string    `json:"id"`
	QuoteID          string    `json:"quote_id"`
	UserID           string    `json:"user_id"`
	FromAccountID    string    `json:"from_account_id"`
	ToAccountID      string    `json:"to_account_id"`
	FromWalletID     string    `json:"from_wallet_id,omitempty"`
	ToWalletID       string    `json:"to_wallet_id,omitempty"`
	FromCurrency     string    `json:"from_currency"`
	ToCurrency       string    `json:"to_currency"`
	FromAmountCents  int64     `json:"from_amount_cents"`
	ToAmountCents    int64     `json:"to_amount_cents"`
	RateMicros       int64     `json:"rate_micros"`
	MarketRateMicros int64     `json:"market_rate_micros,omitempty"`
	SpreadBps        int       `json:"spread_bps"`
	FeeCents         int64     `json:"fee_cents"`
	ExchangeRateID   string    `json:"exchange_rate_id,omitempty"`
	Status           string    `json:"status"`
	CreatedAt        time.Time `json:"created_at"`
}

type QuoteDisclosure struct {
	CustomerRateMicros     int64      `json:"customer_rate_micros"`
	MarketRateMicros       int64      `json:"market_rate_micros,omitempty"`
	SpreadBps              int        `json:"spread_bps"`
	FeeCents               int64      `json:"fee_cents"`
	FeeCurrency            string     `json:"fee_currency"`
	Source                 string     `json:"source,omitempty"`
	Provider               string     `json:"provider,omitempty"`
	SourceTimestamp        *time.Time `json:"source_timestamp,omitempty"`
	RateExpiresAt          time.Time  `json:"rate_expires_at"`
	EstimatedSpreadRevenue int64      `json:"estimated_spread_revenue_quote_cents,omitempty"`
	DisclosureText         string     `json:"disclosure_text"`
}

type RateChangeRequest struct {
	ID               string     `json:"id"`
	RequesterAdminID string     `json:"requester_admin_user_id"`
	RequesterName    string     `json:"requester_name,omitempty"`
	ReviewerAdminID  string     `json:"reviewer_admin_user_id,omitempty"`
	ReviewerName     string     `json:"reviewer_name,omitempty"`
	BaseCurrency     string     `json:"base_currency"`
	QuoteCurrency    string     `json:"quote_currency"`
	RateMicros       int64      `json:"rate_micros"`
	SpreadBps        int        `json:"spread_bps"`
	Source           string     `json:"source"`
	Provider         string     `json:"provider"`
	SourceReference  string     `json:"source_reference,omitempty"`
	SourceTimestamp  *time.Time `json:"source_timestamp,omitempty"`
	ValidFrom        time.Time  `json:"valid_from"`
	StaleAfter       *time.Time `json:"stale_after,omitempty"`
	IsFallback       bool       `json:"is_fallback"`
	Status           string     `json:"status"`
	DecisionNote     string     `json:"decision_note,omitempty"`
	ExchangeRateID   string     `json:"exchange_rate_id,omitempty"`
	IdempotencyKey   string     `json:"-"`
	DecidedAt        *time.Time `json:"decided_at,omitempty"`
	CreatedAt        time.Time  `json:"created_at"`
	UpdatedAt        time.Time  `json:"updated_at"`
}

type RateMonitoring struct {
	ExpiredQuotesMarked int64          `json:"expired_quotes_marked"`
	PendingQuotes       int64          `json:"pending_quotes"`
	ExpiredQuotes       int64          `json:"expired_quotes"`
	StaleRates          []ExchangeRate `json:"stale_rates"`
	FallbackRates       []ExchangeRate `json:"fallback_rates"`
	GeneratedAt         time.Time      `json:"generated_at"`
}

type TreasuryReport struct {
	InventoryPositions []TreasuryInventoryPosition `json:"inventory_positions"`
	PairPnL            []TreasuryPairPnL           `json:"pair_pnl"`
	GeneratedAt        time.Time                   `json:"generated_at"`
}

type TreasuryInventoryPosition struct {
	PositionID   string `json:"position_id"`
	Currency     string `json:"currency"`
	NetAmount    int64  `json:"net_amount_cents"`
	DebitAmount  int64  `json:"debit_amount_cents"`
	CreditAmount int64  `json:"credit_amount_cents"`
}

type TreasuryPairPnL struct {
	Pair                    string     `json:"pair"`
	FromCurrency            string     `json:"from_currency"`
	ToCurrency              string     `json:"to_currency"`
	ConversionCount         int64      `json:"conversion_count"`
	TotalFromAmountCents    int64      `json:"total_from_amount_cents"`
	TotalToAmountCents      int64      `json:"total_to_amount_cents"`
	EstimatedMidToCents     int64      `json:"estimated_mid_to_cents"`
	SpreadRevenueQuoteCents int64      `json:"spread_revenue_quote_cents"`
	TotalFeesFromCents      int64      `json:"total_fees_from_cents"`
	LastConversionAt        *time.Time `json:"last_conversion_at,omitempty"`
}

type FXAuditExport struct {
	RateRequests []RateChangeRequest `json:"rate_change_requests"`
	Rates        []ExchangeRate      `json:"exchange_rates"`
	Conversions  []Conversion        `json:"conversions"`
	GeneratedAt  time.Time           `json:"generated_at"`
}

type CreateRateParams struct {
	BaseCurrency     string
	QuoteCurrency    string
	RateMicros       int64
	SpreadBps        int
	Source           string
	Provider         string
	SourceReference  string
	SourceTimestamp  *time.Time
	ValidFrom        *time.Time
	StaleAfter       *time.Time
	IsFallback       bool
	RequesterAdminID string
	IdempotencyKey   string
}

type CreateQuoteParams struct {
	UserID          string
	FromAccountID   string
	ToAccountID     string
	FromAmountCents int64
	IdempotencyKey  string
}

type ConvertQuoteParams struct {
	UserID  string
	QuoteID string
}

type RateDecisionParams struct {
	RequestID    string
	ActorAdminID string
	Action       string
	DecisionNote string
}

type lockedAccount struct {
	ID           string
	UserID       string
	WalletID     string
	BalanceCents int64
	Currency     string
	Status       string
}

func NewRepository(db *pgxpool.Pool, ledgers ...*ledger.Repository) *Repository {
	ledgerRepo := ledger.NewRepository(db)
	if len(ledgers) > 0 && ledgers[0] != nil {
		ledgerRepo = ledgers[0]
	}
	return &Repository{db: db, ledger: ledgerRepo}
}

func (r *Repository) WithRisk(riskRepo *risk.Repository) *Repository {
	r.risk = riskRepo
	return r
}

func (r *Repository) ListRates(ctx context.Context, baseCurrency, quoteCurrency string, limit int) ([]ExchangeRate, error) {
	if limit <= 0 || limit > 100 {
		limit = 50
	}
	baseCurrency = domain.NormalizeCurrency(baseCurrency)
	quoteCurrency = domain.NormalizeCurrency(quoteCurrency)
	if baseCurrency != "" {
		if err := domain.ValidateCurrency(baseCurrency); err != nil {
			return nil, err
		}
	}
	if quoteCurrency != "" {
		if err := domain.ValidateCurrency(quoteCurrency); err != nil {
			return nil, err
		}
	}

	rows, err := r.db.Query(ctx, `
		SELECT id::text, base_currency, quote_currency, rate_micros, spread_bps, source,
			provider, source_reference, source_timestamp, received_at, valid_from, stale_after, is_fallback, created_at
		FROM exchange_rates
		WHERE ($1 = '' OR base_currency = $1)
			AND ($2 = '' OR quote_currency = $2)
		ORDER BY valid_from DESC, created_at DESC
		LIMIT $3
	`, baseCurrency, quoteCurrency, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	rates := []ExchangeRate{}
	for rows.Next() {
		rate, err := scanRate(rows)
		if err != nil {
			return nil, err
		}
		rates = append(rates, rate)
	}
	return rates, rows.Err()
}

func (r *Repository) CreateRate(ctx context.Context, params CreateRateParams) (ExchangeRate, error) {
	if err := normalizeRateParams(&params); err != nil {
		return ExchangeRate{}, err
	}

	rate, err := insertRate(ctx, r.db, params)
	if err != nil {
		if isForeignKeyViolation(err) {
			return ExchangeRate{}, fmt.Errorf("%w: currency is not enabled", domain.ErrValidation)
		}
		return ExchangeRate{}, err
	}
	return rate, nil
}

func (r *Repository) CreateRateChangeRequest(ctx context.Context, params CreateRateParams) (RateChangeRequest, error) {
	if err := normalizeRateParams(&params); err != nil {
		return RateChangeRequest{}, err
	}
	if err := domain.ValidateUUID("requester_admin_user_id", params.RequesterAdminID); err != nil {
		return RateChangeRequest{}, err
	}
	if len(params.IdempotencyKey) > 128 {
		return RateChangeRequest{}, fmt.Errorf("%w: Idempotency-Key must be 128 characters or fewer", domain.ErrValidation)
	}

	tx, err := r.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.Serializable})
	if err != nil {
		return RateChangeRequest{}, err
	}
	defer tx.Rollback(ctx)

	if params.IdempotencyKey != "" {
		existing, err := findRateChangeRequestByIdempotencyKey(ctx, tx, params.RequesterAdminID, params.IdempotencyKey)
		if err == nil {
			if !rateRequestMatchesParams(existing, params) {
				return RateChangeRequest{}, fmt.Errorf("%w: Idempotency-Key was already used for a different FX rate request", domain.ErrValidation)
			}
			return existing, tx.Commit(ctx)
		}
		if !errors.Is(err, pgx.ErrNoRows) {
			return RateChangeRequest{}, err
		}
	}
	if err := requireLiveAdmin(ctx, tx, params.RequesterAdminID); err != nil {
		return RateChangeRequest{}, err
	}

	var validFrom any
	if params.ValidFrom != nil {
		validFrom = *params.ValidFrom
	}
	var requestID string
	err = tx.QueryRow(ctx, `
		INSERT INTO fx_rate_change_requests (
			requester_admin_user_id, base_currency, quote_currency, rate_micros, spread_bps,
			source, provider, source_reference, source_timestamp, valid_from, stale_after,
			is_fallback, idempotency_key
		)
		VALUES (
			$1, $2, $3, $4, $5, $6, $7, $8, $9, COALESCE($10::timestamptz, now()), $11, $12, NULLIF($13, '')
		)
		RETURNING id::text
	`, params.RequesterAdminID, params.BaseCurrency, params.QuoteCurrency, params.RateMicros,
		params.SpreadBps, params.Source, params.Provider, params.SourceReference, params.SourceTimestamp,
		validFrom, params.StaleAfter, params.IsFallback, params.IdempotencyKey).Scan(&requestID)
	if err != nil {
		if isForeignKeyViolation(err) {
			return RateChangeRequest{}, fmt.Errorf("%w: currency is not enabled", domain.ErrValidation)
		}
		return RateChangeRequest{}, err
	}

	request, err := findRateChangeRequest(ctx, tx, requestID, false)
	if err != nil {
		return RateChangeRequest{}, err
	}
	if err := tx.Commit(ctx); err != nil {
		return RateChangeRequest{}, err
	}
	return request, nil
}

func (r *Repository) ListRateChangeRequests(ctx context.Context, status string, limit int) ([]RateChangeRequest, error) {
	status = strings.ToLower(strings.TrimSpace(status))
	if status != "" && status != "pending" && status != "approved" && status != "rejected" && status != "canceled" {
		return nil, fmt.Errorf("%w: invalid request status", domain.ErrValidation)
	}
	limit = normalizeLimit(limit)

	rows, err := r.db.Query(ctx, rateChangeRequestSelect+`
		WHERE (NULLIF($1, '') IS NULL OR req.status = $1)
		ORDER BY
			CASE WHEN req.status = 'pending' THEN 0 ELSE 1 END,
			req.created_at DESC
		LIMIT $2
	`, status, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	requests := []RateChangeRequest{}
	for rows.Next() {
		request, err := scanRateChangeRequest(rows)
		if err != nil {
			return nil, err
		}
		requests = append(requests, request)
	}
	return requests, rows.Err()
}

func (r *Repository) DecideRateChangeRequest(ctx context.Context, params RateDecisionParams) (RateChangeRequest, error) {
	if err := domain.ValidateUUID("request_id", params.RequestID); err != nil {
		return RateChangeRequest{}, err
	}
	if err := domain.ValidateUUID("actor_admin_user_id", params.ActorAdminID); err != nil {
		return RateChangeRequest{}, err
	}

	tx, err := r.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.Serializable})
	if err != nil {
		return RateChangeRequest{}, err
	}
	defer tx.Rollback(ctx)

	request, err := findRateChangeRequest(ctx, tx, params.RequestID, true)
	if err != nil {
		return RateChangeRequest{}, err
	}
	note, err := validateRateChangeDecision(request, params.ActorAdminID, params.Action, params.DecisionNote)
	if err != nil {
		return RateChangeRequest{}, err
	}
	if err := requireLiveAdmin(ctx, tx, params.ActorAdminID); err != nil {
		return RateChangeRequest{}, err
	}

	switch strings.ToLower(strings.TrimSpace(params.Action)) {
	case "approve":
		if err := requireLiveAdmin(ctx, tx, request.RequesterAdminID); err != nil {
			return RateChangeRequest{}, fmt.Errorf("%w: requester is no longer an active admin", err)
		}
		rate, err := insertRate(ctx, tx, CreateRateParams{
			BaseCurrency:    request.BaseCurrency,
			QuoteCurrency:   request.QuoteCurrency,
			RateMicros:      request.RateMicros,
			SpreadBps:       request.SpreadBps,
			Source:          request.Source,
			Provider:        request.Provider,
			SourceReference: request.SourceReference,
			SourceTimestamp: request.SourceTimestamp,
			ValidFrom:       &request.ValidFrom,
			StaleAfter:      request.StaleAfter,
			IsFallback:      request.IsFallback,
		})
		if err != nil {
			if isForeignKeyViolation(err) {
				return RateChangeRequest{}, fmt.Errorf("%w: currency is not enabled", domain.ErrValidation)
			}
			return RateChangeRequest{}, err
		}
		if _, err := tx.Exec(ctx, `
			UPDATE fx_rate_change_requests
			SET status = 'approved', reviewer_admin_user_id = $1, decision_note = $2,
				exchange_rate_id = $3, decided_at = now(), updated_at = now()
			WHERE id = $4
		`, params.ActorAdminID, note, rate.ID, request.ID); err != nil {
			return RateChangeRequest{}, err
		}
	case "reject":
		if _, err := tx.Exec(ctx, `
			UPDATE fx_rate_change_requests
			SET status = 'rejected', reviewer_admin_user_id = $1, decision_note = $2,
				decided_at = now(), updated_at = now()
			WHERE id = $3
		`, params.ActorAdminID, note, request.ID); err != nil {
			return RateChangeRequest{}, err
		}
	case "cancel":
		if _, err := tx.Exec(ctx, `
			UPDATE fx_rate_change_requests
			SET status = 'canceled', decision_note = $1, decided_at = now(), updated_at = now()
			WHERE id = $2
		`, note, request.ID); err != nil {
			return RateChangeRequest{}, err
		}
	default:
		return RateChangeRequest{}, fmt.Errorf("%w: action must be approve, reject or cancel", domain.ErrValidation)
	}

	updated, err := findRateChangeRequest(ctx, tx, request.ID, false)
	if err != nil {
		return RateChangeRequest{}, err
	}
	if err := tx.Commit(ctx); err != nil {
		return RateChangeRequest{}, err
	}
	return updated, nil
}

func (r *Repository) CreateQuote(ctx context.Context, params CreateQuoteParams) (Quote, error) {
	if err := validateQuoteParams(params); err != nil {
		return Quote{}, err
	}

	tx, err := r.db.BeginTx(ctx, pgx.TxOptions{})
	if err != nil {
		return Quote{}, err
	}
	defer tx.Rollback(ctx)

	if params.IdempotencyKey != "" {
		existing, err := findQuoteByIdempotencyKey(ctx, tx, params.UserID, params.IdempotencyKey)
		if err == nil {
			return existing, tx.Commit(ctx)
		}
		if !errors.Is(err, pgx.ErrNoRows) {
			return Quote{}, err
		}
	}

	accounts, err := lockAccounts(ctx, tx, params.FromAccountID, params.ToAccountID, "FOR SHARE")
	if err != nil {
		return Quote{}, err
	}
	from := accounts[params.FromAccountID]
	to := accounts[params.ToAccountID]
	if from == nil || to == nil {
		return Quote{}, domain.ErrNotFound
	}
	if from.UserID != params.UserID || to.UserID != params.UserID {
		return Quote{}, domain.ErrForbidden
	}
	if from.Status != "active" || to.Status != "active" {
		return Quote{}, fmt.Errorf("%w: both accounts must be active", domain.ErrValidation)
	}
	if from.Currency == to.Currency {
		return Quote{}, fmt.Errorf("%w: FX quote requires different currencies", domain.ErrValidation)
	}

	rate, err := r.latestRate(ctx, tx, from.Currency, to.Currency)
	if err != nil {
		return Quote{}, err
	}
	effectiveRate, err := applySpread(rate.RateMicros, rate.SpreadBps)
	if err != nil {
		return Quote{}, err
	}
	toAmount, err := r.convertAmountForCurrencies(ctx, tx, params.FromAmountCents, effectiveRate, from.Currency, to.Currency)
	if err != nil {
		return Quote{}, err
	}

	row := tx.QueryRow(ctx, `
		INSERT INTO fx_quotes (
			user_id, from_account_id, to_account_id, from_currency, to_currency,
			from_amount_cents, to_amount_cents, rate_micros, spread_bps, fee_cents,
			exchange_rate_id, market_rate_micros, rate_source, rate_provider, rate_source_timestamp,
			rate_stale_after, idempotency_key, expires_at
		)
		VALUES (
			$1, $2, $3, $4, $5, $6, $7, $8, $9, 0, NULLIF($10, '')::uuid, $11, $12, $13,
			$14, $15, NULLIF($16, ''), now() + make_interval(secs => $17)
		)
		RETURNING id::text, user_id::text, from_account_id::text, to_account_id::text,
			from_currency, to_currency, from_amount_cents, to_amount_cents, rate_micros, COALESCE(market_rate_micros, 0),
			spread_bps, fee_cents, COALESCE(exchange_rate_id::text, ''), rate_source, rate_provider,
			rate_source_timestamp, rate_stale_after, status, COALESCE(idempotency_key, ''), expires_at, accepted_at, created_at
	`, params.UserID, from.ID, to.ID, from.Currency, to.Currency, params.FromAmountCents, toAmount,
		effectiveRate, rate.SpreadBps, rate.ID, rate.RateMicros, rate.Source, rate.Provider,
		rate.SourceTimestamp, rate.StaleAfter, params.IdempotencyKey, int(quoteTTL.Seconds()))

	quote, err := scanQuote(row)
	if err != nil {
		if isUniqueViolation(err) && params.IdempotencyKey != "" {
			existing, findErr := findQuoteByIdempotencyKey(ctx, tx, params.UserID, params.IdempotencyKey)
			if findErr == nil {
				return existing, tx.Commit(ctx)
			}
		}
		return Quote{}, err
	}

	if err := tx.Commit(ctx); err != nil {
		return Quote{}, err
	}
	return quote, nil
}

func (r *Repository) ConvertQuote(ctx context.Context, params ConvertQuoteParams) (Conversion, error) {
	var conversion Conversion
	var err error
	for attempt := 0; attempt < 3; attempt++ {
		conversion, err = r.convertQuoteOnce(ctx, params)
		if !isRetryableTxError(err) {
			return conversion, err
		}
		time.Sleep(time.Duration(attempt+1) * 25 * time.Millisecond)
	}
	return conversion, err
}

func (r *Repository) ExpirePendingQuotes(ctx context.Context) (int64, error) {
	tag, err := r.db.Exec(ctx, `
		UPDATE fx_quotes
		SET status = 'expired'
		WHERE status = 'pending' AND expires_at <= now()
	`)
	if err != nil {
		return 0, err
	}
	return tag.RowsAffected(), nil
}

func (r *Repository) Monitoring(ctx context.Context, limit int) (RateMonitoring, error) {
	limit = normalizeLimit(limit)
	marked, err := r.ExpirePendingQuotes(ctx)
	if err != nil {
		return RateMonitoring{}, err
	}

	var pendingQuotes, expiredQuotes int64
	if err := r.db.QueryRow(ctx, `SELECT COUNT(*)::bigint FROM fx_quotes WHERE status = 'pending'`).Scan(&pendingQuotes); err != nil {
		return RateMonitoring{}, err
	}
	if err := r.db.QueryRow(ctx, `SELECT COUNT(*)::bigint FROM fx_quotes WHERE status = 'expired'`).Scan(&expiredQuotes); err != nil {
		return RateMonitoring{}, err
	}

	staleRates, err := r.listRatesByFreshness(ctx, false, limit)
	if err != nil {
		return RateMonitoring{}, err
	}
	fallbackRates, err := r.listRatesByFreshness(ctx, true, limit)
	if err != nil {
		return RateMonitoring{}, err
	}

	return RateMonitoring{
		ExpiredQuotesMarked: marked,
		PendingQuotes:       pendingQuotes,
		ExpiredQuotes:       expiredQuotes,
		StaleRates:          staleRates,
		FallbackRates:       fallbackRates,
		GeneratedAt:         time.Now().UTC(),
	}, nil
}

func (r *Repository) TreasuryReport(ctx context.Context, limit int) (TreasuryReport, error) {
	limit = normalizeLimit(limit)

	positionRows, err := r.db.Query(ctx, `
		SELECT la.reference_id, TRIM(la.currency)::text,
			COALESCE(SUM(CASE WHEN jl.direction = 'debit' THEN jl.amount_cents ELSE -jl.amount_cents END), 0)::bigint AS net_amount,
			COALESCE(SUM(CASE WHEN jl.direction = 'debit' THEN jl.amount_cents ELSE 0 END), 0)::bigint AS debit_amount,
			COALESCE(SUM(CASE WHEN jl.direction = 'credit' THEN jl.amount_cents ELSE 0 END), 0)::bigint AS credit_amount
		FROM ledger_accounts la
		LEFT JOIN ledger_journal_lines jl ON jl.ledger_account_id = la.id
		WHERE la.reference_type = 'fx_inventory'
		GROUP BY la.reference_id, TRIM(la.currency)::text
		ORDER BY la.reference_id, TRIM(la.currency)::text
		LIMIT $1
	`, limit)
	if err != nil {
		return TreasuryReport{}, err
	}
	defer positionRows.Close()

	positions := []TreasuryInventoryPosition{}
	for positionRows.Next() {
		var position TreasuryInventoryPosition
		if err := positionRows.Scan(&position.PositionID, &position.Currency, &position.NetAmount, &position.DebitAmount, &position.CreditAmount); err != nil {
			return TreasuryReport{}, err
		}
		positions = append(positions, position)
	}
	if err := positionRows.Err(); err != nil {
		return TreasuryReport{}, err
	}

	pnlRows, err := r.db.Query(ctx, `
		SELECT from_currency, to_currency,
			COUNT(*)::bigint,
			COALESCE(SUM(from_amount_cents), 0)::bigint,
			COALESCE(SUM(to_amount_cents), 0)::bigint,
			COALESCE(SUM((from_amount_cents * COALESCE(market_rate_micros, rate_micros)) / $1), 0)::bigint,
			COALESCE(SUM(((from_amount_cents * COALESCE(market_rate_micros, rate_micros)) / $1) - to_amount_cents), 0)::bigint,
			COALESCE(SUM(fee_cents), 0)::bigint,
			MAX(created_at)
		FROM fx_conversions
		WHERE status = 'completed'
		GROUP BY from_currency, to_currency
		ORDER BY MAX(created_at) DESC
		LIMIT $2
	`, rateScale, limit)
	if err != nil {
		return TreasuryReport{}, err
	}
	defer pnlRows.Close()

	pairPnL := []TreasuryPairPnL{}
	for pnlRows.Next() {
		var row TreasuryPairPnL
		var lastConversion sql.NullTime
		if err := pnlRows.Scan(
			&row.FromCurrency,
			&row.ToCurrency,
			&row.ConversionCount,
			&row.TotalFromAmountCents,
			&row.TotalToAmountCents,
			&row.EstimatedMidToCents,
			&row.SpreadRevenueQuoteCents,
			&row.TotalFeesFromCents,
			&lastConversion,
		); err != nil {
			return TreasuryReport{}, err
		}
		row.Pair = row.FromCurrency + "/" + row.ToCurrency
		if lastConversion.Valid {
			row.LastConversionAt = &lastConversion.Time
		}
		pairPnL = append(pairPnL, row)
	}
	if err := pnlRows.Err(); err != nil {
		return TreasuryReport{}, err
	}

	return TreasuryReport{InventoryPositions: positions, PairPnL: pairPnL, GeneratedAt: time.Now().UTC()}, nil
}

func (r *Repository) AuditExport(ctx context.Context, from, to *time.Time, limit int) (FXAuditExport, error) {
	limit = normalizeLimit(limit)
	rateRequests, err := r.listRateRequestsForExport(ctx, from, to, limit)
	if err != nil {
		return FXAuditExport{}, err
	}
	rates, err := r.listRatesForExport(ctx, from, to, limit)
	if err != nil {
		return FXAuditExport{}, err
	}
	conversions, err := r.listConversionsForExport(ctx, from, to, limit)
	if err != nil {
		return FXAuditExport{}, err
	}
	return FXAuditExport{
		RateRequests: rateRequests,
		Rates:        rates,
		Conversions:  conversions,
		GeneratedAt:  time.Now().UTC(),
	}, nil
}

func (r *Repository) convertQuoteOnce(ctx context.Context, params ConvertQuoteParams) (Conversion, error) {
	if err := domain.ValidateUUID("quote_id", params.QuoteID); err != nil {
		return Conversion{}, err
	}

	tx, err := r.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.Serializable})
	if err != nil {
		return Conversion{}, err
	}
	defer tx.Rollback(ctx)

	quote, err := lockQuote(ctx, tx, params.QuoteID)
	if err != nil {
		return Conversion{}, err
	}
	if quote.UserID != params.UserID {
		return Conversion{}, domain.ErrForbidden
	}
	if quote.Status == "accepted" {
		existing, err := findConversionByQuote(ctx, tx, quote.ID)
		if err == nil {
			return existing, tx.Commit(ctx)
		}
		if !errors.Is(err, pgx.ErrNoRows) {
			return Conversion{}, err
		}
		return Conversion{}, domain.ErrConflict
	}
	if quote.Status != "pending" {
		return Conversion{}, fmt.Errorf("%w: FX quote is not pending", domain.ErrValidation)
	}
	if isQuoteExpired(time.Now().UTC(), quote.ExpiresAt) {
		if _, err := tx.Exec(ctx, `
			UPDATE fx_quotes
			SET status = 'expired'
			WHERE id = $1 AND status = 'pending'
		`, quote.ID); err != nil {
			return Conversion{}, err
		}
		if err := tx.Commit(ctx); err != nil {
			return Conversion{}, err
		}
		return Conversion{}, fmt.Errorf("%w: FX quote has expired", domain.ErrValidation)
	}

	accounts, err := lockAccounts(ctx, tx, quote.FromAccountID, quote.ToAccountID, "FOR UPDATE")
	if err != nil {
		return Conversion{}, err
	}
	from := accounts[quote.FromAccountID]
	to := accounts[quote.ToAccountID]
	if from == nil || to == nil {
		return Conversion{}, domain.ErrNotFound
	}
	if err := validateConversionAccounts(quote, from, to); err != nil {
		return Conversion{}, err
	}

	totalDebit := quote.FromAmountCents + quote.FeeCents
	if totalDebit <= 0 || from.BalanceCents < totalDebit {
		return Conversion{}, domain.ErrInsufficientFunds
	}
	if r.risk != nil {
		evaluation, err := r.risk.EvaluateAndRecord(ctx, tx, risk.CheckRequest{
			UserID:      params.UserID,
			AccountID:   from.ID,
			Operation:   risk.OperationFXConversion,
			SourceType:  "fx_quote",
			SourceID:    quote.ID,
			Currency:    quote.FromCurrency,
			AmountCents: totalDebit,
			Metadata: map[string]any{
				"quote_id":        quote.ID,
				"from_account_id": from.ID,
				"to_account_id":   to.ID,
				"from_currency":   quote.FromCurrency,
				"to_currency":     quote.ToCurrency,
				"to_amount_cents": quote.ToAmountCents,
				"rate_micros":     quote.RateMicros,
			},
		})
		if err != nil {
			return Conversion{}, err
		}
		if risk.IsBlocked(evaluation) {
			if err := tx.Commit(ctx); err != nil {
				return Conversion{}, err
			}
			return Conversion{}, fmt.Errorf("%w: %s", domain.ErrLimitExceeded, evaluation.Reason)
		}
	}
	if err := updateAccountBalance(ctx, tx, from.ID, -totalDebit); err != nil {
		return Conversion{}, err
	}
	if err := updateAccountBalance(ctx, tx, to.ID, quote.ToAmountCents); err != nil {
		return Conversion{}, err
	}
	if from.WalletID != "" {
		if err := updateWalletBalance(ctx, tx, from.WalletID, from.Currency, -totalDebit); err != nil {
			return Conversion{}, err
		}
	}
	if to.WalletID != "" {
		if err := updateWalletBalance(ctx, tx, to.WalletID, to.Currency, quote.ToAmountCents); err != nil {
			return Conversion{}, err
		}
	}

	conversion, err := insertConversion(ctx, tx, quote, from.WalletID, to.WalletID)
	if err != nil {
		return Conversion{}, err
	}
	if _, err := tx.Exec(ctx, `
		UPDATE fx_quotes
		SET status = 'accepted', accepted_at = now()
		WHERE id = $1
	`, quote.ID); err != nil {
		return Conversion{}, err
	}
	if err := r.postLedger(ctx, tx, conversion, from, to); err != nil {
		return Conversion{}, err
	}

	if err := tx.Commit(ctx); err != nil {
		return Conversion{}, err
	}
	return conversion, nil
}

func (r *Repository) latestRate(ctx context.Context, tx pgx.Tx, baseCurrency, quoteCurrency string) (ExchangeRate, error) {
	row := tx.QueryRow(ctx, `
		SELECT id::text, base_currency, quote_currency, rate_micros, spread_bps, source,
			provider, source_reference, source_timestamp, received_at, valid_from, stale_after, is_fallback, created_at
		FROM exchange_rates
		WHERE base_currency = $1 AND quote_currency = $2 AND valid_from <= now()
			AND (stale_after IS NULL OR stale_after > now() OR is_fallback)
		ORDER BY
			CASE WHEN stale_after IS NULL OR stale_after > now() THEN 0 ELSE 1 END,
			is_fallback ASC,
			valid_from DESC,
			created_at DESC
		LIMIT 1
	`, baseCurrency, quoteCurrency)
	rate, err := scanRate(row)
	if err == nil {
		return rate, nil
	}
	if !errors.Is(err, pgx.ErrNoRows) {
		return ExchangeRate{}, err
	}

	row = tx.QueryRow(ctx, `
		SELECT id::text, base_currency, quote_currency, rate_micros, spread_bps, source,
			provider, source_reference, source_timestamp, received_at, valid_from, stale_after, is_fallback, created_at
		FROM exchange_rates
		WHERE base_currency = $1 AND quote_currency = $2 AND valid_from <= now()
			AND (stale_after IS NULL OR stale_after > now() OR is_fallback)
		ORDER BY
			CASE WHEN stale_after IS NULL OR stale_after > now() THEN 0 ELSE 1 END,
			is_fallback ASC,
			valid_from DESC,
			created_at DESC
		LIMIT 1
	`, quoteCurrency, baseCurrency)
	rate, err = scanRate(row)
	if errors.Is(err, pgx.ErrNoRows) {
		return ExchangeRate{}, fmt.Errorf("%w: no exchange rate available for currency pair", domain.ErrValidation)
	}
	if err != nil {
		return ExchangeRate{}, err
	}
	inverse, err := inverseRateMicros(rate.RateMicros)
	if err != nil {
		return ExchangeRate{}, err
	}
	rate.BaseCurrency = baseCurrency
	rate.QuoteCurrency = quoteCurrency
	rate.RateMicros = inverse
	rate.Source = "inverse:" + rate.Source
	return rate, nil
}

func (r *Repository) postLedger(ctx context.Context, tx pgx.Tx, conversion Conversion, from, to *lockedAccount) error {
	fromLedger, err := r.ledger.EnsureAccount(ctx, tx, ledger.AccountParams{
		OwnerUserID:   from.UserID,
		ReferenceType: "account",
		ReferenceID:   from.ID,
		Currency:      from.Currency,
		NormalBalance: "credit",
	})
	if err != nil {
		return err
	}
	toLedger, err := r.ledger.EnsureAccount(ctx, tx, ledger.AccountParams{
		OwnerUserID:   to.UserID,
		ReferenceType: "account",
		ReferenceID:   to.ID,
		Currency:      to.Currency,
		NormalBalance: "credit",
	})
	if err != nil {
		return err
	}
	positionID := conversion.FromCurrency + "_" + conversion.ToCurrency
	fromPositionLedger, err := r.ledger.EnsureAccount(ctx, tx, ledger.AccountParams{
		ReferenceType: "fx_inventory",
		ReferenceID:   positionID,
		Currency:      conversion.FromCurrency,
		NormalBalance: "credit",
	})
	if err != nil {
		return err
	}
	toPositionLedger, err := r.ledger.EnsureAccount(ctx, tx, ledger.AccountParams{
		ReferenceType: "fx_inventory",
		ReferenceID:   positionID,
		Currency:      conversion.ToCurrency,
		NormalBalance: "debit",
	})
	if err != nil {
		return err
	}

	_, err = r.ledger.Post(ctx, tx, ledger.PostParams{
		EventType:      "fx.conversion.completed",
		SourceType:     "fx_conversion",
		SourceID:       conversion.ID,
		IdempotencyKey: conversion.QuoteID,
		Description:    "Customer FX conversion",
		Metadata: map[string]any{
			"quote_id":           conversion.QuoteID,
			"from_account_id":    conversion.FromAccountID,
			"to_account_id":      conversion.ToAccountID,
			"from_currency":      conversion.FromCurrency,
			"to_currency":        conversion.ToCurrency,
			"rate_micros":        conversion.RateMicros,
			"market_rate_micros": conversion.MarketRateMicros,
			"spread_bps":         conversion.SpreadBps,
			"from_amount_cents":  conversion.FromAmountCents,
			"to_amount_cents":    conversion.ToAmountCents,
		},
		Lines: []ledger.LineParams{
			{
				LedgerAccountID: fromLedger.ID,
				Direction:       "debit",
				AmountCents:     conversion.FromAmountCents + conversion.FeeCents,
				Currency:        conversion.FromCurrency,
			},
			{
				LedgerAccountID: fromPositionLedger.ID,
				Direction:       "credit",
				AmountCents:     conversion.FromAmountCents + conversion.FeeCents,
				Currency:        conversion.FromCurrency,
			},
			{
				LedgerAccountID: toPositionLedger.ID,
				Direction:       "debit",
				AmountCents:     conversion.ToAmountCents,
				Currency:        conversion.ToCurrency,
			},
			{
				LedgerAccountID: toLedger.ID,
				Direction:       "credit",
				AmountCents:     conversion.ToAmountCents,
				Currency:        conversion.ToCurrency,
			},
		},
	})
	return err
}

func normalizeRateParams(params *CreateRateParams) error {
	params.BaseCurrency = domain.NormalizeCurrency(params.BaseCurrency)
	params.QuoteCurrency = domain.NormalizeCurrency(params.QuoteCurrency)
	params.Source = strings.TrimSpace(params.Source)
	if params.Source == "" {
		params.Source = "manual"
	}
	params.Provider = strings.TrimSpace(params.Provider)
	if params.Provider == "" {
		params.Provider = params.Source
	}
	params.SourceReference = strings.TrimSpace(params.SourceReference)
	if len(params.Source) > 80 {
		return fmt.Errorf("%w: source must be 80 characters or fewer", domain.ErrValidation)
	}
	if len(params.Provider) > 80 {
		return fmt.Errorf("%w: provider must be 80 characters or fewer", domain.ErrValidation)
	}
	if len(params.SourceReference) > 160 {
		return fmt.Errorf("%w: source_reference must be 160 characters or fewer", domain.ErrValidation)
	}
	if err := domain.ValidateCurrency(params.BaseCurrency); err != nil {
		return err
	}
	if err := domain.ValidateCurrency(params.QuoteCurrency); err != nil {
		return err
	}
	if params.BaseCurrency == params.QuoteCurrency {
		return fmt.Errorf("%w: base_currency and quote_currency must differ", domain.ErrValidation)
	}
	if params.RateMicros <= 0 || params.RateMicros > maxRateMicros {
		return fmt.Errorf("%w: rate_micros must be between 1 and %d", domain.ErrValidation, maxRateMicros)
	}
	if params.SpreadBps < 0 || params.SpreadBps > 2000 {
		return fmt.Errorf("%w: spread_bps must be between 0 and 2000", domain.ErrValidation)
	}
	validFrom := time.Now().UTC()
	if params.ValidFrom != nil {
		validFrom = params.ValidFrom.UTC()
		*params.ValidFrom = validFrom
	}
	if params.SourceTimestamp != nil {
		sourceTimestamp := params.SourceTimestamp.UTC()
		params.SourceTimestamp = &sourceTimestamp
	}
	if params.StaleAfter != nil {
		staleAfter := params.StaleAfter.UTC()
		params.StaleAfter = &staleAfter
		if !staleAfter.After(validFrom) {
			return fmt.Errorf("%w: stale_after must be after valid_from", domain.ErrValidation)
		}
	}
	return nil
}

func validateQuoteParams(params CreateQuoteParams) error {
	if err := domain.ValidateUUID("from_account_id", params.FromAccountID); err != nil {
		return err
	}
	if err := domain.ValidateUUID("to_account_id", params.ToAccountID); err != nil {
		return err
	}
	if params.FromAccountID == params.ToAccountID {
		return fmt.Errorf("%w: from_account_id and to_account_id must be different", domain.ErrValidation)
	}
	if err := domain.ValidateAmount(params.FromAmountCents); err != nil {
		return err
	}
	if len(params.IdempotencyKey) > 128 {
		return fmt.Errorf("%w: Idempotency-Key must be 128 characters or fewer", domain.ErrValidation)
	}
	return nil
}

func validateConversionAccounts(quote Quote, from, to *lockedAccount) error {
	if from.UserID != quote.UserID || to.UserID != quote.UserID {
		return domain.ErrForbidden
	}
	if from.Status != "active" || to.Status != "active" {
		return fmt.Errorf("%w: both accounts must be active", domain.ErrValidation)
	}
	if from.Currency != quote.FromCurrency || to.Currency != quote.ToCurrency {
		return fmt.Errorf("%w: account currencies no longer match the FX quote", domain.ErrValidation)
	}
	return nil
}

func applySpread(rateMicros int64, spreadBps int) (int64, error) {
	if rateMicros <= 0 || rateMicros > maxRateMicros {
		return 0, fmt.Errorf("%w: invalid rate_micros", domain.ErrValidation)
	}
	if spreadBps < 0 || spreadBps > 2000 {
		return 0, fmt.Errorf("%w: invalid spread_bps", domain.ErrValidation)
	}
	effective := rateMicros * (basisPointBase - int64(spreadBps)) / basisPointBase
	if effective <= 0 {
		return 0, fmt.Errorf("%w: effective FX rate is too small", domain.ErrValidation)
	}
	return effective, nil
}

func convertAmount(amountCents, effectiveRateMicros int64) (int64, error) {
	if effectiveRateMicros <= 0 {
		return 0, fmt.Errorf("%w: effective FX rate is too small", domain.ErrValidation)
	}
	return convertAmountWithPolicies(
		amountCents,
		effectiveRateMicros,
		domain.DefaultRoundingPolicy("EUR"),
		domain.DefaultRoundingPolicy("USD"),
	)
}

func (r *Repository) convertAmountForCurrencies(ctx context.Context, tx pgx.Tx, amountMinor, effectiveRateMicros int64, fromCurrency, toCurrency string) (int64, error) {
	if effectiveRateMicros <= 0 {
		return 0, fmt.Errorf("%w: effective FX rate is too small", domain.ErrValidation)
	}
	fromPolicy, err := currencyRoundingPolicy(ctx, tx, fromCurrency)
	if err != nil {
		return 0, err
	}
	toPolicy, err := currencyRoundingPolicy(ctx, tx, toCurrency)
	if err != nil {
		return 0, err
	}
	return convertAmountWithPolicies(amountMinor, effectiveRateMicros, fromPolicy, toPolicy)
}

func convertAmountWithPolicies(amountMinor, effectiveRateMicros int64, fromPolicy, toPolicy domain.RoundingPolicy) (int64, error) {
	return domain.ConvertMinorUnits(amountMinor, effectiveRateMicros, fromPolicy, toPolicy)
}

func currencyRoundingPolicy(ctx context.Context, tx pgx.Tx, currency string) (domain.RoundingPolicy, error) {
	currency = domain.NormalizeCurrency(currency)
	if err := domain.ValidateCurrency(currency); err != nil {
		return domain.RoundingPolicy{}, err
	}
	policy := domain.DefaultRoundingPolicy(currency)
	var mode string
	err := tx.QueryRow(ctx, `
		SELECT minor_unit, fx_rounding_mode, cash_rounding_increment_minor_units
		FROM currency_rounding_policies
		WHERE currency = $1
	`, currency).Scan(&policy.MinorUnit, &mode, &policy.CashRoundingIncrementMinorUnits)
	if errors.Is(err, pgx.ErrNoRows) {
		err = tx.QueryRow(ctx, `SELECT minor_unit FROM currencies WHERE code = $1`, currency).Scan(&policy.MinorUnit)
		mode = string(domain.RoundingModeHalfUp)
	}
	if err != nil {
		return domain.RoundingPolicy{}, err
	}
	policy.Mode = domain.RoundingMode(mode)
	return domain.NormalizeRoundingPolicy(policy)
}

func inverseRateMicros(rateMicros int64) (int64, error) {
	if rateMicros <= 0 {
		return 0, fmt.Errorf("%w: invalid inverse FX rate", domain.ErrValidation)
	}
	inverse := rateScale * rateScale / rateMicros
	if inverse <= 0 {
		return 0, fmt.Errorf("%w: inverse FX rate is too small", domain.ErrValidation)
	}
	return inverse, nil
}

func lockAccounts(ctx context.Context, tx pgx.Tx, fromID, toID, lockClause string) (map[string]*lockedAccount, error) {
	rows, err := tx.Query(ctx, fmt.Sprintf(`
		SELECT id::text, user_id::text, COALESCE(wallet_id::text, ''), balance_cents, currency, status
		FROM accounts
		WHERE id = $1 OR id = $2
		ORDER BY id
		%s
	`, lockClause), fromID, toID)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	accounts := map[string]*lockedAccount{}
	for rows.Next() {
		var account lockedAccount
		if err := rows.Scan(&account.ID, &account.UserID, &account.WalletID, &account.BalanceCents, &account.Currency, &account.Status); err != nil {
			return nil, err
		}
		accounts[account.ID] = &account
	}
	return accounts, rows.Err()
}

func lockQuote(ctx context.Context, tx pgx.Tx, quoteID string) (Quote, error) {
	row := tx.QueryRow(ctx, `
		SELECT id::text, user_id::text, from_account_id::text, to_account_id::text,
			from_currency, to_currency, from_amount_cents, to_amount_cents, rate_micros, COALESCE(market_rate_micros, 0),
			spread_bps, fee_cents, COALESCE(exchange_rate_id::text, ''), rate_source, rate_provider,
			rate_source_timestamp, rate_stale_after, status, COALESCE(idempotency_key, ''), expires_at, accepted_at, created_at
		FROM fx_quotes
		WHERE id = $1
		FOR UPDATE
	`, quoteID)
	quote, err := scanQuote(row)
	if errors.Is(err, pgx.ErrNoRows) {
		return Quote{}, domain.ErrNotFound
	}
	return quote, err
}

func findQuoteByIdempotencyKey(ctx context.Context, tx pgx.Tx, userID, key string) (Quote, error) {
	row := tx.QueryRow(ctx, `
		SELECT id::text, user_id::text, from_account_id::text, to_account_id::text,
			from_currency, to_currency, from_amount_cents, to_amount_cents, rate_micros, COALESCE(market_rate_micros, 0),
			spread_bps, fee_cents, COALESCE(exchange_rate_id::text, ''), rate_source, rate_provider,
			rate_source_timestamp, rate_stale_after, status, COALESCE(idempotency_key, ''), expires_at, accepted_at, created_at
		FROM fx_quotes
		WHERE user_id = $1 AND idempotency_key = $2
	`, userID, key)
	return scanQuote(row)
}

func insertConversion(ctx context.Context, tx pgx.Tx, quote Quote, fromWalletID, toWalletID string) (Conversion, error) {
	row := tx.QueryRow(ctx, `
		INSERT INTO fx_conversions (
			quote_id, user_id, from_account_id, to_account_id, from_wallet_id, to_wallet_id,
			from_currency, to_currency, from_amount_cents, to_amount_cents,
			rate_micros, market_rate_micros, spread_bps, fee_cents, exchange_rate_id
		)
		VALUES (
			$1, $2, $3, $4, NULLIF($5, '')::uuid, NULLIF($6, '')::uuid,
			$7, $8, $9, $10, $11, NULLIF($12, 0), $13, $14, NULLIF($15, '')::uuid
		)
		RETURNING id::text, quote_id::text, user_id::text, from_account_id::text, to_account_id::text,
			COALESCE(from_wallet_id::text, ''), COALESCE(to_wallet_id::text, ''),
			from_currency, to_currency, from_amount_cents, to_amount_cents, rate_micros,
			COALESCE(market_rate_micros, 0), spread_bps, fee_cents, COALESCE(exchange_rate_id::text, ''), status, created_at
	`, quote.ID, quote.UserID, quote.FromAccountID, quote.ToAccountID, fromWalletID, toWalletID,
		quote.FromCurrency, quote.ToCurrency, quote.FromAmountCents, quote.ToAmountCents,
		quote.RateMicros, quote.MarketRateMicros, quote.SpreadBps, quote.FeeCents, quote.ExchangeRateID)
	return scanConversion(row)
}

func findConversionByQuote(ctx context.Context, tx pgx.Tx, quoteID string) (Conversion, error) {
	row := tx.QueryRow(ctx, `
		SELECT id::text, quote_id::text, user_id::text, from_account_id::text, to_account_id::text,
			COALESCE(from_wallet_id::text, ''), COALESCE(to_wallet_id::text, ''),
			from_currency, to_currency, from_amount_cents, to_amount_cents, rate_micros,
			COALESCE(market_rate_micros, 0), spread_bps, fee_cents, COALESCE(exchange_rate_id::text, ''), status, created_at
		FROM fx_conversions
		WHERE quote_id = $1
	`, quoteID)
	return scanConversion(row)
}

func updateAccountBalance(ctx context.Context, tx pgx.Tx, accountID string, delta int64) error {
	tag, err := tx.Exec(ctx, `
		UPDATE accounts
		SET balance_cents = balance_cents + $1
		WHERE id = $2 AND balance_cents + $1 >= 0
	`, delta, accountID)
	if err != nil {
		return err
	}
	if tag.RowsAffected() == 0 {
		return domain.ErrInsufficientFunds
	}
	return nil
}

func updateWalletBalance(ctx context.Context, tx pgx.Tx, walletID, currency string, delta int64) error {
	if _, err := tx.Exec(ctx, `
		INSERT INTO wallet_balances (wallet_id, currency)
		VALUES ($1, $2)
		ON CONFLICT (wallet_id, currency) DO NOTHING
	`, walletID, currency); err != nil {
		return err
	}
	tag, err := tx.Exec(ctx, `
		UPDATE wallet_balances
		SET available_balance_cents = available_balance_cents + $1
		WHERE wallet_id = $2 AND currency = $3 AND available_balance_cents + $1 >= 0
	`, delta, walletID, currency)
	if err != nil {
		return err
	}
	if tag.RowsAffected() == 0 {
		return domain.ErrInsufficientFunds
	}
	return nil
}

func (r *Repository) listRatesByFreshness(ctx context.Context, fallbackOnly bool, limit int) ([]ExchangeRate, error) {
	condition := "stale_after <= now() AND is_fallback = false"
	if fallbackOnly {
		condition = "is_fallback = true"
	}
	rows, err := r.db.Query(ctx, fmt.Sprintf(`
		SELECT id::text, base_currency, quote_currency, rate_micros, spread_bps, source,
			provider, source_reference, source_timestamp, received_at, valid_from, stale_after, is_fallback, created_at
		FROM exchange_rates
		WHERE %s
		ORDER BY COALESCE(stale_after, valid_from) DESC, created_at DESC
		LIMIT $1
	`, condition), limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	rates := []ExchangeRate{}
	for rows.Next() {
		rate, err := scanRate(rows)
		if err != nil {
			return nil, err
		}
		rates = append(rates, rate)
	}
	return rates, rows.Err()
}

func (r *Repository) listRateRequestsForExport(ctx context.Context, from, to *time.Time, limit int) ([]RateChangeRequest, error) {
	fromArg, toArg := nullableTimeArgs(from, to)
	rows, err := r.db.Query(ctx, rateChangeRequestSelect+`
		WHERE ($1::timestamptz IS NULL OR req.created_at >= $1)
			AND ($2::timestamptz IS NULL OR req.created_at <= $2)
		ORDER BY req.created_at DESC
		LIMIT $3
	`, fromArg, toArg, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	requests := []RateChangeRequest{}
	for rows.Next() {
		request, err := scanRateChangeRequest(rows)
		if err != nil {
			return nil, err
		}
		requests = append(requests, request)
	}
	return requests, rows.Err()
}

func (r *Repository) listRatesForExport(ctx context.Context, from, to *time.Time, limit int) ([]ExchangeRate, error) {
	fromArg, toArg := nullableTimeArgs(from, to)
	rows, err := r.db.Query(ctx, `
		SELECT id::text, base_currency, quote_currency, rate_micros, spread_bps, source,
			provider, source_reference, source_timestamp, received_at, valid_from, stale_after, is_fallback, created_at
		FROM exchange_rates
		WHERE ($1::timestamptz IS NULL OR created_at >= $1)
			AND ($2::timestamptz IS NULL OR created_at <= $2)
		ORDER BY created_at DESC
		LIMIT $3
	`, fromArg, toArg, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	rates := []ExchangeRate{}
	for rows.Next() {
		rate, err := scanRate(rows)
		if err != nil {
			return nil, err
		}
		rates = append(rates, rate)
	}
	return rates, rows.Err()
}

func (r *Repository) listConversionsForExport(ctx context.Context, from, to *time.Time, limit int) ([]Conversion, error) {
	fromArg, toArg := nullableTimeArgs(from, to)
	rows, err := r.db.Query(ctx, `
		SELECT id::text, quote_id::text, user_id::text, from_account_id::text, to_account_id::text,
			COALESCE(from_wallet_id::text, ''), COALESCE(to_wallet_id::text, ''),
			from_currency, to_currency, from_amount_cents, to_amount_cents, rate_micros,
			COALESCE(market_rate_micros, 0), spread_bps, fee_cents, COALESCE(exchange_rate_id::text, ''), status, created_at
		FROM fx_conversions
		WHERE ($1::timestamptz IS NULL OR created_at >= $1)
			AND ($2::timestamptz IS NULL OR created_at <= $2)
		ORDER BY created_at DESC
		LIMIT $3
	`, fromArg, toArg, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	conversions := []Conversion{}
	for rows.Next() {
		conversion, err := scanConversion(rows)
		if err != nil {
			return nil, err
		}
		conversions = append(conversions, conversion)
	}
	return conversions, rows.Err()
}

func nullableTimeArgs(from, to *time.Time) (any, any) {
	var fromArg any
	if from != nil {
		fromArg = *from
	}
	var toArg any
	if to != nil {
		toArg = *to
	}
	return fromArg, toArg
}

func findRateChangeRequest(ctx context.Context, tx pgx.Tx, requestID string, forUpdate bool) (RateChangeRequest, error) {
	query := rateChangeRequestSelect + ` WHERE req.id = $1`
	if forUpdate {
		query += ` FOR UPDATE OF req`
	}
	request, err := scanRateChangeRequest(tx.QueryRow(ctx, query, requestID))
	if errors.Is(err, pgx.ErrNoRows) {
		return RateChangeRequest{}, domain.ErrNotFound
	}
	return request, err
}

func findRateChangeRequestByIdempotencyKey(ctx context.Context, tx pgx.Tx, requesterID, key string) (RateChangeRequest, error) {
	return scanRateChangeRequest(tx.QueryRow(ctx, rateChangeRequestSelect+`
		WHERE req.requester_admin_user_id = $1 AND req.idempotency_key = $2
	`, requesterID, key))
}

func validateRateChangeDecision(request RateChangeRequest, actorID, action, rawNote string) (string, error) {
	action = strings.ToLower(strings.TrimSpace(action))
	note := strings.TrimSpace(rawNote)
	if request.Status != "pending" {
		return "", fmt.Errorf("%w: only pending FX rate requests can be decided", domain.ErrValidation)
	}
	if action == "approve" || action == "reject" {
		if actorID == request.RequesterAdminID {
			return "", fmt.Errorf("%w: requester cannot approve or reject their own FX rate request", domain.ErrForbidden)
		}
	}
	if action == "cancel" && actorID != request.RequesterAdminID {
		return "", fmt.Errorf("%w: only the requester can cancel this FX rate request", domain.ErrForbidden)
	}
	if action != "approve" && action != "reject" && action != "cancel" {
		return "", fmt.Errorf("%w: action must be approve, reject or cancel", domain.ErrValidation)
	}
	if (action == "reject" || action == "cancel") && (len(note) < 8 || len(note) > 500) {
		return "", fmt.Errorf("%w: decision_note must be between 8 and 500 characters", domain.ErrValidation)
	}
	if action == "approve" && len(note) > 500 {
		return "", fmt.Errorf("%w: decision_note must be 500 characters or fewer", domain.ErrValidation)
	}
	return note, nil
}

func requireLiveAdmin(ctx context.Context, tx pgx.Tx, adminUserID string) error {
	var role string
	err := tx.QueryRow(ctx, `SELECT role FROM users WHERE id = $1 FOR SHARE`, adminUserID).Scan(&role)
	if errors.Is(err, pgx.ErrNoRows) {
		return domain.ErrUnauthorized
	}
	if err != nil {
		return err
	}
	if role != "admin" {
		return domain.ErrForbidden
	}
	return nil
}

func rateRequestMatchesParams(request RateChangeRequest, params CreateRateParams) bool {
	return request.BaseCurrency == params.BaseCurrency &&
		request.QuoteCurrency == params.QuoteCurrency &&
		request.RateMicros == params.RateMicros &&
		request.SpreadBps == params.SpreadBps &&
		request.Source == params.Source &&
		request.Provider == params.Provider &&
		request.SourceReference == params.SourceReference &&
		sameOptionalTime(request.SourceTimestamp, params.SourceTimestamp) &&
		(params.ValidFrom == nil || sameOptionalTime(&request.ValidFrom, params.ValidFrom)) &&
		sameOptionalTime(request.StaleAfter, params.StaleAfter) &&
		request.IsFallback == params.IsFallback
}

func sameOptionalTime(left, right *time.Time) bool {
	if left == nil || right == nil {
		return left == nil && right == nil
	}
	return left.Equal(*right)
}

func buildQuoteDisclosure(quote Quote) *QuoteDisclosure {
	disclosure := &QuoteDisclosure{
		CustomerRateMicros: quote.RateMicros,
		MarketRateMicros:   quote.MarketRateMicros,
		SpreadBps:          quote.SpreadBps,
		FeeCents:           quote.FeeCents,
		FeeCurrency:        quote.FromCurrency,
		Source:             quote.RateSource,
		Provider:           quote.RateProvider,
		SourceTimestamp:    quote.RateSourceTimestamp,
		RateExpiresAt:      quote.ExpiresAt,
		DisclosureText:     "The FX quote locks the shown customer rate until rate_expires_at. Fees are charged in the source currency.",
	}
	if quote.MarketRateMicros > 0 {
		disclosure.EstimatedSpreadRevenue = spreadRevenueQuoteCents(quote.FromAmountCents, quote.MarketRateMicros, quote.ToAmountCents)
	}
	return disclosure
}

func spreadRevenueQuoteCents(fromAmountCents, marketRateMicros, toAmountCents int64) int64 {
	if fromAmountCents <= 0 || marketRateMicros <= 0 || toAmountCents <= 0 || fromAmountCents > math.MaxInt64/marketRateMicros {
		return 0
	}
	midToAmount := fromAmountCents * marketRateMicros / rateScale
	if midToAmount <= toAmountCents {
		return 0
	}
	return midToAmount - toAmountCents
}

func isQuoteExpired(now, expiresAt time.Time) bool {
	return !now.Before(expiresAt)
}

func normalizeLimit(limit int) int {
	if limit <= 0 || limit > 100 {
		return 50
	}
	return limit
}

const rateChangeRequestSelect = `
	SELECT req.id::text, req.requester_admin_user_id::text, requester.full_name,
		COALESCE(req.reviewer_admin_user_id::text, ''), COALESCE(reviewer.full_name, ''),
		req.base_currency, req.quote_currency, req.rate_micros, req.spread_bps,
		req.source, req.provider, req.source_reference, req.source_timestamp,
		req.valid_from, req.stale_after, req.is_fallback, req.status,
		req.decision_note, COALESCE(req.exchange_rate_id::text, ''), COALESCE(req.idempotency_key, ''),
		req.decided_at, req.created_at, req.updated_at
	FROM fx_rate_change_requests req
	JOIN users requester ON requester.id = req.requester_admin_user_id
	LEFT JOIN users reviewer ON reviewer.id = req.reviewer_admin_user_id
`

type scanner interface {
	Scan(dest ...any) error
}

type queryRower interface {
	QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
}

func insertRate(ctx context.Context, q queryRower, params CreateRateParams) (ExchangeRate, error) {
	var validFrom any
	if params.ValidFrom != nil {
		validFrom = *params.ValidFrom
	}
	row := q.QueryRow(ctx, `
		INSERT INTO exchange_rates (
			base_currency, quote_currency, rate_micros, spread_bps, source, provider,
			source_reference, source_timestamp, valid_from, stale_after, is_fallback
		)
		VALUES (
			$1, $2, $3, $4, $5, $6, $7, $8, COALESCE($9::timestamptz, now()), $10, $11
		)
		RETURNING id::text, base_currency, quote_currency, rate_micros, spread_bps, source,
			provider, source_reference, source_timestamp, received_at, valid_from, stale_after, is_fallback, created_at
	`, params.BaseCurrency, params.QuoteCurrency, params.RateMicros, params.SpreadBps, params.Source,
		params.Provider, params.SourceReference, params.SourceTimestamp, validFrom, params.StaleAfter, params.IsFallback)
	return scanRate(row)
}

func scanRate(row scanner) (ExchangeRate, error) {
	var rate ExchangeRate
	var sourceTimestamp sql.NullTime
	var staleAfter sql.NullTime
	err := row.Scan(
		&rate.ID,
		&rate.BaseCurrency,
		&rate.QuoteCurrency,
		&rate.RateMicros,
		&rate.SpreadBps,
		&rate.Source,
		&rate.Provider,
		&rate.SourceReference,
		&sourceTimestamp,
		&rate.ReceivedAt,
		&rate.ValidFrom,
		&staleAfter,
		&rate.IsFallback,
		&rate.CreatedAt,
	)
	if sourceTimestamp.Valid {
		rate.SourceTimestamp = &sourceTimestamp.Time
	}
	if staleAfter.Valid {
		rate.StaleAfter = &staleAfter.Time
	}
	return rate, err
}

func scanQuote(row scanner) (Quote, error) {
	var quote Quote
	var acceptedAt sql.NullTime
	var rateSourceTimestamp sql.NullTime
	var rateStaleAfter sql.NullTime
	err := row.Scan(
		&quote.ID,
		&quote.UserID,
		&quote.FromAccountID,
		&quote.ToAccountID,
		&quote.FromCurrency,
		&quote.ToCurrency,
		&quote.FromAmountCents,
		&quote.ToAmountCents,
		&quote.RateMicros,
		&quote.MarketRateMicros,
		&quote.SpreadBps,
		&quote.FeeCents,
		&quote.ExchangeRateID,
		&quote.RateSource,
		&quote.RateProvider,
		&rateSourceTimestamp,
		&rateStaleAfter,
		&quote.Status,
		&quote.IdempotencyKey,
		&quote.ExpiresAt,
		&acceptedAt,
		&quote.CreatedAt,
	)
	if rateSourceTimestamp.Valid {
		quote.RateSourceTimestamp = &rateSourceTimestamp.Time
	}
	if rateStaleAfter.Valid {
		quote.RateStaleAfter = &rateStaleAfter.Time
	}
	if acceptedAt.Valid {
		quote.AcceptedAt = &acceptedAt.Time
	}
	quote.Disclosure = buildQuoteDisclosure(quote)
	return quote, err
}

func scanConversion(row scanner) (Conversion, error) {
	var conversion Conversion
	err := row.Scan(
		&conversion.ID,
		&conversion.QuoteID,
		&conversion.UserID,
		&conversion.FromAccountID,
		&conversion.ToAccountID,
		&conversion.FromWalletID,
		&conversion.ToWalletID,
		&conversion.FromCurrency,
		&conversion.ToCurrency,
		&conversion.FromAmountCents,
		&conversion.ToAmountCents,
		&conversion.RateMicros,
		&conversion.MarketRateMicros,
		&conversion.SpreadBps,
		&conversion.FeeCents,
		&conversion.ExchangeRateID,
		&conversion.Status,
		&conversion.CreatedAt,
	)
	return conversion, err
}

func scanRateChangeRequest(row scanner) (RateChangeRequest, error) {
	var request RateChangeRequest
	var sourceTimestamp sql.NullTime
	var staleAfter sql.NullTime
	var decidedAt sql.NullTime
	err := row.Scan(
		&request.ID,
		&request.RequesterAdminID,
		&request.RequesterName,
		&request.ReviewerAdminID,
		&request.ReviewerName,
		&request.BaseCurrency,
		&request.QuoteCurrency,
		&request.RateMicros,
		&request.SpreadBps,
		&request.Source,
		&request.Provider,
		&request.SourceReference,
		&sourceTimestamp,
		&request.ValidFrom,
		&staleAfter,
		&request.IsFallback,
		&request.Status,
		&request.DecisionNote,
		&request.ExchangeRateID,
		&request.IdempotencyKey,
		&decidedAt,
		&request.CreatedAt,
		&request.UpdatedAt,
	)
	if sourceTimestamp.Valid {
		request.SourceTimestamp = &sourceTimestamp.Time
	}
	if staleAfter.Valid {
		request.StaleAfter = &staleAfter.Time
	}
	if decidedAt.Valid {
		request.DecidedAt = &decidedAt.Time
	}
	return request, err
}

func isForeignKeyViolation(err error) bool {
	var pgErr *pgconn.PgError
	return errors.As(err, &pgErr) && pgErr.Code == "23503"
}

func isRetryableTxError(err error) bool {
	var pgErr *pgconn.PgError
	return errors.As(err, &pgErr) && (pgErr.Code == "40001" || pgErr.Code == "40P01")
}

func isUniqueViolation(err error) bool {
	var pgErr *pgconn.PgError
	return errors.As(err, &pgErr) && pgErr.Code == "23505"
}
