import { Prisma } from '@prisma/client';
import type { ExceptionRule } from '../services/exceptionService';
import { randomUUID } from 'crypto';
import { translationQueue } from '../queue';
import { logger } from '../utils/logger';

const BULK_QUEUE_LEASE_MS = 120_000;
const BULK_QUEUE_BASE_DELAY_MS = 5_000;
const BULK_QUEUE_MAX_DELAY_MS = 60_000;
const BULK_QUEUE_BATCH_SIZE = 100;

export type BulkContentItem = {
  ref: string;
  targetLang: string;
  title?: string;
  excerpt?: string;
  content?: string;
};

export type BulkQueuePayload =
  | {
    type: 'bulk-strings'; jobId: string; clientJobId: string | null; userId: string; plugin: string;
    strings: Array<{ id: string; content: string }>; sourceLang: string; targetLangs: string[];
    tone: string; exceptions?: ExceptionRule[]; callbackUrl: string; callbackSecret: string;
  }
  | {
    type: 'bulk-content'; jobId: string; clientJobId: string | null; userId: string; plugin: string;
    items: BulkContentItem[]; sourceLang: string;
    tone: string; callbackUrl: string; callbackSecret: string;
  };

// Distributive omit — a plain `Omit` over a union collapses to the key
// intersection and silently drops variant-specific fields (e.g. `items`,
// `strings`). This preserves each member's own shape minus `callbackSecret`.
export type StoredBulkQueuePayload = BulkQueuePayload extends infer T
  ? T extends { callbackSecret: string }
    ? Omit<T, 'callbackSecret'>
    : never
  : never;

export type BulkQueueCandidate = {
  id: string; status: string; queued_at: Date | null; queue_lease_id: string | null;
  queue_lease_expires_at: Date | null; queue_payload: unknown; callback_secret?: string | null;
  created_at: Date;
};

type ClaimedBulkQueueJob = { callbackSecret: string | null };

export type BulkQueueDispatcherDependencies = {
  listCandidates: (_jobId?: string, _cursor?: Pick<BulkQueueCandidate, 'id' | 'created_at'>) => Promise<BulkQueueCandidate[]>;
  claim: (_jobId: string, _leaseId: string, _now: Date) => Promise<ClaimedBulkQueueJob | null>;
  markQueued: (_jobId: string, _leaseId: string) => Promise<boolean>;
  release: (_jobId: string, _leaseId: string) => Promise<void>;
  add: (_name: string, _payload: BulkQueuePayload, _options: { jobId: string }) => Promise<unknown>;
  createLeaseId: () => string; now: () => Date;
};

export type BulkQueueDrainResult = { claimed: number; queued: number; failed: number; skipped: number };

export function createBulkQueueDispatcher(dependencies: BulkQueueDispatcherDependencies) {
  let drainPromise: Promise<BulkQueueDrainResult> | undefined;
  let timer: ReturnType<typeof setTimeout> | undefined;
  let started = false; let stopRequested = false; let delayMs = BULK_QUEUE_BASE_DELAY_MS;

  async function runDrain(jobId?: string): Promise<BulkQueueDrainResult> {
    const result: BulkQueueDrainResult = { claimed: 0, queued: 0, failed: 0, skipped: 0 };
    let cursor: Pick<BulkQueueCandidate, 'id' | 'created_at'> | undefined;
    do {
      const candidates = await dependencies.listCandidates(jobId, cursor);
      for (const job of candidates) {
        cursor = { id: job.id, created_at: job.created_at };
        if (job.status !== 'pending' || job.queued_at) { result.skipped += 1; continue; }
        const payload = job.queue_payload as StoredBulkQueuePayload;
        const leaseId = dependencies.createLeaseId();
        const claimed = await dependencies.claim(job.id, leaseId, dependencies.now());
        if (!claimed?.callbackSecret) { result.skipped += 1; continue; }
        result.claimed += 1;
        try {
          await dependencies.add(payload.type, { ...payload, callbackSecret: claimed.callbackSecret } as BulkQueuePayload, { jobId: job.id });
          if (await dependencies.markQueued(job.id, leaseId)) result.queued += 1;
          else result.skipped += 1;
        } catch (error) {
          await dependencies.release(job.id, leaseId); result.failed += 1;
          logger.error('Bulk queue dispatch failed', { error, jobId: job.id });
        }
      }
      if (jobId || candidates.length < BULK_QUEUE_BATCH_SIZE) break;
    } while (cursor);
    return result;
  }

  async function drain(jobId?: string): Promise<BulkQueueDrainResult> {
    if (!drainPromise) drainPromise = runDrain(jobId).finally(() => { drainPromise = undefined; });
    return drainPromise;
  }
  function schedule(nextDelayMs: number): void {
    if (!started || stopRequested) return;
    timer = setTimeout(() => { timer = undefined; void drain().then((result) => {
      delayMs = result.failed > 0 ? Math.min(delayMs * 2, BULK_QUEUE_MAX_DELAY_MS) : BULK_QUEUE_BASE_DELAY_MS;
    }).catch((error) => { logger.error('Bulk queue recurring dispatch failed', { error }); delayMs = Math.min(delayMs * 2, BULK_QUEUE_MAX_DELAY_MS); }).finally(() => schedule(delayMs)); }, nextDelayMs);
    timer.unref?.();
  }
  return { drain, async start(): Promise<void> { if (started) return; started = true; stopRequested = false; try {
    const result = await drain(); delayMs = result.failed > 0 ? BULK_QUEUE_BASE_DELAY_MS * 2 : BULK_QUEUE_BASE_DELAY_MS;
  } catch (error) { logger.error('Bulk queue startup dispatch failed', { error }); delayMs = BULK_QUEUE_BASE_DELAY_MS * 2; } schedule(delayMs); },
  async stop(): Promise<void> { stopRequested = true; started = false; if (timer) { clearTimeout(timer); timer = undefined; } await drainPromise; } };
}

type TranslationJobDelegate = {
  findMany(_args: object): Promise<BulkQueueCandidate[]>; findFirst(_args: object): Promise<{ callback_secret: string | null } | null>;
  updateMany(_args: object): Promise<{ count: number }>;
};

export function createPrismaBulkQueueDispatcher(translationJob: TranslationJobDelegate) {
  return createBulkQueueDispatcher({
    listCandidates: (jobId, cursor) => translationJob.findMany({ where: { ...(jobId ? { id: jobId } : {}), status: 'pending', queued_at: null, queue_payload: { not: Prisma.DbNull }, OR: [{ queue_lease_expires_at: null }, { queue_lease_expires_at: { lt: new Date() } }] }, select: { id: true, status: true, queued_at: true, queue_lease_id: true, queue_lease_expires_at: true, queue_payload: true, created_at: true }, orderBy: [{ created_at: 'asc' }, { id: 'asc' }], ...(cursor ? { cursor: { id: cursor.id }, skip: 1 } : {}), take: jobId ? 1 : BULK_QUEUE_BATCH_SIZE }),
    async claim(jobId, leaseId, now) {
      const claimed = await translationJob.updateMany({ where: { id: jobId, status: 'pending', queued_at: null, OR: [{ queue_lease_expires_at: null }, { queue_lease_expires_at: { lt: now } }] }, data: { queue_lease_id: leaseId, queue_lease_expires_at: new Date(now.getTime() + BULK_QUEUE_LEASE_MS) } });
      if (claimed.count !== 1) return null;
      const job = await translationJob.findFirst({ where: { id: jobId, status: 'pending', queued_at: null, queue_lease_id: leaseId }, select: { callback_secret: true } });
      return job ? { callbackSecret: job.callback_secret } : null;
    },
    async markQueued(jobId, leaseId) { const result = await translationJob.updateMany({ where: { id: jobId, status: 'pending', queued_at: null, queue_lease_id: leaseId }, data: { queued_at: new Date(), queue_lease_id: null, queue_lease_expires_at: null } }); return result.count === 1; },
    async release(jobId, leaseId) { await translationJob.updateMany({ where: { id: jobId, status: 'pending', queued_at: null, queue_lease_id: leaseId }, data: { queue_lease_id: null, queue_lease_expires_at: null } }); },
    add: (name, payload, options) => translationQueue.add(name, payload, options), createLeaseId: randomUUID, now: () => new Date(),
  });
}

let dispatcher: ReturnType<typeof createBulkQueueDispatcher> | undefined;
export function configureBulkQueueDispatcher(translationJob: TranslationJobDelegate): void { dispatcher = createPrismaBulkQueueDispatcher(translationJob); }
function requireDispatcher() { if (!dispatcher) throw new Error('Bulk queue dispatcher is not configured for this runtime'); return dispatcher; }
export const drainBulkQueue = (jobId?: string) => requireDispatcher().drain(jobId);
export const startBulkQueueDispatcher = () => requireDispatcher().start();
export const stopBulkQueueDispatcher = () => dispatcher?.stop() ?? Promise.resolve();
