import { createScheduledDbService } from '@/server/services/db.js';
/**
 * Cloudflare Queue consumer for `multideal-outbox-preview`.
 *
 * Processes outbox rows per-message (Phase 2 cutover). Each message carries
 * `{ outboxId: string }` — the consumer loads the outbox row, dispatches the
 * side effect via shared `dispatchOutboxRow`, then acks or retries the message.
 *
 * Runs local to this worker (multideal-preview) — all secrets (DATABASE_URL,
 * RESEND_API_KEY, PII_KEY, VAPID_*) are already configured, so handlers
 * dispatch email + loyalty events end-to-end. The `process-outbox` cron in
 * this worker continues running as a safety backstop and picks up any row
 * this consumer did not fully dispatch.
 *
 * Row-not-found: ack silently (row deleted, never committed, or duplicate
 * delivery after cron already processed it). Do NOT retry — there is nothing
 * to dispatch and retrying won't help.
 */

import { z } from 'zod';
import { eq } from 'drizzle-orm';
import type { MultidealEnv } from '../lib/env.js';
import { dispatchOutboxRow } from '@/server/workflows/outbox/dispatcher';
import {
  markOutboxEventProcessed,
  markOutboxEventFailed,
  outboxRetryLimitForError,
} from '@/server/db/queries/outbox';

import { outbox } from '@/server/db/schema';
import { captureCaught } from '@/server/observability/capture.server';

// ─── Zod validation ───────────────────────────────────────────────────────────

const OutboxMessageSchema = z.object({
  outboxId: z.uuid(),
});

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

export async function handleOutboxBatch(
  batch: MessageBatch<{ outboxId: string }>,
  env: MultidealEnv,
  _ctx: ExecutionContext,
): Promise<void> {
  const db = createScheduledDbService(env);

  for (const message of batch.messages) {
    // 1. Zod-validate message body (§1 Zod-at-the-boundary law)
    const parsed = OutboxMessageSchema.safeParse(message.body);
    if (!parsed.success) {
      console.error(
        JSON.stringify({
          event: 'outbox_invalid_message',
          error: parsed.error.message,
          body: message.body,
        }),
      );
      // Malformed message — ack to prevent infinite retry on structurally bad data
      message.ack();
      continue;
    }

    const { outboxId } = parsed.data;

    // 2. Load outbox row
    const rows = await db.select().from(outbox).where(eq(outbox.id, outboxId)).limit(1);

    const row = rows[0];

    if (!row) {
      console.warn(
        JSON.stringify({
          event: 'outbox_row_not_found',
          outboxId,
          note: 'Row missing — likely already processed by cron backstop or never committed. Acking.',
        }),
      );
      message.ack();
      continue;
    }

    // 3. Skip already-processed rows (race with cron or duplicate queue delivery)
    if (row.processedAt) {
      console.info(
        JSON.stringify({
          event: 'outbox_already_processed',
          outboxId,
          processedAt: row.processedAt,
        }),
      );
      message.ack();
      continue;
    }

    // Terminal (dead) rows have exhausted retries — never re-attempt the side effect.
    // Surfaced by the outbox_abandoned check; recovered only via admin re-drive.
    if (row.deadAt) {
      console.info(
        JSON.stringify({
          event: 'outbox_dead_skipped',
          outboxId,
          deadAt: row.deadAt,
        }),
      );
      message.ack();
      continue;
    }

    // 4. Dispatch side effect — strict=false so missing secrets warn+ack rather than throw
    try {
      await dispatchOutboxRow(
        {
          DATABASE_URL: env.DATABASE_URL,
          PAYMENT_PROVIDER: env.PAYMENT_PROVIDER,
          STRIPE_SECRET_KEY: env.STRIPE_SECRET_KEY,
          RESEND_API_KEY: env.RESEND_API_KEY,
          RESEND_FROM_EMAIL: env.RESEND_FROM_EMAIL,
          // Non-prod mock sink — without these, resend.ts really sends every
          // transactional email (receipts etc.) on preview to decrypted real
          // addresses. This is the real-traffic path. See design addendum 2026-06-17.
          ENVIRONMENT: env.ENVIRONMENT,
          EMAIL_MOCK_DB: env.EMAIL_MOCK_DB,
          VAPID_PUBLIC_KEY: env.VAPID_PUBLIC_KEY,
          VAPID_PRIVATE_KEY: env.VAPID_PRIVATE_KEY,
          VAPID_SUBJECT: env.VAPID_SUBJECT,
          PII_KEY: env.PII_KEY,
          GOOGLE_API_KEY: env.GOOGLE_API_KEY,
          TRANSLATION_QUEUE: env.TRANSLATION_QUEUE,
          USER_SESSION_DO: env.USER_SESSION_DO,
          TOPIC_DO: env.TOPIC_DO,
          CACHE_EPOCH_DO: env.CACHE_EPOCH_DO,
          PUBLIC_SITE_URL: env.PUBLIC_SITE_URL,
          strict: false,
        },
        row,
      );

      await markOutboxEventProcessed(db, outboxId);
      message.ack();
    } catch (err) {
      const errorMsg = err instanceof Error ? err.message : String(err);
      console.error(
        JSON.stringify({
          event: 'outbox_consume_failed',
          outboxId,
          kind: row.eventType,
          error: errorMsg,
        }),
      );

      // Record failure in DB for admin visibility + cron retry; cap sets dead_at (outbox_abandoned pages).
      await markOutboxEventFailed(db, outboxId, errorMsg, outboxRetryLimitForError(err)).catch(
        (markErr) => {
          // Non-fatal — DB write failure shouldn't prevent queue retry
          captureCaught(markErr, {
            scope: 'server.do-host.queues.outbox-consumer',
            severity: 'warning',
          });
        },
      );

      // Let the queue's max_retries (3) + DLQ handle the failure
      message.retry();
    }
  }
}
