/**
 * Cron: Process Translation Jobs — runs every 30 minutes (backstop; queue consumer is primary).
 *
 * Re-drives translation_jobs rows whose queue messages were lost (DLQ, healthcheck
 * ack-without-process, failed self-requeue send). The consumer's atomic claim is the
 * concurrency guard — duplicate sends are idempotent (claim returns 0 rows, message acks).
 *
 * Step 0: reset RUNNING jobs abandoned by a killed worker (>10 min).
 * Step 1: re-enqueue stuck PENDING / PENDING_BUDGET jobs (5-min idle guard).
 * Step 2: terminally fail jobs that exhausted TRANSLATION_RETRY_MAX.
 */

import { createDbService } from '@/server/services/db.js';
import { and, lt, lte, asc, sql, inArray, isNull, or, gte } from 'drizzle-orm';
import { translationJobs } from '../db/schema.js';
import { TranslationEnabledSchema } from '../env.js';
import { TRANSLATION_RETRY_MAX } from '../translation/jobs/constants.js';
import { failJob } from '../translation/jobs/deadletter.js';
import { captureCaught } from '../observability/capture.server.js';
import { scrubErrorForLog } from '../observability/pii-scrub.js';
import type { CronEnv } from './deal-expiry.js';
import { reapRunningTranslationJobs } from '../db/queries/cron-boundary.js';

export interface TranslationCronEnv extends CronEnv {
  /** Feature flag for deal auto-translation pipeline. */
  TRANSLATION_ENABLED?: string;
  /** Producer binding for translation job re-enqueue. */
  TRANSLATION_QUEUE?: Queue<{ jobId: string }>;
}

const MAX_REQUEUE_PER_RUN = 25;

export async function runProcessTranslationJobs(env: TranslationCronEnv): Promise<void> {
  try {
    const translationEnabled = TranslationEnabledSchema.parse(env.TRANSLATION_ENABLED);
    if (!translationEnabled) {
      return;
    }

    const db = env.db ?? createDbService({ DATABASE_URL: env.DATABASE_URL });

    // ── Step 0: reap RUNNING jobs abandoned by a killed worker ─────────────────
    // Worker death (CF CPU/wall-clock limit) leaves jobs RUNNING forever — step 1
    // only re-enqueues PENDING. Reset any job stuck RUNNING for >10 min back to PENDING.
    await reapRunningTranslationJobs(db);

    // ── Step 1: re-enqueue stuck due PENDING / PENDING_BUDGET jobs ─────────────
    if (env.TRANSLATION_QUEUE) {
      const stuckJobs = await db
        .select({ id: translationJobs.id })
        .from(translationJobs)
        .where(
          and(
            inArray(translationJobs.status, ['PENDING', 'PENDING_BUDGET']),
            or(isNull(translationJobs.scheduledAt), lte(translationJobs.scheduledAt, sql`now()`)),
            lt(translationJobs.attempt, TRANSLATION_RETRY_MAX),
            isNull(translationJobs.startedAt),
            lt(translationJobs.createdAt, sql`now() - interval '5 minutes'`),
          ),
        )
        .orderBy(asc(translationJobs.createdAt))
        .limit(MAX_REQUEUE_PER_RUN);

      for (const job of stuckJobs) {
        await env.TRANSLATION_QUEUE.send({ jobId: job.id }).catch((e: unknown) => {
          captureCaught(e, {
            scope: 'cron.translation.queue_send',
            severity: 'warning',
            extra: { jobId: job.id },
          });
        });
      }
    } else {
      console.warn(
        JSON.stringify({
          event: 'process_translation_jobs_skip_requeue',
          reason: 'TRANSLATION_QUEUE binding not set',
        }),
      );
    }

    // ── Step 2: terminally fail exhausted jobs ─────────────────────────────────
    const exhaustedJobs = await db
      .select({ id: translationJobs.id })
      .from(translationJobs)
      .where(
        and(
          inArray(translationJobs.status, ['PENDING', 'PENDING_BUDGET']),
          gte(translationJobs.attempt, TRANSLATION_RETRY_MAX),
          isNull(translationJobs.startedAt),
        ),
      )
      .orderBy(asc(translationJobs.createdAt))
      .limit(MAX_REQUEUE_PER_RUN);

    for (const job of exhaustedJobs) {
      await failJob(db, job.id, 'max attempts exhausted (reaper)');
    }
  } catch (err) {
    console.error(
      JSON.stringify({
        event: 'process_translation_jobs_fatal',
        error: scrubErrorForLog(err),
      }),
    );
  }
}
