import { createScheduledDbService, type DoDbClient } from '@/server/services/db.js';
/**
 * Cloudflare Queue consumer for `multideal-settlements-preview`.
 *
 * Processes settlement release batches. Each message carries a list of
 * vendorPayoutRelease IDs to claim and pay out via Stripe.
 *
 * Flow:
 *   1. Load rows with status='enqueued' matching releaseIds
 *      — 0 rows → recover 'releasing' rows with the original idempotency key
 *   2. Mark rows as 'releasing' + increment claimAttempts
 *   3. Check Stripe balance for the vendor Connect account
 *   4a. Insufficient funds → mark 'held' + lastError; retry via queue
 *   4b. Max retries reached → re-enqueue (cron backstop will re-sweep)
 *   5. Persist a release-ID-derived idempotency key, then create Stripe payout
 *   6. Mark rows 'released' with payoutId + releasedAt
 *   7. Permanent payout error → hold; transient error → same-key recovery on redelivery
 *
 * Message format: { kind: 'settlement.release', batchKey: string, releaseIds: string[] }
 */

import Stripe from 'stripe';
import { z } from 'zod';
import type { MultidealEnv } from '../lib/env.js';
import {
  findEnqueuedReleases,
  findReleasingReleases,
  claimReleases,
  persistPayoutIdempotencyKey,
  holdReleases,
  settleReleases,
  releaseReleases,
} from '@/server/db/queries/settlement-queue.js';

// ─── Types ────────────────────────────────────────────────────────────────────

const SettlementReleaseMessageSchema = z.object({
  kind: z.literal('settlement.release'),
  batchKey: z.string().min(1).max(200),
  releaseIds: z.array(z.uuid()).min(1).max(500),
});

export type SettlementReleaseMessage = z.infer<typeof SettlementReleaseMessageSchema>;

/** Minimal Stripe surface used by this consumer — keeps it dependency-free from the stripe package. */
export interface StripePayoutClient {
  balance: {
    retrieve(
      params?: Record<string, unknown>,
      options?: { stripeAccount?: string },
    ): Promise<{ available: Array<{ currency: string; amount: number }> }>;
  };
  payouts: {
    create(
      params: { amount: number; currency: string; metadata?: Record<string, string> },
      options?: { stripeAccount?: string; idempotencyKey?: string },
    ): Promise<{ id: string }>;
  };
}

// ─── Constants ────────────────────────────────────────────────────────────────

const MAX_PENDING_RETRIES = 10;
const STRIPE_IDEMPOTENCY_RETENTION_MS = 24 * 60 * 60 * 1_000;

class SettlementBatchInvariantError extends Error {
  readonly statusCode = 400;
}

class SettlementPayoutStateConflictError extends Error {}

async function buildPayoutIdempotencyKey(ids: string[]): Promise<string> {
  const canonicalIds = [...ids].sort().join(',');
  const digest = await crypto.subtle.digest('SHA-256', new TextEncoder().encode(canonicalIds));
  const hex = Array.from(new Uint8Array(digest), (byte) => byte.toString(16).padStart(2, '0')).join(
    '',
  );
  return `settlement-payout-${hex}`;
}

type ReleasePayoutKeyState = {
  payoutIdempotencyKey: string | null;
  payoutIdempotencyKeyCreatedAt: Date | null;
};

async function validatePayoutKeyState(
  db: DoDbClient,
  ids: string[],
  releases: ReleasePayoutKeyState[],
  requirePersisted: boolean,
): Promise<string | undefined> {
  const persistedKeys = new Set(
    releases
      .map((release) => release.payoutIdempotencyKey)
      .filter((key): key is string => key != null),
  );
  const hasUnkeyedRelease = releases.some((release) => release.payoutIdempotencyKey == null);
  const hasTimestampWithoutKey = releases.some(
    (release) =>
      release.payoutIdempotencyKey == null && release.payoutIdempotencyKeyCreatedAt != null,
  );
  if (
    persistedKeys.size > 1 ||
    (persistedKeys.size === 1 && hasUnkeyedRelease) ||
    hasTimestampWithoutKey
  ) {
    await holdReleases(db, ids, 'mixed_payout_idempotency_keys');
    throw new SettlementBatchInvariantError('settlement batch has mixed payout keys');
  }

  const idempotencyKey = persistedKeys.values().next().value as string | undefined;
  if (idempotencyKey === undefined) {
    if (requirePersisted) {
      await holdReleases(db, ids, 'missing_payout_idempotency_key');
      throw new SettlementBatchInvariantError('releasing settlement batch has no payout key');
    }
    return undefined;
  }

  const keyCreatedAtValues = releases.map((release) => release.payoutIdempotencyKeyCreatedAt);
  const keyCreatedAtEpochs = new Set(
    keyCreatedAtValues
      .filter((value): value is Date => value != null)
      .map((value) => value.getTime()),
  );
  if (keyCreatedAtValues.some((value) => value == null) || keyCreatedAtEpochs.size !== 1) {
    await holdReleases(db, ids, 'mixed_payout_idempotency_keys');
    throw new SettlementBatchInvariantError(
      'settlement batch has incomplete payout key timestamps',
    );
  }
  if (Date.now() - keyCreatedAtValues[0]!.getTime() >= STRIPE_IDEMPOTENCY_RETENTION_MS) {
    await holdReleases(db, ids, 'payout_idempotency_key_expired');
    throw new SettlementBatchInvariantError(
      'persisted payout key is outside Stripe idempotency retention',
    );
  }
  return idempotencyKey;
}

/**
 * Classify a thrown Stripe/runtime error as permanent (no point retrying) vs
 * transient (worth a queue retry). Permanent = HTTP 4xx other than 429
 * (invalid account, bad request, auth) — these fail identically on every retry,
 * so retrying only burns Queue read ops until the message dead-letters. Transient
 * = 429 rate-limit, 5xx, and network errors with no statusCode. Unknown → treated
 * as transient (safe: bounded by max_retries then DLQ, vs silently dropping work).
 */
function isPermanentError(err: unknown): boolean {
  const e = err as { statusCode?: number; status?: number };
  const code = typeof e?.statusCode === 'number' ? e.statusCode : e?.status;
  if (typeof code !== 'number') return false;
  if (code === 429) return false;
  return code >= 400 && code < 500;
}

// ─── Consumer ─────────────────────────────────────────────────────────────────

/**
 * Batch handler — called by CF Queue consumer binding.
 * Accepts an optional `stripeClient` for testing (injected mock).
 */
export async function handleSettlementBatch(
  batch: MessageBatch<unknown>,
  env: Pick<MultidealEnv, 'DATABASE_URL' | 'STRIPE_SECRET_KEY' | 'PAYMENT_PROVIDER'>,
  stripeClient?: StripePayoutClient,
): Promise<void> {
  const validMessages: Array<{
    message: (typeof batch.messages)[number];
    body: SettlementReleaseMessage;
  }> = [];

  for (const message of batch.messages) {
    const parsed = SettlementReleaseMessageSchema.safeParse(message.body);
    if (!parsed.success) {
      console.error(
        JSON.stringify({
          event: 'settlement_invalid_message',
          messageId: message.id,
          issue: parsed.error.issues[0]?.message ?? 'invalid payload',
          issueCount: parsed.error.issues.length,
        }),
      );
      message.ack();
      continue;
    }
    validMessages.push({ message, body: parsed.data });
  }

  if (validMessages.length === 0) return;

  // Permanent-config guard: a missing STRIPE_SECRET_KEY cannot be fixed by retrying.
  // Constructing `new Stripe('')` would auth-fail on every balance.retrieve, so each
  // message would burn max_retries+1 reads → DLQ — pure waste with no progress.
  // Mirror the outbox consumer's `strict=false` house convention: warn + ack, leave the
  // rows untouched (still 'enqueued'), and let the */30 settlement cron re-sweep once a
  // key is configured. (Preview is a DEPLOYED worker → must run real `sk_test_…`, never
  // PAYMENT_PROVIDER=mock; an absent secret here is a misconfig to fix, not a code path.)
  if (!stripeClient && !env.STRIPE_SECRET_KEY) {
    for (const { message, body } of validMessages) {
      console.warn(
        JSON.stringify({
          event: 'settlement_payouts_unconfigured',
          batchKey: body.batchKey,
          releaseIds: body.releaseIds,
          note: 'STRIPE_SECRET_KEY not set on this worker — acking without payout; set the secret to enable settlement releases.',
        }),
      );
      message.ack();
    }
    return;
  }

  const db = createScheduledDbService(env);
  const stripe: StripePayoutClient =
    stripeClient ??
    new Stripe(env.STRIPE_SECRET_KEY ?? '', {
      apiVersion: '2026-04-22.dahlia',
    });

  for (const { message, body } of validMessages) {
    try {
      await processMessage(body, db, stripe);
      message.ack();
    } catch (err) {
      const errorMsg = err instanceof Error ? err.message : String(err);
      const permanent = isPermanentError(err);
      console.error(
        JSON.stringify({
          event: 'settlement_consumer_error',
          batchKey: body.batchKey,
          error: errorMsg,
          permanent,
        }),
      );
      // Permanent errors (invalid account, 4xx) can't be fixed by retrying — ack to
      // stop a retry->DLQ->sweep->re-enqueue storm. Claimed rows stay visible as held;
      // pre-claim invariant failures leave rows enqueued. Only transient failures
      // (429/5xx/network) get a queue retry.
      if (permanent) {
        message.ack();
      } else {
        message.retry();
      }
    }
  }
}

// ─── Core processor ───────────────────────────────────────────────────────────

async function processMessage(
  msg: SettlementReleaseMessage,
  db: DoDbClient,
  stripe: StripePayoutClient,
): Promise<void> {
  // 1. Load enqueued releases
  const candidates = await findEnqueuedReleases(db, msg.releaseIds);

  if (candidates.length === 0) {
    await recoverStrandedReleasing(msg, db, stripe);
    return;
  }

  if (candidates.some((release) => release.payoutBatchKey !== msg.batchKey)) {
    throw new SettlementBatchInvariantError(
      'enqueued settlement rows do not match message batch key',
    );
  }

  const candidateIds = candidates.map((release) => release.id);

  // 2. Mark releasing + bump claimAttempts
  const claimed = await claimReleases(db, candidateIds);
  const claimedById = new Map(claimed.map((release) => [release.id, release]));
  const releases = candidates
    .filter((release) => claimedById.has(release.id))
    .map((release) => ({
      ...release,
      payoutIdempotencyKey: claimedById.get(release.id)!.payoutIdempotencyKey,
      payoutIdempotencyKeyCreatedAt: claimedById.get(release.id)!.payoutIdempotencyKeyCreatedAt,
    }));

  if (releases.length === 0) {
    return;
  }

  const ids = releases.map((r) => r.id);
  const vendorAcctIds = new Set(releases.map((release) => release.vendorAcctId));
  if (vendorAcctIds.size !== 1) {
    await holdReleases(db, ids, 'mixed_vendor_batch');
    throw new SettlementBatchInvariantError('claimed settlement batch spans multiple vendors');
  }

  const vendorAcctId = releases[0]!.vendorAcctId;
  const totalAgorot = releases.reduce((s, r) => s + r.netAmountAgorot, 0);
  let idempotencyKey = await validatePayoutKeyState(db, ids, releases, false);

  // 3. Check Stripe balance
  let ilsAvailable: number;
  try {
    const bal = await stripe.balance.retrieve({}, { stripeAccount: vendorAcctId });
    ilsAvailable = bal.available.find((b) => b.currency === 'ils')?.amount ?? 0;
  } catch (err) {
    // Balance check failed — treat as held, let queue retry
    const errorMsg = err instanceof Error ? err.message : String(err);
    await holdReleases(db, ids, `balance_check_failed: ${errorMsg}`);
    throw err;
  }

  // 4. Insufficient funds path
  if (ilsAvailable < totalAgorot) {
    // claimAttempts was just incremented; compute new value for threshold check
    const maxAttemptsSoFar = Math.max(...releases.map((r) => (r.claimAttempts ?? 0) + 1));
    const dlqOnNextSweep = maxAttemptsSoFar >= MAX_PENDING_RETRIES;

    await settleReleases(
      db,
      ids,
      dlqOnNextSweep ? 'enqueued' : 'held',
      `funds_pending: available=${ilsAvailable} needed=${totalAgorot}`,
    );

    if (dlqOnNextSweep) {
      // At max retries — ack the queue message and let the cron backstop re-sweep
      console.warn(
        JSON.stringify({
          event: 'settlement_max_retries_reached',
          batchKey: msg.batchKey,
          vendorAcctId,
          totalAgorot,
          ilsAvailable,
        }),
      );
      return; // caller acks
    }

    // Below max retries — throw so caller retries the queue message
    throw new Error(
      `funds_pending: available=${ilsAvailable} needed=${totalAgorot} vendor=${vendorAcctId}`,
    );
  }

  // 5. Create Stripe payout
  let payout: { id: string };
  try {
    if (idempotencyKey === undefined) {
      idempotencyKey = await buildPayoutIdempotencyKey(ids);
      const persisted = await persistPayoutIdempotencyKey(db, ids, idempotencyKey);
      if (persisted.length !== ids.length) {
        await holdReleases(db, ids, 'payout_idempotency_key_write_conflict');
        throw new SettlementPayoutStateConflictError(
          `payout key state conflict: ${ids.length} releases claimed but ${persisted.length} persisted the payout key`,
        );
      }
    }
    payout = await stripe.payouts.create(
      {
        amount: totalAgorot,
        currency: 'ils',
        metadata: {
          batchKey: msg.batchKey,
          releaseCount: String(ids.length),
        },
      },
      {
        stripeAccount: vendorAcctId,
        idempotencyKey,
      },
    );
  } catch (err) {
    const errorMsg = err instanceof Error ? err.message : String(err);
    if (isPermanentError(err)) {
      await holdReleases(db, ids, errorMsg);
    }
    throw err;
  }

  // 6. Mark released
  const transitioned = await releaseReleases(db, ids, payout.id);
  if (transitioned.length !== ids.length) {
    throw new SettlementPayoutStateConflictError(
      `payout state conflict: Stripe paid ${ids.length} releases but ${transitioned.length} accepted the released transition`,
    );
  }
}

async function recoverStrandedReleasing(
  msg: SettlementReleaseMessage,
  db: DoDbClient,
  stripe: StripePayoutClient,
): Promise<void> {
  const releases = await findReleasingReleases(db, msg.releaseIds);
  if (releases.length === 0) return;

  const ids = releases.map((release) => release.id);
  if (releases.some((release) => release.payoutBatchKey !== msg.batchKey)) {
    await holdReleases(db, ids, 'batch_key_mismatch');
    throw new SettlementBatchInvariantError(
      'releasing settlement rows do not match message batch key',
    );
  }

  const vendorAcctIds = new Set(releases.map((release) => release.vendorAcctId));
  if (vendorAcctIds.size !== 1) {
    await holdReleases(db, ids, 'mixed_vendor_batch');
    throw new SettlementBatchInvariantError('releasing settlement batch spans multiple vendors');
  }

  const vendorAcctId = releases[0]!.vendorAcctId;
  const totalAgorot = releases.reduce((sum, release) => sum + release.netAmountAgorot, 0);
  const idempotencyKey = await validatePayoutKeyState(db, ids, releases, true);

  let payout: { id: string };
  try {
    payout = await stripe.payouts.create(
      {
        amount: totalAgorot,
        currency: 'ils',
        metadata: { batchKey: msg.batchKey, releaseIds: ids.join(',') },
      },
      {
        stripeAccount: vendorAcctId,
        idempotencyKey,
      },
    );
  } catch (err) {
    if (isPermanentError(err)) {
      const errorMsg = err instanceof Error ? err.message : String(err);
      await holdReleases(db, ids, `recovery_payout_failed: ${errorMsg}`);
    }
    throw err;
  }

  const transitioned = await releaseReleases(db, ids, payout.id);
  if (transitioned.length !== ids.length) {
    throw new SettlementPayoutStateConflictError(
      `recovery payout state conflict: Stripe paid ${ids.length} releases but ${transitioned.length} accepted the released transition`,
    );
  }
}
