import {
  getRedisUrl,
  initializeConfigFromDatabase,
  startConfigSubscription,
  useEnvironmentConfigFallback,
} from './config';
import dns from 'dns';
import Redis from 'ioredis';
import { randomUUID } from 'crypto';

// Set Google DNS for resolution as Tailscale might not resolve dev3.press.zone correctly
dns.setServers(['8.8.8.8', '4.4.4.4']);

/**
 * Background Job Worker
 *
 * Processes async translation jobs using Bull queue
 */

import { PrismaClient } from '@prisma/client';
import { logger } from './utils/logger';
import { geminiClient } from './services/geminiClient';
import { exceptionService } from './services/exceptionService';
import { verifyExceptions } from './services/exceptionVerifier';
import { deliverWebhook } from './services/webhookService';
import { mapStoredTranslation } from './services/translationService';
import {
  countSiteContentCharacters,
  siteContentCallbackMetadata,
  SiteContentJobSubmitRequest,
  translateSiteContentResource,
} from './services/siteContentTranslation';
import { countSourceCharacters, calculateCost, calculateCustomerCostByChars } from './utils/tokenCalculation';
import {
  countStructuredCharacters,
  serializeStructuredFields,
  splitStructuredTranslation,
  StructuredFields,
} from './utils/structuredFields';
import { trackTranslationJob, trackTokensProcessed } from './utils/metrics';
import { TranslationStatus, BulkTranslationJobData, Tone, WebhookPayload } from './types';
import { translationQueue } from './queue';
import { reconcilePermanentOrphanDisputes } from './workers/orphanDisputeReconciler';
import { WorkerHeartbeat, getWorkerHeartbeatKey } from './workerHeartbeat';
import { shutdownWorkerRuntime, startWorkerRuntime } from './workerLifecycle';
import { JobCancellationWonError, lockProcessingJob } from './workerLock';
import { configureBulkQueueDispatcher, startBulkQueueDispatcher, stopBulkQueueDispatcher } from './queue/bulkQueueDispatcher';

const prisma = new PrismaClient();
configureBulkQueueDispatcher(prisma.translationJob);
const workerId = (process.env.WORKER_ID || process.env.HOSTNAME)?.trim();
if (!workerId) {
  throw new Error('WORKER_ID or HOSTNAME is required for worker identity');
}
let configSubscription: ReturnType<typeof startConfigSubscription> | undefined;
const workerHeartbeat = new WorkerHeartbeat({
  redis: new Redis(getRedisUrl()),
  checkQueueReady: () => translationQueue.isReady(),
  checkRuntimeReady: () => {
    if (!configSubscription) {
      return Promise.reject(new Error('Config subscription is not initialized'));
    }
    return configSubscription.checkReady();
  },
  checkDatabase: () => prisma.$queryRaw`SELECT 1`,
  key: getWorkerHeartbeatKey(workerId),
  ttlSeconds: Number(process.env.WORKER_HEARTBEAT_TTL_SECONDS || 15),
});

interface TranslationJobData {
  jobId: string;
  clientJobId?: string;
  userId: string;
  content?: string;
  fields?: StructuredFields;
  siteContent?: SiteContentJobSubmitRequest;
  sourceLang: string;
  targetLang: string;
  tone: string;
  exceptions?: unknown;
  callbackUrl?: string;
  callbackSecret?: string;
}

type WebhookOutboxEntry = {
  id: string;
  job_id: string;
  attempts: number;
  payload: unknown;
  job: { callback_url: string | null; callback_secret: string | null };
};

type WebhookOutboxDelegate = {
  findMany(_args: object): Promise<WebhookOutboxEntry[]>;
  updateMany(_args: object): Promise<{ count: number }>;
  create(_args: object): Promise<unknown>;
};

type PrismaWithWebhookOutbox = PrismaClient & { webhookOutbox: WebhookOutboxDelegate };

export const OUTBOX_MAX_ATTEMPTS = 5;
const OUTBOX_LEASE_MS = 60_000;
const OUTBOX_BASE_DELAY_MS = 5_000;
const OUTBOX_MAX_DELAY_MS = 60_000;
const OUTBOX_BATCH_SIZE = 100;

type OutboxDrainResult = {
  claimed: number;
  delivered: number;
  failed: number;
  deadLettered: number;
};

let outboxDrainPromise: Promise<OutboxDrainResult> | undefined;
let outboxStartPromise: Promise<void> | undefined;
let outboxTimer: ReturnType<typeof setTimeout> | undefined;
let outboxStarted = false;
let outboxStopRequested = false;
let outboxBackoffMs = OUTBOX_BASE_DELAY_MS;

const PROCESSING_LEASE_MS = 120_000;
const PROCESSING_LEASE_RENEWAL_MS = 30_000;

export type ProcessingLeaseClaim = {
  claimed: boolean;
  leaseId: string;
};

export async function claimTranslationJob(
  jobId: string,
  leaseId = randomUUID(),
  now = new Date()
): Promise<ProcessingLeaseClaim> {
  const expiresAt = new Date(now.getTime() + PROCESSING_LEASE_MS);
  const result = await prisma.translationJob.updateMany({
    where: {
      id: jobId,
      OR: [
        { status: 'pending' },
        { status: 'processing', processing_lease_expires_at: { lt: now } },
      ],
    },
    data: {
      status: 'processing' as TranslationStatus,
      processing_lease_id: leaseId,
      processing_lease_expires_at: expiresAt,
    },
  });
  return { claimed: result.count === 1, leaseId };
}

function startProcessingLeaseHeartbeat(jobId: string, leaseId: string): () => void {
  const timer = setInterval(() => {
    void prisma.translationJob.updateMany({
      where: { id: jobId, status: 'processing', processing_lease_id: leaseId },
      data: { processing_lease_expires_at: new Date(Date.now() + PROCESSING_LEASE_MS) },
    }).catch((error) => logger.error('Translation processing lease renewal failed', { jobId, error }));
  }, PROCESSING_LEASE_RENEWAL_MS);
  timer.unref?.();
  return () => clearInterval(timer);
}

async function releaseProcessingLease(jobId: string, leaseId: string): Promise<boolean> {
  const result = await prisma.translationJob.updateMany({
    where: { id: jobId, status: 'processing', processing_lease_id: leaseId },
    data: { status: 'pending' as TranslationStatus, processing_lease_id: null, processing_lease_expires_at: null },
  });
  return result.count === 1;
}

function isTerminalWorkerError(error: unknown): boolean {
  const message = error instanceof Error ? error.message : '';
  return /insufficient credits|validation|invalid/i.test(message);
}

async function settleOutboxEntry(
  outboxPrisma: PrismaWithWebhookOutbox,
  entryId: string,
  claimedAt: Date,
  data: object
): Promise<boolean> {
  const result = await outboxPrisma.webhookOutbox.updateMany({
    where: { id: entryId, status: 'delivering', claimed_at: claimedAt },
    data,
  });
  return result.count === 1;
}

async function runCompletionOutboxDrain(jobId?: string): Promise<OutboxDrainResult> {
  const result: OutboxDrainResult = { claimed: 0, delivered: 0, failed: 0, deadLettered: 0 };
  const outboxPrisma = prisma as PrismaWithWebhookOutbox;
  const leaseCutoff = new Date(Date.now() - OUTBOX_LEASE_MS);
  const jobFilter = jobId ? { job_id: jobId } : {};

  await outboxPrisma.webhookOutbox.updateMany({
    where: { ...jobFilter, status: 'delivering', claimed_at: { lt: leaseCutoff } },
    data: { status: 'pending', claimed_at: null },
  });
  const exhausted = await outboxPrisma.webhookOutbox.updateMany({
    where: { ...jobFilter, status: 'pending', attempts: { gte: OUTBOX_MAX_ATTEMPTS } },
    data: { status: 'dead', last_error: `Webhook delivery exhausted after ${OUTBOX_MAX_ATTEMPTS} attempts` },
  });
  result.deadLettered += exhausted.count;

  const outboxEntries = await outboxPrisma.webhookOutbox.findMany({
    where: {
      ...(jobId ? { job_id: jobId } : {}),
      event: { in: ['translation.completed', 'bulk_translation.completed'] },
      status: 'pending',
      attempts: { lt: OUTBOX_MAX_ATTEMPTS },
    },
    orderBy: { created_at: 'asc' },
    take: OUTBOX_BATCH_SIZE,
    include: { job: { select: { callback_url: true, callback_secret: true } } },
  });

  for (const entry of outboxEntries) {
    const claimedAt = new Date();
    const claimed = await outboxPrisma.webhookOutbox.updateMany({
      where: { id: entry.id, status: 'pending', attempts: { lt: OUTBOX_MAX_ATTEMPTS } },
      data: { status: 'delivering', attempts: { increment: 1 }, claimed_at: claimedAt },
    });
    if (claimed.count !== 1) {
      continue;
    }
    result.claimed += 1;
    const attemptNumber = (entry.attempts || 0) + 1;

    if (!entry.job.callback_url || !entry.job.callback_secret) {
      const deadLettered = await settleOutboxEntry(outboxPrisma, entry.id, claimedAt, {
        status: 'dead',
        last_error: 'Webhook callback configuration is missing',
      });
      if (deadLettered) {
        result.deadLettered += 1;
      }
      continue;
    }

    try {
      const delivery = await deliverWebhook(
        entry.job_id,
        entry.payload as unknown as WebhookPayload,
        entry.job.callback_url,
        entry.job.callback_secret,
        entry.id
      );
      if (delivery.success) {
        if (await settleOutboxEntry(outboxPrisma, entry.id, claimedAt, {
          status: 'delivered',
          delivered_at: new Date(),
          claimed_at: null,
          last_error: null,
        })) {
          result.delivered += 1;
        }
      } else {
        const status = attemptNumber >= OUTBOX_MAX_ATTEMPTS ? 'dead' : 'pending';
        if (await settleOutboxEntry(outboxPrisma, entry.id, claimedAt, {
          status,
          claimed_at: null,
          last_error: delivery.errorMessage || 'Webhook delivery failed',
        })) {
          result.failed += 1;
          if (status === 'dead') {
            result.deadLettered += 1;
          }
        }
      }
    } catch (error) {
      const status = attemptNumber >= OUTBOX_MAX_ATTEMPTS ? 'dead' : 'pending';
      if (await settleOutboxEntry(outboxPrisma, entry.id, claimedAt, {
        status,
        claimed_at: null,
        last_error: error instanceof Error ? error.message : 'Webhook delivery failed',
      })) {
        result.failed += 1;
        if (status === 'dead') {
          result.deadLettered += 1;
        }
      }
    }
  }

  return result;
}

export async function drainCompletionOutbox(jobId?: string): Promise<OutboxDrainResult> {
  if (!outboxDrainPromise) {
    outboxDrainPromise = runCompletionOutboxDrain(jobId).catch((error) => {
      logger.error('Completion webhook outbox drain failed', { error, jobId });
      return { claimed: 0, delivered: 0, failed: 1, deadLettered: 0 };
    });
  }
  try {
    return await outboxDrainPromise;
  } finally {
    outboxDrainPromise = undefined;
  }
}

function scheduleCompletionOutboxDrain(delayMs: number): void {
  if (!outboxStarted || outboxStopRequested) {
    return;
  }
  outboxTimer = setTimeout(() => {
    outboxTimer = undefined;
    void drainCompletionOutbox().then((result) => {
      outboxBackoffMs = result.failed > 0
        ? Math.min(outboxBackoffMs * 2, OUTBOX_MAX_DELAY_MS)
        : OUTBOX_BASE_DELAY_MS;
    }).catch((error) => {
      logger.error('Completion webhook outbox recurring drain failed', { error });
      outboxBackoffMs = Math.min(outboxBackoffMs * 2, OUTBOX_MAX_DELAY_MS);
    }).finally(() => scheduleCompletionOutboxDrain(outboxBackoffMs));
  }, delayMs);
  outboxTimer.unref?.();
}

export async function startCompletionOutboxDrainer(): Promise<void> {
  if (outboxStarted) {
    if (outboxStartPromise) {
      await outboxStartPromise;
    }
    return;
  }
  outboxStarted = true;
  outboxStopRequested = false;
  outboxBackoffMs = OUTBOX_BASE_DELAY_MS;
  outboxStartPromise = (async () => {
    const result = await drainCompletionOutbox();
    outboxBackoffMs = result.failed > 0
      ? Math.min(OUTBOX_BASE_DELAY_MS * 2, OUTBOX_MAX_DELAY_MS)
      : OUTBOX_BASE_DELAY_MS;
    scheduleCompletionOutboxDrain(outboxBackoffMs);
  })();
  try {
    await outboxStartPromise;
  } finally {
    outboxStartPromise = undefined;
  }
}

export async function stopCompletionOutboxDrainer(): Promise<void> {
  outboxStopRequested = true;
  outboxStarted = false;
  if (outboxTimer) {
    clearTimeout(outboxTimer);
    outboxTimer = undefined;
  }
  await outboxStartPromise;
  await outboxDrainPromise;
}

export function registerProcessors(): void {

/**
 * Process translation job
 */
translationQueue.process('translate', async (job) => {
  const startTime = Date.now();
  const data: TranslationJobData = job.data;
  const exceptionRules = exceptionService.normalizeRules(data.exceptions);

  logger.info('Processing translation job', { jobId: data.jobId, clientJobId: data.clientJobId });

  const claim = await claimTranslationJob(data.jobId);
  if (!claim.claimed) {
    const current = await prisma.translationJob.findUnique({
      where: { id: data.jobId },
      select: { status: true, translation: true },
    });
    if (current?.status === 'completed') {
      await drainCompletionOutbox(data.jobId);
      return { translation: current.translation || undefined };
    }
    return { status: current?.status || 'cancelled' };
  }
  const stopLeaseHeartbeat = startProcessingLeaseHeartbeat(data.jobId, claim.leaseId);

  try {

    let storedTranslation: string;
    let translatedFields: StructuredFields | undefined;
    let translatedOutputFields: StructuredFields | undefined;
    let tokensUsed: number;
    let inputTokens: number;
    let outputTokens: number;
    let modelUsed: string;

    if (data.siteContent) {
      const siteContentResult = await translateSiteContentResource(
        data.siteContent,
        async (fields) => {
          const processedFields: StructuredFields = {};
          const replacements = new Map<string, Map<string, string>>();
          for (const [key, value] of Object.entries(fields)) {
            const processed = exceptionService.replaceExceptions(exceptionRules, value);
            processedFields[key] = processed.processedText;
            if (processed.replacements.size > 0) {
              replacements.set(key, processed.replacements);
            }
          }

          const structuredResult = await geminiClient.translateStructured(
            processedFields,
            data.sourceLang,
            data.targetLang,
            data.tone as Tone
          );
          const outputFields = structuredResult.translatedFields || structuredResult.fields;
          for (const [key, fieldReplacements] of replacements.entries()) {
            const translated = outputFields[key];
            if (translated !== undefined) {
              outputFields[key] = exceptionService.restoreExceptions(translated, fieldReplacements);
            }
          }

          return {
            translatedFields: outputFields,
            tokens_used: structuredResult.tokens_used,
            input_tokens: structuredResult.input_tokens,
            output_tokens: structuredResult.output_tokens,
            processing_time_ms: structuredResult.processing_time_ms,
            model_used: structuredResult.model_used,
          };
        }
      );
      storedTranslation = siteContentResult.translation;
      tokensUsed = siteContentResult.tokens_used;
      inputTokens = siteContentResult.input_tokens;
      outputTokens = siteContentResult.output_tokens;
      modelUsed = siteContentResult.model_used;
    } else if (data.fields) {
      const processedFields: StructuredFields = {};
      const replacements = new Map<string, Map<string, string>>();

      for (const [key, value] of Object.entries(data.fields)) {
        const processed = exceptionService.replaceExceptions(exceptionRules, value);
        processedFields[key] = processed.processedText;
        if (processed.replacements.size > 0) {
          replacements.set(key, processed.replacements);
        }
      }

      const structuredResult = await geminiClient.translateStructured(
        processedFields,
        data.sourceLang,
        data.targetLang,
        data.tone as Tone
      );
      translatedOutputFields = structuredResult.translatedFields || structuredResult.fields;

      for (const [key, fieldReplacements] of replacements.entries()) {
        const translated = translatedOutputFields[key];
        if (translated !== undefined) {
          translatedOutputFields[key] = exceptionService.restoreExceptions(translated, fieldReplacements);
        }
      }

      translatedFields = splitStructuredTranslation(translatedOutputFields).translatedFields;
      storedTranslation = serializeStructuredFields(translatedOutputFields);
      tokensUsed = structuredResult.tokens_used;
      inputTokens = structuredResult.input_tokens;
      outputTokens = structuredResult.output_tokens;
      modelUsed = structuredResult.model_used;
    } else {
      const content = data.content || '';
      const { processedText, replacements, matchedExceptions } =
        exceptionService.replaceExceptions(exceptionRules, content);

      const result = await geminiClient.translate(
        processedText,
        data.sourceLang,
        data.targetLang,
        data.tone as Tone
      );

      storedTranslation = result.translation;
      if (replacements.size > 0) {
        storedTranslation = exceptionService.restoreExceptions(storedTranslation, replacements);

        if (matchedExceptions.length > 0) {
          const verification = verifyExceptions(storedTranslation, matchedExceptions);
          if (!verification.passed) {
            logger.warn('Worker: Exception verification warnings', {
              jobId: data.jobId,
              violations: verification.violations,
            });
          }
        }
      }

      tokensUsed = result.tokens_used;
      inputTokens = result.input_tokens;
      outputTokens = result.output_tokens;
      modelUsed = result.model_used;
    }

    const processingTime = Date.now() - startTime;
    const charactersUsed = data.siteContent
      ? countSiteContentCharacters(data.siteContent)
      : data.fields
        ? countStructuredCharacters(data.fields)
        : countSourceCharacters(data.content || '');
    const internalCost = calculateCost(inputTokens, outputTokens, modelUsed);

    // Look up user's subscription to get customer cost per character
    const subscription = await prisma.subscription.findFirst({
      where: { user_id: data.userId },
      select: { customer_cost_per_char: true },
    });
    const costPerChar = subscription?.customer_cost_per_char
      ? Number(subscription.customer_cost_per_char)
      : 0;
    const customerCost = calculateCustomerCostByChars(charactersUsed, costPerChar);

    const completionPayload = data.callbackUrl && data.callbackSecret
      ? {
        event: 'translation.completed' as const,
        jobId: data.jobId,
        clientJobId: data.clientJobId,
        status: 'completed' as TranslationStatus,
        ...(data.siteContent ? siteContentCallbackMetadata(data.siteContent) : {}),
        ...mapStoredTranslation(storedTranslation),
        charactersUsed,
        cost: customerCost,
        processingTimeMs: processingTime,
        timestamp: new Date().toISOString(),
      }
      : null;

    await prisma.$transaction(async (tx) => {
      await lockProcessingJob(tx, data.jobId, claim.leaseId);

      const user = await tx.user.findUnique({
        where: { id: data.userId },
        select: { id: true },
      });

      if (!user) {
        throw new Error('User not found');
      }

      const latestTransaction = await tx.creditTransaction.findFirst({
        where: { user_id: data.userId },
        orderBy: { created_at: 'desc' },
        select: { balance_after: true },
      });

      const currentBalance = latestTransaction?.balance_after ?? 0;
      const newBalance = currentBalance - charactersUsed;

      await tx.creditTransaction.create({
        data: {
          user_id: data.userId,
          type: 'deduction',
          amount: -charactersUsed,
          balance_after: newBalance,
          description: `Translation job ${data.jobId}: ${data.sourceLang} → ${data.targetLang}`,
          related_job_id: data.jobId,
        },
      });

      await tx.user.update({
        where: { id: data.userId },
        data: {
          updated_at: new Date(),
        },
      });

      await tx.translationJob.update({
        where: { id: data.jobId },
        data: {
          status: 'completed' as TranslationStatus,
          translation: storedTranslation,
          model: modelUsed,
          characters_used: charactersUsed,
          tokens_used: tokensUsed,
          input_tokens: inputTokens,
          output_tokens: outputTokens,
          cost: internalCost,
          customer_cost: customerCost,
          processing_time_ms: processingTime,
          completed_at: new Date(),
          processing_lease_id: null,
          processing_lease_expires_at: null,
        },
      });

      if (completionPayload) {
        const txWithOutbox = tx as typeof tx & { webhookOutbox: WebhookOutboxDelegate };
        await txWithOutbox.webhookOutbox.create({
          data: {
            job_id: data.jobId,
            event: 'translation.completed',
            payload: completionPayload,
          },
        });
      }
    });

    // Track metrics
    trackTranslationJob('gemini', 'completed', 'async', processingTime);
    trackTokensProcessed('gemini', tokensUsed);

    await drainCompletionOutbox(data.jobId);

    logger.info('Translation job completed', {
      jobId: data.jobId,
      characters_used: charactersUsed,
      tokens_used: tokensUsed,
      cost: internalCost,
      processingTime,
    });

    return {
      translation: storedTranslation,
      translatedFields,
      tokensUsed,
      processingTimeMs: processingTime,
    };
  } catch (error: any) {
    const processingTime = Date.now() - startTime;

    if (error instanceof JobCancellationWonError) {
      return { status: 'cancelled' };
    }

    const publicErrorMessage = data.siteContent
      ? 'Site Content translation failed.'
      : error instanceof Error
        ? error.message
        : 'Translation failed';
    const outwardError = data.siteContent ? new Error(publicErrorMessage) : error;
    logger.error('Translation job failed', {
      jobId: data.jobId,
      error: data.siteContent ? publicErrorMessage : error,
    });

    const isFinalAttempt = (job.opts?.attempts ?? 1) <= job.attemptsMade + 1;
    if (!isFinalAttempt && !isTerminalWorkerError(error)) {
      await releaseProcessingLease(data.jobId, claim.leaseId);
      throw outwardError;
    }

    const failed = await prisma.translationJob.updateMany({
      where: { id: data.jobId, status: 'processing', processing_lease_id: claim.leaseId },
      data: {
        status: 'failed' as TranslationStatus,
        processing_lease_id: null,
        processing_lease_expires_at: null,
        error_message: publicErrorMessage,
        processing_time_ms: processingTime,
        completed_at: new Date(),
      },
    });
    if (failed.count !== 1) {
      const current = await prisma.translationJob.findUnique({
        where: { id: data.jobId },
        select: { status: true, translation: true },
      });
      if (current?.status === 'cancelled' || current?.status === 'completed') {
        return { status: current.status, translation: current.translation || undefined };
      }
      throw outwardError;
    }

    // Track metrics
    trackTranslationJob('gemini', 'failed', 'async', processingTime);

    // Deliver failure webhook
    if (data.callbackUrl && data.callbackSecret) {
      try {
        await deliverWebhook(
          data.jobId,
          {
            event: 'translation.failed',
            jobId: data.jobId,
            clientJobId: data.clientJobId,
            status: 'failed' as TranslationStatus,
            ...(data.siteContent ? siteContentCallbackMetadata(data.siteContent) : {}),
            errorMessage: publicErrorMessage,
            processingTimeMs: processingTime,
            timestamp: new Date().toISOString(),
          },
          data.callbackUrl,
          data.callbackSecret
        );
      } catch (webhookError) {
        logger.error('Webhook delivery failed', { jobId: data.jobId, error: webhookError });
      }
    }

    throw outwardError;
  } finally {
    stopLeaseHeartbeat();
  }
});

const BULK_BATCH_SIZE = 50;

/**
 * Process async bulk-strings translation job.
 * Translates all strings for each target language (in batches of 50 per Gemini call),
 * deducts credits once for the entire job, then delivers one webhook with full results.
 */
translationQueue.process('bulk-strings', async (job) => {
  const startTime = Date.now();
  const data: BulkTranslationJobData = job.data;
  const exceptionRules = exceptionService.normalizeRules(data.exceptions);

  logger.info('Processing bulk-strings job', {
    jobId: data.jobId,
    stringCount: data.strings.length,
    targetLangs: data.targetLangs,
  });

  const claim = await claimTranslationJob(data.jobId);
  if (!claim.claimed) {
    const current = await prisma.translationJob.findUnique({
      where: { id: data.jobId },
      select: { status: true, translation: true },
    });
    if (current?.status === 'completed') {
      await drainCompletionOutbox(data.jobId);
      return { translation: current.translation || undefined };
    }
    return { status: current?.status || 'cancelled' };
  }
  const stopLeaseHeartbeat = startProcessingLeaseHeartbeat(data.jobId, claim.leaseId);

  const resultsByLang: Record<string, Array<{ id: string; translation: string; success: boolean }>> = {};
  let totalCharactersUsed = 0;
  let totalFailedCount = 0;

  try {
  // Process each target language in series
  for (const targetLang of data.targetLangs) {
    const langResults: Array<{ id: string; translation: string; success: boolean }> = [];
    const langChars = data.strings.reduce((sum, s) => sum + countSourceCharacters(s.content), 0);
    totalCharactersUsed += langChars;

    // Pre-process: replace exceptions in each string for this language
    const perStringReplacements: Map<string, Map<string, string>> = new Map();
    const processedStrings = [];

    for (const s of data.strings) {
      const { processedText, replacements } =
        exceptionService.replaceExceptions(exceptionRules, s.content);
      processedStrings.push({ id: s.id, content: processedText });
      if (replacements.size > 0) {
        perStringReplacements.set(s.id, replacements);
      }
    }

    // Split into batches of 50
    for (let i = 0; i < processedStrings.length; i += BULK_BATCH_SIZE) {
      const batch = processedStrings.slice(i, i + BULK_BATCH_SIZE);
      const bulkResponse = await geminiClient.translateBulk(
        batch,
        data.sourceLang,
        targetLang,
        (data.tone as Tone) || Tone.NEUTRAL
      );

      // Post-process: restore exceptions in each result
      for (const result of bulkResponse.results) {
        const reps = perStringReplacements.get(result.id);
        if (reps && result.success && result.translation) {
          result.translation = exceptionService.restoreExceptions(result.translation, reps);
        }
      }

      langResults.push(...bulkResponse.results);
    }

    totalFailedCount += langResults.filter(r => !r.success).length;
    resultsByLang[targetLang] = langResults;
  }

  const processingTime = Date.now() - startTime;
  const storedTranslation = JSON.stringify(resultsByLang);
  const completionPayload = {
    event: 'bulk_translation.completed' as const,
    jobId: data.jobId,
    clientJobId: data.clientJobId,
    status: 'completed' as TranslationStatus,
    translation: storedTranslation,
    job_id: data.jobId,
    results_by_lang: resultsByLang,
    total_characters_used: totalCharactersUsed,
    failed_count: totalFailedCount,
    timestamp: new Date().toISOString(),
  };

  try {
    await prisma.$transaction(async (tx) => {
      await lockProcessingJob(tx, data.jobId, claim.leaseId);
      const latestTransaction = await tx.creditTransaction.findFirst({
        where: { user_id: data.userId },
        orderBy: { created_at: 'desc' },
        select: { balance_after: true },
      });
      const currentBalance = latestTransaction?.balance_after ?? 0;
      await tx.creditTransaction.create({
        data: {
          user_id: data.userId,
          type: 'deduction',
          amount: -totalCharactersUsed,
          balance_after: currentBalance - totalCharactersUsed,
          description: `Bulk-strings job ${data.jobId}: ${data.sourceLang} → ${data.targetLangs.join(', ')}`,
          related_job_id: data.jobId,
        },
      });
      await tx.user.update({ where: { id: data.userId }, data: { updated_at: new Date() } });
      await tx.translationJob.update({
        where: { id: data.jobId },
        data: {
          status: 'completed' as TranslationStatus,
          translation: storedTranslation,
          characters_used: totalCharactersUsed,
          processing_time_ms: processingTime,
          completed_at: new Date(),
          processing_lease_id: null,
          processing_lease_expires_at: null,
        },
      });
      const txWithOutbox = tx as typeof tx & { webhookOutbox: WebhookOutboxDelegate };
      await txWithOutbox.webhookOutbox.create({
        data: {
          job_id: data.jobId,
          event: 'bulk_translation.completed',
          payload: completionPayload,
        },
      });
    });
  } catch (error) {
    if (error instanceof JobCancellationWonError) {
      return { status: 'cancelled' };
    }
    throw error;
  }

  await drainCompletionOutbox(data.jobId);

  logger.info('Bulk-strings job completed', {
    jobId: data.jobId,
    totalCharactersUsed,
    totalFailedCount,
    processingTimeMs: processingTime,
  });

  return {
    translation: storedTranslation,
    totalCharactersUsed,
    totalFailedCount,
    processingTimeMs: processingTime,
  };
  } catch (error) {
    if (error instanceof JobCancellationWonError) {
      return { status: 'cancelled' };
    }
    const isFinalAttempt = (job.opts?.attempts ?? 1) <= job.attemptsMade + 1;
    if (!isFinalAttempt && !isTerminalWorkerError(error)) {
      await releaseProcessingLease(data.jobId, claim.leaseId);
      throw error;
    }
    const failed = await prisma.translationJob.updateMany({
      where: { id: data.jobId, status: 'processing', processing_lease_id: claim.leaseId },
      data: {
        status: 'failed' as TranslationStatus,
        processing_lease_id: null,
        processing_lease_expires_at: null,
        error_message: error instanceof Error ? error.message : 'Bulk translation failed',
        processing_time_ms: Date.now() - startTime,
        completed_at: new Date(),
      },
    });
    if (failed.count !== 1) {
      const current = await prisma.translationJob.findUnique({
        where: { id: data.jobId },
        select: { status: true, translation: true },
      });
      if (current?.status === 'cancelled' || current?.status === 'completed') {
        return { status: current.status, translation: current.translation || undefined };
      }
    }
    throw error;
  } finally {
    stopLeaseHeartbeat();
  }
});
}

if (process.env.NODE_ENV !== 'test') {
  registerProcessors();
}

// ---------------------------------------------------------------------------
// Orphan Dispute Reconciler — runs every hour (at :17 to avoid stampede
// with other hourly crons). The named jobId makes Bull dedupe this repeatable
// across worker restarts: re-adding the same cron+jobId is idempotent.
// ---------------------------------------------------------------------------
translationQueue.process('reconcile-orphan-disputes', async () => {
  const summary = await reconcilePermanentOrphanDisputes();
  logger.info('Orphan dispute reconciliation pass complete', summary);
  return summary;
});

translationQueue
  .add(
    'reconcile-orphan-disputes',
    {},
    {
      repeat: { cron: '17 * * * *' },
      jobId: 'reconcile-orphan-disputes-cron',
    }
  )
  .catch((err) => logger.error('Failed to schedule orphan dispute reconciler', { error: err }));

// Queue event handlers
translationQueue.on('completed', (job) => {
  logger.info('Job completed', { jobId: job.id });
});

translationQueue.on('failed', (job, error) => {
  logger.error('Job failed', { jobId: job?.id, error: error.message });
});

translationQueue.on('error', (error) => {
  logger.error('Queue error', { error: error.message });
});

// Graceful shutdown
let shutdownStarted = false;

async function gracefulShutdown(signal: string) {
  if (shutdownStarted) {
    return;
  }
  shutdownStarted = true;
  logger.info(`Received ${signal}, shutting down worker...`);

  try {
    await shutdownWorkerRuntime({
      closeQueue: () => translationQueue.close(),
      stopHeartbeat: () => workerHeartbeat.stop(),
      closeConfigSubscription: () => configSubscription?.close() ?? Promise.resolve(),
      disconnectDatabase: () => prisma.$disconnect(),
      stopOutboxDrainer: async () => {
        await stopBulkQueueDispatcher();
        await stopCompletionOutboxDrainer();
      },
    });
    logger.info('Worker shut down successfully');
    process.exit(0);
  } catch (error) {
    logger.error('Worker shutdown completed with errors', { error });
    process.exit(1);
  }
}

if (process.env.NODE_ENV !== 'test') {
  process.on('SIGTERM', () => gracefulShutdown('SIGTERM'));
  process.on('SIGINT', () => gracefulShutdown('SIGINT'));

  startWorkerRuntime({
    initializeConfig: initializeConfigFromDatabase,
    waitForQueue: () => translationQueue.isReady(),
    startOutboxDrainer: async () => {
      await startCompletionOutboxDrainer();
      await startBulkQueueDispatcher();
    },
    startHeartbeat: () => {
      configSubscription = startConfigSubscription({
        workerId,
        categories: ['gemini', 'pricing', 'credits', 'rateLimits'],
      });
      workerHeartbeat.start(5000, (error) => {
        logger.error('Worker heartbeat failed', { error });
      });
    },
    onConfigFallback: (error) => {
      useEnvironmentConfigFallback();
      logger.error('Worker config init from database failed, using explicit .env fallback', { error });
    },
  })
    .then(() => {
      logger.info('Translation worker started', {
        redis: getRedisUrl().replace(/:[^:]*@/, ':***@'),
        workerId,
      });
    })
    .catch((error) => {
      logger.error('Translation worker startup failed', { error });
      process.exit(1);
    });
}

export { translationQueue };
