package reconciliation

import (
	"context"
	"database/sql"
	"encoding/json"
	"errors"
	"fmt"
	"strings"
	"time"

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

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

type Repository struct {
	db *pgxpool.Pool
}

type Run struct {
	ID                   string     `json:"id"`
	RunType              string     `json:"run_type"`
	Status               string     `json:"status"`
	StartedByAdminUserID string     `json:"started_by_admin_user_id,omitempty"`
	TotalChecks          int        `json:"total_checks"`
	MatchedCount         int        `json:"matched_count"`
	BreakCount           int        `json:"break_count"`
	ErrorMessage         string     `json:"error_message,omitempty"`
	StartedAt            time.Time  `json:"started_at"`
	CompletedAt          *time.Time `json:"completed_at,omitempty"`
}

type Break struct {
	ID                    string          `json:"id"`
	RunID                 string          `json:"run_id"`
	BreakType             string          `json:"break_type"`
	Severity              string          `json:"severity"`
	Status                string          `json:"status"`
	ReferenceType         string          `json:"reference_type"`
	ReferenceID           string          `json:"reference_id"`
	OwnerUserID           string          `json:"owner_user_id,omitempty"`
	Currency              string          `json:"currency"`
	ExpectedAmountCents   int64           `json:"expected_amount_cents"`
	ActualAmountCents     int64           `json:"actual_amount_cents"`
	DifferenceCents       int64           `json:"difference_cents"`
	Description           string          `json:"description"`
	Metadata              json.RawMessage `json:"metadata"`
	ResolvedByAdminUserID string          `json:"resolved_by_admin_user_id,omitempty"`
	ResolutionNote        string          `json:"resolution_note,omitempty"`
	ResolvedAt            *time.Time      `json:"resolved_at,omitempty"`
	CreatedAt             time.Time       `json:"created_at"`
	UpdatedAt             time.Time       `json:"updated_at"`
}

type ProviderBalanceSnapshot struct {
	ID                string          `json:"id"`
	Provider          string          `json:"provider"`
	ExternalAccountID string          `json:"external_account_id"`
	ReferenceType     string          `json:"reference_type"`
	ReferenceID       string          `json:"reference_id,omitempty"`
	Currency          string          `json:"currency"`
	BalanceCents      int64           `json:"balance_cents"`
	AsOf              time.Time       `json:"as_of"`
	Source            string          `json:"source"`
	RawPayload        json.RawMessage `json:"raw_payload"`
	CreatedAt         time.Time       `json:"created_at"`
}

type Metrics struct {
	OpenBreaks          int64      `json:"open_breaks"`
	InvestigatingBreaks int64      `json:"investigating_breaks"`
	CriticalOpenBreaks  int64      `json:"critical_open_breaks"`
	HighOpenBreaks      int64      `json:"high_open_breaks"`
	Resolved24h         int64      `json:"resolved_24h"`
	IgnoredBreaks       int64      `json:"ignored_breaks"`
	LastRunID           string     `json:"last_run_id,omitempty"`
	LastRunStatus       string     `json:"last_run_status,omitempty"`
	LastRunBreakCount   int        `json:"last_run_break_count,omitempty"`
	LastRunTotalChecks  int        `json:"last_run_total_checks,omitempty"`
	LastRunCompletedAt  *time.Time `json:"last_run_completed_at,omitempty"`
}

type Dashboard struct {
	Metrics Metrics `json:"metrics"`
	Runs    []Run   `json:"runs"`
	Breaks  []Break `json:"breaks"`
}

type CreateSnapshotParams struct {
	Provider          string
	ExternalAccountID string
	ReferenceType     string
	ReferenceID       string
	Currency          string
	BalanceCents      int64
	AsOf              *time.Time
	Source            string
	RawPayload        json.RawMessage
}

type BreakDecisionParams struct {
	BreakID        string
	AdminUserID    string
	Status         string
	ResolutionNote string
}

type candidateBreak struct {
	BreakType           string
	Severity            string
	ReferenceType       string
	ReferenceID         string
	OwnerUserID         string
	Currency            string
	ExpectedAmountCents int64
	ActualAmountCents   int64
	Description         string
	Metadata            map[string]any
}

type runStats struct {
	total   int
	matched int
	breaks  []candidateBreak
}

func NewRepository(db *pgxpool.Pool) *Repository {
	return &Repository{db: db}
}

func (r *Repository) Dashboard(ctx context.Context, status string, limit int) (Dashboard, error) {
	metrics, err := r.Metrics(ctx)
	if err != nil {
		return Dashboard{}, err
	}
	runs, err := r.ListRuns(ctx, limit)
	if err != nil {
		return Dashboard{}, err
	}
	breaks, err := r.ListBreaks(ctx, status, limit)
	if err != nil {
		return Dashboard{}, err
	}
	return Dashboard{Metrics: metrics, Runs: runs, Breaks: breaks}, nil
}

func (r *Repository) Run(ctx context.Context, adminUserID, runType string) (Run, error) {
	runType = normalizeRunType(runType)
	run, err := r.createRun(ctx, adminUserID, runType)
	if err != nil {
		return Run{}, err
	}

	tx, err := r.db.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.RepeatableRead})
	if err != nil {
		_ = r.failRun(ctx, run.ID, err)
		return Run{}, err
	}
	defer tx.Rollback(ctx)

	stats := runStats{}
	for _, check := range []func(context.Context, pgx.Tx, *runStats) error{
		r.checkJournalEntriesBalanced,
		r.checkAccountLedgerBalances,
		r.checkWalletAvailableBalances,
		r.checkWalletReservedBalances,
		r.checkProviderSnapshots,
	} {
		if err := check(ctx, tx, &stats); err != nil {
			_ = r.failRun(ctx, run.ID, err)
			return Run{}, err
		}
	}

	for _, item := range stats.breaks {
		if err := insertBreak(ctx, tx, run.ID, item); err != nil {
			_ = r.failRun(ctx, run.ID, err)
			return Run{}, err
		}
	}

	completed, err := completeRun(ctx, tx, run.ID, stats.total, stats.matched, len(stats.breaks))
	if err != nil {
		_ = r.failRun(ctx, run.ID, err)
		return Run{}, err
	}
	if err := tx.Commit(ctx); err != nil {
		_ = r.failRun(ctx, run.ID, err)
		return Run{}, err
	}
	return completed, nil
}

func (r *Repository) Metrics(ctx context.Context) (Metrics, error) {
	var metrics Metrics
	if err := r.db.QueryRow(ctx, `
		SELECT
			(SELECT COUNT(*) FROM reconciliation_breaks WHERE status = 'open')::bigint,
			(SELECT COUNT(*) FROM reconciliation_breaks WHERE status = 'investigating')::bigint,
			(SELECT COUNT(*) FROM reconciliation_breaks WHERE status = 'open' AND severity = 'critical')::bigint,
			(SELECT COUNT(*) FROM reconciliation_breaks WHERE status = 'open' AND severity = 'high')::bigint,
			(SELECT COUNT(*) FROM reconciliation_breaks WHERE status = 'resolved' AND resolved_at >= now() - interval '24 hours')::bigint,
			(SELECT COUNT(*) FROM reconciliation_breaks WHERE status = 'ignored')::bigint
	`).Scan(
		&metrics.OpenBreaks,
		&metrics.InvestigatingBreaks,
		&metrics.CriticalOpenBreaks,
		&metrics.HighOpenBreaks,
		&metrics.Resolved24h,
		&metrics.IgnoredBreaks,
	); err != nil {
		return Metrics{}, err
	}

	last, err := scanRun(r.db.QueryRow(ctx, `
		SELECT id::text, run_type, status, COALESCE(started_by_admin_user_id::text, ''),
			total_checks, matched_count, break_count, COALESCE(error_message, ''),
			started_at, completed_at
		FROM reconciliation_runs
		ORDER BY started_at DESC
		LIMIT 1
	`))
	if err == nil {
		metrics.LastRunID = last.ID
		metrics.LastRunStatus = last.Status
		metrics.LastRunBreakCount = last.BreakCount
		metrics.LastRunTotalChecks = last.TotalChecks
		metrics.LastRunCompletedAt = last.CompletedAt
	} else if !errors.Is(err, pgx.ErrNoRows) {
		return Metrics{}, err
	}
	return metrics, nil
}

func (r *Repository) ListRuns(ctx context.Context, limit int) ([]Run, error) {
	limit = normalizeLimit(limit)
	rows, err := r.db.Query(ctx, `
		SELECT id::text, run_type, status, COALESCE(started_by_admin_user_id::text, ''),
			total_checks, matched_count, break_count, COALESCE(error_message, ''),
			started_at, completed_at
		FROM reconciliation_runs
		ORDER BY started_at DESC
		LIMIT $1
	`, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	runs := []Run{}
	for rows.Next() {
		run, err := scanRun(rows)
		if err != nil {
			return nil, err
		}
		runs = append(runs, run)
	}
	return runs, rows.Err()
}

func (r *Repository) ListBreaks(ctx context.Context, status string, limit int) ([]Break, error) {
	limit = normalizeLimit(limit)
	status = strings.ToLower(strings.TrimSpace(status))
	if status == "" {
		status = "open"
	}
	if status != "all" && status != "open" && status != "investigating" && status != "resolved" && status != "ignored" {
		return nil, fmt.Errorf("%w: invalid reconciliation break status", domain.ErrValidation)
	}

	rows, err := r.db.Query(ctx, `
		SELECT id::text, run_id::text, break_type, severity, status, reference_type, reference_id,
			COALESCE(owner_user_id::text, ''), currency, expected_amount_cents, actual_amount_cents,
			difference_cents, description, metadata, COALESCE(resolved_by_admin_user_id::text, ''),
			COALESCE(resolution_note, ''), resolved_at, created_at, updated_at
		FROM reconciliation_breaks
		WHERE ($1 = 'all' OR status = $1)
		ORDER BY
			CASE severity WHEN 'critical' THEN 1 WHEN 'high' THEN 2 WHEN 'medium' THEN 3 ELSE 4 END,
			created_at DESC
		LIMIT $2
	`, status, limit)
	if err != nil {
		return nil, err
	}
	defer rows.Close()

	breaks := []Break{}
	for rows.Next() {
		item, err := scanBreak(rows)
		if err != nil {
			return nil, err
		}
		breaks = append(breaks, item)
	}
	return breaks, rows.Err()
}

func (r *Repository) CreateProviderBalanceSnapshot(ctx context.Context, params CreateSnapshotParams) (ProviderBalanceSnapshot, error) {
	if err := normalizeSnapshotParams(&params); err != nil {
		return ProviderBalanceSnapshot{}, err
	}
	raw := params.RawPayload
	if len(raw) == 0 {
		raw = []byte("{}")
	}
	if !json.Valid(raw) {
		return ProviderBalanceSnapshot{}, fmt.Errorf("%w: raw_payload must be valid JSON", domain.ErrValidation)
	}
	var asOf any
	if params.AsOf != nil {
		asOf = *params.AsOf
	}

	row := r.db.QueryRow(ctx, `
		INSERT INTO provider_balance_snapshots (
			provider, external_account_id, reference_type, reference_id, currency,
			balance_cents, as_of, source, raw_payload
		)
		VALUES ($1, $2, $3, NULLIF($4, ''), $5, $6, COALESCE($7::timestamptz, now()), $8, $9)
		RETURNING id::text, provider, external_account_id, reference_type, COALESCE(reference_id, ''),
			currency, balance_cents, as_of, source, raw_payload, created_at
	`, params.Provider, params.ExternalAccountID, params.ReferenceType, params.ReferenceID,
		params.Currency, params.BalanceCents, asOf, params.Source, raw)
	return scanSnapshot(row)
}

func (r *Repository) DecideBreak(ctx context.Context, params BreakDecisionParams) (Break, error) {
	if err := normalizeBreakDecision(&params); err != nil {
		return Break{}, err
	}

	var resolvedBy any
	var resolvedAt any
	if params.Status == "resolved" || params.Status == "ignored" {
		resolvedBy = params.AdminUserID
		resolvedAt = time.Now().UTC()
	}
	row := r.db.QueryRow(ctx, `
		UPDATE reconciliation_breaks
		SET status = $2,
			resolution_note = NULLIF($3, ''),
			resolved_by_admin_user_id = $4,
			resolved_at = $5::timestamptz
		WHERE id = $1
		RETURNING id::text, run_id::text, break_type, severity, status, reference_type, reference_id,
			COALESCE(owner_user_id::text, ''), currency, expected_amount_cents, actual_amount_cents,
			difference_cents, description, metadata, COALESCE(resolved_by_admin_user_id::text, ''),
			COALESCE(resolution_note, ''), resolved_at, created_at, updated_at
	`, params.BreakID, params.Status, params.ResolutionNote, resolvedBy, resolvedAt)

	item, err := scanBreak(row)
	if errors.Is(err, pgx.ErrNoRows) {
		return Break{}, domain.ErrNotFound
	}
	return item, err
}

func (r *Repository) createRun(ctx context.Context, adminUserID, runType string) (Run, error) {
	row := r.db.QueryRow(ctx, `
		INSERT INTO reconciliation_runs (run_type, started_by_admin_user_id)
		VALUES ($1, NULLIF($2, '')::uuid)
		RETURNING id::text, run_type, status, COALESCE(started_by_admin_user_id::text, ''),
			total_checks, matched_count, break_count, COALESCE(error_message, ''),
			started_at, completed_at
	`, runType, adminUserID)
	return scanRun(row)
}

func (r *Repository) failRun(ctx context.Context, runID string, cause error) error {
	_, err := r.db.Exec(ctx, `
		UPDATE reconciliation_runs
		SET status = 'failed', error_message = $2, completed_at = now()
		WHERE id = $1
	`, runID, cause.Error())
	return err
}

func (r *Repository) checkAccountLedgerBalances(ctx context.Context, tx pgx.Tx, stats *runStats) error {
	rows, err := tx.Query(ctx, `
		WITH ledger_balances AS (
			SELECT la.reference_id AS account_id,
				TRIM(la.currency)::text AS currency,
				COALESCE(SUM(CASE WHEN jl.direction = 'credit' THEN jl.amount_cents ELSE -jl.amount_cents END), 0)::bigint AS ledger_balance
			FROM ledger_accounts la
			LEFT JOIN ledger_journal_lines jl ON jl.ledger_account_id = la.id
			WHERE la.reference_type = 'account'
			GROUP BY la.reference_id, TRIM(la.currency)::text
		)
		SELECT a.id::text, a.user_id::text, TRIM(a.currency)::text,
			COALESCE(lb.ledger_balance, 0)::bigint AS expected_amount_cents,
			a.balance_cents::bigint AS actual_amount_cents
		FROM accounts a
		LEFT JOIN ledger_balances lb ON lb.account_id = a.id::text AND lb.currency = TRIM(a.currency)::text
		ORDER BY a.created_at, a.id
	`)
	if err != nil {
		return err
	}
	defer rows.Close()

	for rows.Next() {
		var accountID, userID, currency string
		var expected, actual int64
		if err := rows.Scan(&accountID, &userID, &currency, &expected, &actual); err != nil {
			return err
		}
		stats.total++
		if expected == actual {
			stats.matched++
			continue
		}
		stats.breaks = append(stats.breaks, candidateBreak{
			BreakType:           "account_ledger_mismatch",
			Severity:            "critical",
			ReferenceType:       "account",
			ReferenceID:         accountID,
			OwnerUserID:         userID,
			Currency:            currency,
			ExpectedAmountCents: expected,
			ActualAmountCents:   actual,
			Description:         "Account balance does not match account ledger balance",
			Metadata: map[string]any{
				"account_id": accountID,
				"source":     "accounts_vs_ledger_accounts",
			},
		})
	}
	return rows.Err()
}

func (r *Repository) checkWalletAvailableBalances(ctx context.Context, tx pgx.Tx, stats *runStats) error {
	rows, err := tx.Query(ctx, `
		WITH keys AS (
			SELECT wallet_id, TRIM(currency)::text AS currency FROM wallet_balances
			UNION
			SELECT wallet_id, TRIM(currency)::text FROM accounts WHERE wallet_id IS NOT NULL
		),
		account_totals AS (
			SELECT wallet_id, TRIM(currency)::text AS currency, COALESCE(SUM(balance_cents), 0)::bigint AS amount
			FROM accounts
			WHERE wallet_id IS NOT NULL
			GROUP BY wallet_id, TRIM(currency)::text
		)
		SELECT k.wallet_id::text, w.user_id::text, k.currency,
			COALESCE(a.amount, 0)::bigint AS expected_amount_cents,
			COALESCE(wb.available_balance_cents, 0)::bigint AS actual_amount_cents
		FROM keys k
		JOIN wallets w ON w.id = k.wallet_id
		LEFT JOIN wallet_balances wb ON wb.wallet_id = k.wallet_id AND TRIM(wb.currency)::text = k.currency
		LEFT JOIN account_totals a ON a.wallet_id = k.wallet_id AND a.currency = k.currency
		ORDER BY w.created_at, k.currency
	`)
	if err != nil {
		return err
	}
	defer rows.Close()

	for rows.Next() {
		var walletID, userID, currency string
		var expected, actual int64
		if err := rows.Scan(&walletID, &userID, &currency, &expected, &actual); err != nil {
			return err
		}
		stats.total++
		if expected == actual {
			stats.matched++
			continue
		}
		stats.breaks = append(stats.breaks, candidateBreak{
			BreakType:           "wallet_available_mismatch",
			Severity:            "high",
			ReferenceType:       "wallet",
			ReferenceID:         walletID,
			OwnerUserID:         userID,
			Currency:            currency,
			ExpectedAmountCents: expected,
			ActualAmountCents:   actual,
			Description:         "Wallet available balance does not match linked account balances",
			Metadata: map[string]any{
				"wallet_id": walletID,
				"source":    "wallet_balances_available_vs_accounts",
			},
		})
	}
	return rows.Err()
}

func (r *Repository) checkWalletReservedBalances(ctx context.Context, tx pgx.Tx, stats *runStats) error {
	rows, err := tx.Query(ctx, `
		WITH keys AS (
			SELECT wallet_id, TRIM(currency)::text AS currency FROM wallet_balances
			UNION
			SELECT a.wallet_id, TRIM(sg.currency)::text
			FROM savings_goals sg
			JOIN accounts a ON a.id = sg.account_id
			WHERE a.wallet_id IS NOT NULL AND sg.status <> 'closed'
		),
		savings_totals AS (
			SELECT a.wallet_id, TRIM(sg.currency)::text AS currency, COALESCE(SUM(sg.current_amount_cents), 0)::bigint AS amount
			FROM savings_goals sg
			JOIN accounts a ON a.id = sg.account_id
			WHERE a.wallet_id IS NOT NULL AND sg.status <> 'closed'
			GROUP BY a.wallet_id, TRIM(sg.currency)::text
		)
		SELECT k.wallet_id::text, w.user_id::text, k.currency,
			COALESCE(s.amount, 0)::bigint AS expected_amount_cents,
			COALESCE(wb.reserved_balance_cents, 0)::bigint AS actual_amount_cents
		FROM keys k
		JOIN wallets w ON w.id = k.wallet_id
		LEFT JOIN wallet_balances wb ON wb.wallet_id = k.wallet_id AND TRIM(wb.currency)::text = k.currency
		LEFT JOIN savings_totals s ON s.wallet_id = k.wallet_id AND s.currency = k.currency
		ORDER BY w.created_at, k.currency
	`)
	if err != nil {
		return err
	}
	defer rows.Close()

	for rows.Next() {
		var walletID, userID, currency string
		var expected, actual int64
		if err := rows.Scan(&walletID, &userID, &currency, &expected, &actual); err != nil {
			return err
		}
		stats.total++
		if expected == actual {
			stats.matched++
			continue
		}
		stats.breaks = append(stats.breaks, candidateBreak{
			BreakType:           "wallet_reserved_mismatch",
			Severity:            "high",
			ReferenceType:       "wallet",
			ReferenceID:         walletID,
			OwnerUserID:         userID,
			Currency:            currency,
			ExpectedAmountCents: expected,
			ActualAmountCents:   actual,
			Description:         "Wallet reserved balance does not match open savings goal balances",
			Metadata: map[string]any{
				"wallet_id": walletID,
				"source":    "wallet_balances_reserved_vs_savings_goals",
			},
		})
	}
	return rows.Err()
}

func (r *Repository) checkJournalEntriesBalanced(ctx context.Context, tx pgx.Tx, stats *runStats) error {
	rows, err := tx.Query(ctx, `
		SELECT je.id::text, TRIM(jl.currency)::text,
			COALESCE(SUM(CASE WHEN jl.direction = 'debit' THEN jl.amount_cents ELSE 0 END), 0)::bigint AS debit_total,
			COALESCE(SUM(CASE WHEN jl.direction = 'credit' THEN jl.amount_cents ELSE 0 END), 0)::bigint AS credit_total
		FROM ledger_journal_entries je
		LEFT JOIN ledger_journal_lines jl ON jl.journal_entry_id = je.id
		GROUP BY je.id, TRIM(jl.currency)::text
		ORDER BY je.created_at, je.id
	`)
	if err != nil {
		return err
	}
	defer rows.Close()

	for rows.Next() {
		var entryID string
		var currency sql.NullString
		var debit, credit int64
		if err := rows.Scan(&entryID, &currency, &debit, &credit); err != nil {
			return err
		}
		stats.total++
		currencyValue := strings.TrimSpace(currency.String)
		if currencyValue == "" {
			currencyValue = "EUR"
		}
		if debit == credit && debit > 0 {
			stats.matched++
			continue
		}
		stats.breaks = append(stats.breaks, candidateBreak{
			BreakType:           "journal_entry_unbalanced",
			Severity:            "critical",
			ReferenceType:       "ledger_journal_entry",
			ReferenceID:         entryID,
			Currency:            currencyValue,
			ExpectedAmountCents: debit,
			ActualAmountCents:   credit,
			Description:         "Ledger journal entry is not balanced for currency",
			Metadata: map[string]any{
				"journal_entry_id": entryID,
				"debit_total":      debit,
				"credit_total":     credit,
			},
		})
	}
	return rows.Err()
}

func (r *Repository) checkProviderSnapshots(ctx context.Context, tx pgx.Tx, stats *runStats) error {
	rows, err := tx.Query(ctx, `
		WITH ranked AS (
			SELECT pbs.*,
				ROW_NUMBER() OVER (
					PARTITION BY provider, external_account_id, currency
					ORDER BY as_of DESC, created_at DESC
				) AS rn
			FROM provider_balance_snapshots pbs
		),
		ledger_balances AS (
			SELECT la.reference_id AS account_id,
				TRIM(la.currency)::text AS currency,
				COALESCE(SUM(CASE WHEN jl.direction = 'credit' THEN jl.amount_cents ELSE -jl.amount_cents END), 0)::bigint AS ledger_balance
			FROM ledger_accounts la
			LEFT JOIN ledger_journal_lines jl ON jl.ledger_account_id = la.id
			WHERE la.reference_type = 'account'
			GROUP BY la.reference_id, TRIM(la.currency)::text
		)
		SELECT s.id::text, s.provider, s.external_account_id, TRIM(s.currency)::text,
			s.balance_cents::bigint AS expected_amount_cents,
			COALESCE(lb.ledger_balance, 0)::bigint AS actual_amount_cents,
			COALESCE(a.id::text, '') AS account_id,
			COALESCE(a.user_id::text, '') AS user_id
		FROM ranked s
		LEFT JOIN accounts a ON a.bank_provider = s.provider
			AND a.external_account_id = s.external_account_id
			AND TRIM(a.currency)::text = TRIM(s.currency)::text
		LEFT JOIN ledger_balances lb ON lb.account_id = a.id::text
			AND lb.currency = TRIM(s.currency)::text
		WHERE s.rn = 1
		ORDER BY s.provider, s.external_account_id, s.currency
	`)
	if err != nil {
		return err
	}
	defer rows.Close()

	for rows.Next() {
		var snapshotID, provider, externalAccountID, currency, accountID, userID string
		var expected, actual int64
		if err := rows.Scan(&snapshotID, &provider, &externalAccountID, &currency, &expected, &actual, &accountID, &userID); err != nil {
			return err
		}
		stats.total++
		if accountID == "" {
			stats.breaks = append(stats.breaks, candidateBreak{
				BreakType:           "provider_snapshot_orphan",
				Severity:            "medium",
				ReferenceType:       "provider_balance_snapshot",
				ReferenceID:         snapshotID,
				Currency:            currency,
				ExpectedAmountCents: expected,
				ActualAmountCents:   0,
				Description:         "Provider balance snapshot could not be matched to an internal account",
				Metadata: map[string]any{
					"provider":            provider,
					"external_account_id": externalAccountID,
				},
			})
			continue
		}
		if expected == actual {
			stats.matched++
			continue
		}
		stats.breaks = append(stats.breaks, candidateBreak{
			BreakType:           "provider_ledger_mismatch",
			Severity:            "critical",
			ReferenceType:       "account",
			ReferenceID:         accountID,
			OwnerUserID:         userID,
			Currency:            currency,
			ExpectedAmountCents: expected,
			ActualAmountCents:   actual,
			Description:         "Latest provider balance snapshot does not match internal ledger balance",
			Metadata: map[string]any{
				"provider":            provider,
				"external_account_id": externalAccountID,
				"snapshot_id":         snapshotID,
				"source":              "provider_balance_snapshot_vs_account_ledger",
			},
		})
	}
	return rows.Err()
}

func insertBreak(ctx context.Context, tx pgx.Tx, runID string, item candidateBreak) error {
	metadata := []byte("{}")
	if item.Metadata != nil {
		encoded, err := json.Marshal(item.Metadata)
		if err != nil {
			return err
		}
		metadata = encoded
	}
	_, err := tx.Exec(ctx, `
		INSERT INTO reconciliation_breaks (
			run_id, break_type, severity, reference_type, reference_id, owner_user_id, currency,
			expected_amount_cents, actual_amount_cents, difference_cents, description, metadata
		)
		VALUES ($1, $2, $3, $4, $5, NULLIF($6, '')::uuid, $7, $8, $9, $10, $11, $12)
	`, runID, item.BreakType, item.Severity, item.ReferenceType, item.ReferenceID, item.OwnerUserID,
		item.Currency, item.ExpectedAmountCents, item.ActualAmountCents,
		item.ActualAmountCents-item.ExpectedAmountCents, item.Description, metadata)
	return err
}

func completeRun(ctx context.Context, tx pgx.Tx, runID string, totalChecks, matchedCount, breakCount int) (Run, error) {
	row := tx.QueryRow(ctx, `
		UPDATE reconciliation_runs
		SET status = 'completed',
			total_checks = $2,
			matched_count = $3,
			break_count = $4,
			completed_at = now()
		WHERE id = $1
		RETURNING id::text, run_type, status, COALESCE(started_by_admin_user_id::text, ''),
			total_checks, matched_count, break_count, COALESCE(error_message, ''),
			started_at, completed_at
	`, runID, totalChecks, matchedCount, breakCount)
	return scanRun(row)
}

func normalizeRunType(runType string) string {
	runType = strings.ToLower(strings.TrimSpace(runType))
	if runType == "scheduled" || runType == "provider_snapshot" || runType == "daily" {
		return runType
	}
	return "manual"
}

func normalizeSnapshotParams(params *CreateSnapshotParams) error {
	params.Provider = strings.TrimSpace(params.Provider)
	params.ExternalAccountID = strings.TrimSpace(params.ExternalAccountID)
	params.ReferenceType = strings.ToLower(strings.TrimSpace(params.ReferenceType))
	params.ReferenceID = strings.TrimSpace(params.ReferenceID)
	params.Currency = domain.NormalizeCurrency(params.Currency)
	params.Source = strings.TrimSpace(params.Source)
	if params.Provider == "" || len(params.Provider) > 80 {
		return fmt.Errorf("%w: provider is required and must be 80 characters or fewer", domain.ErrValidation)
	}
	if params.ExternalAccountID == "" || len(params.ExternalAccountID) > 160 {
		return fmt.Errorf("%w: external_account_id is required and must be 160 characters or fewer", domain.ErrValidation)
	}
	if params.ReferenceType == "" {
		params.ReferenceType = "account"
	}
	if params.ReferenceType != "account" && params.ReferenceType != "wallet" && params.ReferenceType != "provider_account" {
		return fmt.Errorf("%w: invalid reference_type", domain.ErrValidation)
	}
	if err := domain.ValidateCurrency(params.Currency); err != nil {
		return err
	}
	if params.BalanceCents < 0 {
		return fmt.Errorf("%w: balance_cents must be zero or greater", domain.ErrValidation)
	}
	if params.Source == "" {
		params.Source = "manual"
	}
	if len(params.Source) > 80 {
		return fmt.Errorf("%w: source must be 80 characters or fewer", domain.ErrValidation)
	}
	return nil
}

func normalizeBreakDecision(params *BreakDecisionParams) error {
	if err := domain.ValidateUUID("id", params.BreakID); err != nil {
		return err
	}
	params.Status = strings.ToLower(strings.TrimSpace(params.Status))
	params.ResolutionNote = strings.TrimSpace(params.ResolutionNote)
	switch params.Status {
	case "open", "investigating":
		return nil
	case "resolved", "ignored":
		if len(params.ResolutionNote) < 8 || len(params.ResolutionNote) > 500 {
			return fmt.Errorf("%w: resolution_note must be between 8 and 500 characters", domain.ErrValidation)
		}
		return nil
	default:
		return fmt.Errorf("%w: invalid reconciliation break status", domain.ErrValidation)
	}
}

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

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

func scanRun(row scanner) (Run, error) {
	var run Run
	var completedAt sql.NullTime
	err := row.Scan(
		&run.ID,
		&run.RunType,
		&run.Status,
		&run.StartedByAdminUserID,
		&run.TotalChecks,
		&run.MatchedCount,
		&run.BreakCount,
		&run.ErrorMessage,
		&run.StartedAt,
		&completedAt,
	)
	if completedAt.Valid {
		run.CompletedAt = &completedAt.Time
	}
	return run, err
}

func scanBreak(row scanner) (Break, error) {
	var item Break
	var resolvedAt sql.NullTime
	err := row.Scan(
		&item.ID,
		&item.RunID,
		&item.BreakType,
		&item.Severity,
		&item.Status,
		&item.ReferenceType,
		&item.ReferenceID,
		&item.OwnerUserID,
		&item.Currency,
		&item.ExpectedAmountCents,
		&item.ActualAmountCents,
		&item.DifferenceCents,
		&item.Description,
		&item.Metadata,
		&item.ResolvedByAdminUserID,
		&item.ResolutionNote,
		&resolvedAt,
		&item.CreatedAt,
		&item.UpdatedAt,
	)
	if resolvedAt.Valid {
		item.ResolvedAt = &resolvedAt.Time
	}
	return item, err
}

func scanSnapshot(row scanner) (ProviderBalanceSnapshot, error) {
	var snapshot ProviderBalanceSnapshot
	err := row.Scan(
		&snapshot.ID,
		&snapshot.Provider,
		&snapshot.ExternalAccountID,
		&snapshot.ReferenceType,
		&snapshot.ReferenceID,
		&snapshot.Currency,
		&snapshot.BalanceCents,
		&snapshot.AsOf,
		&snapshot.Source,
		&snapshot.RawPayload,
		&snapshot.CreatedAt,
	)
	return snapshot, err
}
