/**
 * Cron: Audience Sync — runs daily at 03:00 UTC.
 *
 * Sweeps all opted-in users and upserts them to the Brevo contact list
 * identified by BREVO_LIST_ID. Uses cursor pagination (PAGE_SIZE per batch)
 * to stay within CF Worker CPU limits.
 *
 * users.email is stored encrypted (pgcrypto). Each row is decrypted via
 * decryptUserEmail() before being passed to Brevo.
 */

import { createDbService } from '@/server/services/db.js';
import { withSentry } from '@/server/observability/with-sentry';
import { captureCaught } from '@/server/observability/capture.server';
import { BrevoMailer } from '../email/bulk-mailer.js';
import { getCursor, setCursor } from '../db/queries/marketing.js';
import { decryptUserEmail } from '../workflows/outbox/_helpers/decrypt-email.js';
import type { CronEnv } from './deal-expiry.js';
import { sql } from 'drizzle-orm';

const PAGE_SIZE = 200;
const CURSOR_KEY = 'brevo.audience.sync.cursor';

export const runAudienceSync = withSentry(
  async function runAudienceSync(env: CronEnv): Promise<void> {
    if (!env.BREVO_API_KEY || !env.BREVO_LIST_ID) {
      console.warn('[audience-sync] BREVO_API_KEY or BREVO_LIST_ID not set — skipping');
      return;
    }

    const listId = parseInt(env.BREVO_LIST_ID, 10);
    if (isNaN(listId)) {
      console.error('[audience-sync] BREVO_LIST_ID is not a valid integer');
      return;
    }

    const mailer = new BrevoMailer(env.BREVO_API_KEY);
    const db = createDbService({ DATABASE_URL: env.DATABASE_URL });

    let cursor = await getCursor(db, CURSOR_KEY);
    let synced = 0;
    let errors = 0;

    try {
      while (true) {
        const rows = await db.execute<{ id: string; email: string }>(
          sql`SELECT u.id, u.email
              FROM users u
              WHERE u.email IS NOT NULL
                AND ${cursor ? sql`u.id > ${cursor}` : sql`TRUE`}
                AND (u.notif_prefs->'marketing'->>'all')::boolean = true
              ORDER BY u.id
              LIMIT ${PAGE_SIZE}`,
        );

        if (!rows.rows.length) break;

        for (const row of rows.rows) {
          try {
            const email = await decryptUserEmail(env, row.id);
            if (!email) {
              console.warn('[audience-sync] email decrypt failed — skipping', {
                userId_hash: row.id.slice(0, 8),
              });
              errors++;
              continue;
            }
            await mailer.upsertContact({ email }, listId);
            synced++;
          } catch (err) {
            errors++;
            console.error('[audience-sync] upsert error', {
              userId_hash: row.id.slice(0, 8),
              err: String(err),
            });
          }
        }

        cursor = rows.rows[rows.rows.length - 1]!.id;
        await setCursor(db, CURSOR_KEY, cursor);

        if (rows.rows.length < PAGE_SIZE) break;
      }
    } catch (err) {
      captureCaught(err, { scope: 'server.cron.audience-sync', severity: 'warning' });
    }

    // Reset cursor after full sweep so next run re-syncs all opted-in users
    await setCursor(db, CURSOR_KEY, '');

    console.warn(JSON.stringify({ event: 'audience_sync_complete', synced, errors }));
  },
  { name: 'cron.audience-sync', kind: 'cron' },
);
