/** Deterministic Generate All scale proof: one manifest fans out to max-20 backend jobs. */
import type { Job } from 'bull';
import type { Prisma } from '@prisma/client';

const LIVE = /(api\.press\.zone|generativelanguage\.googleapis\.com|aiplatform\.googleapis\.com)/i;
const TEST_KEY = /^(?:test|fake|mock|fixture|dummy|placeholder)[-_]/i;
function rejectLiveConfig(): void {
  if (process.env.NODE_ENV !== 'test') throw new Error('This suite requires NODE_ENV=test');
  for (const key of ['API_URL', 'BACKEND_URL', 'TRANSLATION_API_URL', 'GEMINI_API_URL', 'GOOGLE_GEMINI_API_URL']) {
    if (LIVE.test(process.env[key] || '')) throw new Error(`Live translation URL is forbidden: ${key}`);
  }
  for (const key of ['GEMINI_API_KEY', 'GOOGLE_API_KEY', 'GOOGLE_GENERATIVE_AI_API_KEY']) {
    const value = process.env[key];
    if (value && !TEST_KEY.test(value)) throw new Error(`Live provider key is forbidden: ${key}`);
  }
}
rejectLiveConfig();

type FixtureJob = Pick<Job<BulkContentTranslationJobData>, 'data' | 'opts' | 'attemptsMade'>;
type BulkProcessor = (job: FixtureJob) => Promise<unknown>;
const processors = new Map<string, BulkProcessor>();
const queueAdd = jest.fn<Promise<void>, [string, BulkContentTranslationJobData, { jobId: string }]>().mockResolvedValue(undefined);
const queueProcess = jest.fn((name: string, concurrencyOrProcessor: number | BulkProcessor, maybeProcessor?: BulkProcessor) => {
  const processor = typeof concurrencyOrProcessor === 'function' ? concurrencyOrProcessor : maybeProcessor;
  if (!processor) throw new Error('processor fixture is not callable');
  processors.set(name, processor);
});
const drainBulkQueue = jest.fn();
const translateStructured = jest.fn();
const deliverWebhook = jest.fn();
const lockProcessingJob = jest.fn();
const loggerError = jest.fn();

jest.mock('../../queue', () => ({ translationQueue: {
  process: queueProcess, add: queueAdd, on: jest.fn(), close: jest.fn(), isReady: jest.fn().mockResolvedValue(true),
} }));
jest.mock('../../services/geminiClient', () => ({ geminiClient: { translate: jest.fn(), translateBulk: jest.fn(), translateStructured } }));
jest.mock('../../services/webhookService', () => ({ deliverWebhook }));
jest.mock('../../utils/logger', () => ({ logger: { info: jest.fn(), warn: jest.fn(), error: loggerError } }));
jest.mock('../../utils/metrics', () => ({ trackTranslationJob: jest.fn(), trackTokensProcessed: jest.fn() }));
jest.mock('../../workers/orphanDisputeReconciler', () => ({ reconcilePermanentOrphanDisputes: jest.fn() }));
jest.mock('../../queue/bulkQueueDispatcher', () => ({ configureBulkQueueDispatcher: jest.fn(), drainBulkQueue, startBulkQueueDispatcher: jest.fn(), stopBulkQueueDispatcher: jest.fn() }));
jest.mock('../../workerLifecycle', () => ({ startWorkerRuntime: jest.fn(), shutdownWorkerRuntime: jest.fn() }));
jest.mock('../../workerLock', () => ({ JobCancellationWonError: class extends Error {}, lockProcessingJob }));
jest.mock('../../workerHeartbeat', () => ({ getWorkerHeartbeatKey: jest.fn(() => 'worker:fixture'), WorkerHeartbeat: class { start(): void {} stop(): Promise<void> { return Promise.resolve(); } } }));
jest.mock('../../config', () => ({ config: { redisHost: 'fixture.invalid', redisPort: 6379, redisPassword: undefined, redisDb: 0 }, getRedisUrl: jest.fn(() => 'redis://fixture.invalid'), initializeConfigFromDatabase: jest.fn(), startConfigSubscription: jest.fn(), useEnvironmentConfigFallback: jest.fn() }));

process.env.WORKER_ID = 'bulk-content-scale-fixture';
process.env.WORKER_GEMINI_CONCURRENCY = '8';

import { prismaMock } from '../setup';
import { bulkContentJobSchema, submitBulkContentJob } from '../../routes/jobs';
import { registerProcessors } from '../../worker';
import { BULK_CONTENT_MAX_ITEMS, Tone } from '../../types';
import type {
  BulkContentQueueItem,
  BulkContentQueuePayload,
  BulkContentTranslationJobData,
  BulkContentWebhookPayload,
} from '../../types';
import { countStructuredCharacters } from '../../utils/structuredFields';

const CHUNK = BULK_CONTENT_MAX_ITEMS;
const MANIFEST = 'generate-all-fixture-manifest';
const CALLBACK = 'https://fixture.invalid/callback';
const SECRET = 'fixture-secret-0123456789';
const itemsFor = (n: number): BulkContentQueueItem[] => Array.from({ length: n }, (_, i) => ({
  ref: `content:${i + 1}:es`, targetLang: 'es', fields: {
    title: `Title ${i + 1}`, excerpt: `Excerpt ${i + 1}`, content: `Body ${i + 1}`,
    'acf.subtitle': `Subtitle ${i + 1}`, 'seo.description': `Description ${i + 1}`,
  },
}));
const submissionId = (i: number): string => `00000000-0000-4000-8000-${String(i).padStart(12, '0')}`;
const deterministic = (fields: Record<string, string>, lang: string): Record<string, string> =>
  Object.fromEntries(Object.entries(fields).map(([key, value]) => [key, `fixture:${lang}:${key}:${value}`]));

type StoredJob = {
  id: string; status: 'pending' | 'processing' | 'completed'; submissionId: string; clientJobId: string;
  secret: string; payload: BulkContentQueuePayload; db: Record<string, unknown>; queuedAt?: Date; translation?: string; characters?: number;
};
type OutboxJob = { callback_url: string | null; callback_secret: string | null };
type Outbox = {
  jobId: string;
  event: 'bulk_content_translation.completed';
  payload: BulkContentWebhookPayload;
  status: 'pending' | 'delivering' | 'delivered';
  attempts: number;
  claimedAt?: Date;
  job: OutboxJob;
};
let jobs: Map<string, StoredJob>;
let outbox: Outbox[];
let bull: Array<{ jobId: string; data: BulkContentTranslationJobData }>;
let nextJob: number;
let activeCalls: number;
let maxActiveCalls: number;

type OutboxUpdateManyArgs = {
  where: { status?: Outbox['status']; attempts?: { lt: number }; claimed_at?: Date };
  data: { status?: Outbox['status']; claimed_at?: Date };
};
type OutboxCreateArgs = { data: { job_id: string; event: Outbox['event']; payload: BulkContentWebhookPayload } };
type PersistedResult = { ref: string; status: 'completed' | 'failed'; fields?: Record<string, string>; error?: string };

function isRecord(value: unknown): value is Record<string, unknown> {
  return typeof value === 'object' && value !== null && !Array.isArray(value);
}
function isQueueItem(value: unknown): value is BulkContentQueueItem & Prisma.JsonObject {
  return isRecord(value) && typeof value.ref === 'string' && typeof value.targetLang === 'string' &&
    isRecord(value.fields) && Object.values(value.fields).every((field) => typeof field === 'string');
}
function isBulkContentQueuePayload(value: unknown): value is BulkContentQueuePayload & Prisma.JsonObject {
  return isRecord(value) && value.type === 'bulk-content' && typeof value.jobId === 'string' &&
    typeof value.submissionId === 'string' && typeof value.clientJobId === 'string' &&
    typeof value.userId === 'string' && typeof value.plugin === 'string' && Array.isArray(value.items) &&
    value.items.every(isQueueItem) && typeof value.sourceLang === 'string' && typeof value.tone === 'string' &&
    typeof value.callbackUrl === 'string';
}

function resetFakes(): void {
  jobs = new Map(); outbox = []; bull = []; nextJob = 1; activeCalls = 0; maxActiveCalls = 0;
  processors.clear(); queueProcess.mockClear(); queueAdd.mockReset(); queueAdd.mockResolvedValue(undefined);
  drainBulkQueue.mockReset(); translateStructured.mockReset(); loggerError.mockReset(); lockProcessingJob.mockReset(); lockProcessingJob.mockResolvedValue(undefined);
  deliverWebhook.mockReset(); deliverWebhook.mockImplementation(async (_id: string, _payload: unknown, url: string) => {
    if (LIVE.test(url)) throw new Error('live callback URL in deterministic fixture');
    return { success: true };
  });
  const webhookOutbox = {
    findMany: jest.fn(async () => outbox.filter((row) => row.status === 'pending').map((row) => ({
      ...row,
      id: `${row.jobId}:outbox`,
      job: { callback_url: CALLBACK, callback_secret: SECRET },
    }))),
    updateMany: jest.fn(async (args: OutboxUpdateManyArgs) => {
      if (args.where?.status === 'pending' && args.where.attempts?.lt !== undefined) {
        const row = outbox.find((candidate) => candidate.status === 'pending');
        if (!row) return { count: 0 }; row.status = 'delivering'; row.attempts += 1; row.claimedAt = args.data.claimed_at; return { count: 1 };
      }
      if (args.where?.status === 'delivering' && args.where.claimed_at) {
        const row = outbox.find((candidate) => candidate.status === 'delivering');
        if (!row) return { count: 0 }; row.status = args.data.status || 'delivered'; return { count: 1 };
      }
      return { count: 0 };
    }),
    create: jest.fn(async (args: OutboxCreateArgs) => {
      const row: Outbox = { jobId: args.data.job_id, event: args.data.event, payload: args.data.payload, status: 'pending', attempts: 0, job: { callback_url: CALLBACK, callback_secret: SECRET } };
      outbox.push(row); return row;
    }),
  };
  Object.assign(prismaMock, { webhookOutbox });
  prismaMock.translationJob.findUnique.mockImplementation((args) => {
    let result: Record<string, unknown> | StoredJob | null;
    if (args.where.submission_id) {
      const job = [...jobs.values()].find((candidate) => candidate.submissionId === args.where.submission_id);
      result = job ? { ...job.db, status: job.status, translation: job.translation } : null;
    } else {
      result = args.where.id ? jobs.get(args.where.id) || null : null;
    }
    return Promise.resolve(result) as never;
  });
  prismaMock.translationJob.create.mockImplementation(({ data }) => {
    const id = typeof data.id === 'string' ? data.id : `fixture-job-${nextJob++}`;
    if (typeof data.submission_id !== 'string' || typeof data.client_job_id !== 'string' ||
        typeof data.callback_secret !== 'string' || !isBulkContentQueuePayload(data.queue_payload)) {
      throw new Error('invalid canonical fixture persistence payload');
    }
    const payload = data.queue_payload;
    jobs.set(id, { id, status: 'pending', submissionId: data.submission_id, clientJobId: data.client_job_id, secret: data.callback_secret, payload, db: { ...data, id } });
    return Promise.resolve({ ...data, id, status: 'pending' }) as never;
  });
  prismaMock.translationJob.updateMany.mockImplementation((args) => {
    const id = args.where?.id;
    const job = typeof id === 'string' ? jobs.get(id) : undefined;
    let count = 0;
    if (job && Array.isArray(args.where?.OR) && job.status === 'pending') {
      job.status = 'processing'; count = 1;
    }
    return Promise.resolve({ count }) as never;
  });
  prismaMock.translationJob.update.mockImplementation((args) => {
    const id = args.where.id;
    const job = typeof id === 'string' ? jobs.get(id) : undefined;
    if (!job) throw new Error('missing fixture job');
    const data = args.data;
    if (typeof data.status === 'string') job.status = data.status as StoredJob['status'];
    if (typeof data.translation === 'string') job.translation = data.translation;
    if (typeof data.characters_used === 'number') job.characters = data.characters_used;
    return Promise.resolve(job) as never;
  });
  prismaMock.creditTransaction.findFirst.mockResolvedValue({ balance_after: 1_000_000 } as never);
  prismaMock.translationJob.aggregate.mockResolvedValue({ _sum: { reserved_characters: 0 } } as never);
  prismaMock.creditTransaction.create.mockResolvedValue({} as never); prismaMock.user.update.mockResolvedValue({} as never);
  type TransactionCallback = (tx: typeof prismaMock) => Promise<unknown>;
  const transaction = jest.fn((callback: TransactionCallback) => callback(prismaMock));
  Object.defineProperty(prismaMock, '$transaction', { value: transaction, configurable: true });
  drainBulkQueue.mockImplementation(async (jobId: string) => {
    const job = jobs.get(jobId); if (!job) throw new Error('missing fixture queue job');
    // Mirror the production dispatcher's queued_at guard: replay may request a drain,
    // but an already-admitted row must never be added to Bull twice.
    if (job.queuedAt) return;
    job.queuedAt = new Date();
    const data: BulkContentTranslationJobData = {
      jobId: job.payload.jobId,
      submissionId: job.payload.submissionId,
      clientJobId: job.payload.clientJobId,
      userId: job.payload.userId,
      items: job.payload.items,
      sourceLang: job.payload.sourceLang,
      tone: job.payload.tone,
      callbackUrl: job.payload.callbackUrl,
      callbackSecret: job.secret,
    };
    bull.push({ jobId, data });
    await queueAdd('bulk-content', data, { jobId });
  });
  translateStructured.mockImplementation(async (fields: Record<string, string>, _source: string, lang: string) => {
    activeCalls += 1; maxActiveCalls = Math.max(maxActiveCalls, activeCalls); await Promise.resolve(); activeCalls -= 1;
    const translatedFields = deterministic(fields, lang); return { fields: translatedFields, translatedFields };
  });
}

function worker(): BulkProcessor {
  registerProcessors(); const value = processors.get('bulk-content');
  if (!value) throw new Error('bulk-content processor not registered'); return value;
}

type ManifestStatus = {
  status: 'queued' | 'processing' | 'completed'; manifestId: string; backendJobs: number;
  completedBackendJobs: number; storedCount: number; total: number;
};
function status(itemCount: number): ManifestStatus {
  const completed = [...jobs.values()].filter((job) => job.status === 'completed').length;
  return {
    status: completed === 0 ? 'queued' : completed === jobs.size ? 'completed' : 'processing',
    manifestId: MANIFEST, backendJobs: jobs.size, completedBackendJobs: completed,
    storedCount: [...jobs.values()].reduce((n, job) => n + (job.translation ? JSON.parse(job.translation).length : 0), 0),
    total: itemCount,
  };
}

describe('bulk-content worker canonical submission scale', () => {
  beforeEach(() => { rejectLiveConfig(); resetFakes(); });
  it.each([[100, 5], [1000, 50]])('admits %d items as exactly %d max-20 submissions, and reconciles all results', async (itemCount, expectedBackendJobs) => {
    const items = itemsFor(itemCount);
    const expectedChars = items.reduce((n, item) => n + countStructuredCharacters(item.fields), 0);
    const fixtureChunks = Array.from({ length: expectedBackendJobs }, (_, index) =>
      items.slice(index * CHUNK, (index + 1) * CHUNK)
    );
    expect(fixtureChunks).toHaveLength(expectedBackendJobs);
    expect(fixtureChunks.every((chunk) => chunk.length === BULK_CONTENT_MAX_ITEMS)).toBe(true);
    expect(fixtureChunks.flatMap((chunk) => chunk.map((item) => item.ref))).toEqual(items.map((item) => item.ref));
    let clicks = 0;
    const manifestRequests: Array<{ manifestId: string; items: BulkContentQueueItem[] }> = [];
    const submissions: Array<Parameters<typeof submitBulkContentJob>[0]> = [];
    clicks += 1; manifestRequests.push({ manifestId: MANIFEST, items });
    for (const [chunkOffset, chunk] of fixtureChunks.entries()) {
      const index = chunkOffset + 1;
      const submissionItems = chunk.map((item) => {
        const { title, excerpt, content, ...fields } = item.fields;
        return { ref: item.ref, targetLang: item.targetLang, title, excerpt, content, fields };
      });
      const request = { userId: 'fixture-user', plugin: 'international', siteUrl: 'https://fixture.invalid/site', submissionId: submissionId(index), clientJobId: `${MANIFEST}:${index}`, sourceLang: 'en', items: submissionItems, callbackUrl: CALLBACK, callbackSecret: SECRET, tone: Tone.NEUTRAL };
      submissions.push(request);
      const { userId: _userId, plugin: _plugin, siteUrl: _siteUrl, callbackSecret: _callbackSecret, ...publicRequest } = request;
      expect(bulkContentJobSchema.parse(publicRequest).items).toHaveLength(chunk.length);
      await submitBulkContentJob(request);
    }
    expect(clicks).toBe(1); expect(manifestRequests).toHaveLength(1); expect(manifestRequests[0].items).toHaveLength(itemCount);
    expect(jobs.size).toBe(expectedBackendJobs); expect(bull).toHaveLength(expectedBackendJobs); expect(queueAdd).toHaveBeenCalledTimes(expectedBackendJobs);
    expect(new Set([...jobs.values()].map((job) => job.submissionId)).size).toBe(expectedBackendJobs);
    expect(queueAdd.mock.calls.map((call) => call[1].items.length)).toEqual(Array(expectedBackendJobs).fill(CHUNK));
    expect(queueAdd.mock.calls.every((call) => typeof call[1].submissionId === 'string' && call[1].submissionId.length > 0 && call[1].callbackSecret === SECRET)).toBe(true);
    const persistedPayloads = [...jobs.values()].map((job) => job.payload);
    expect(persistedPayloads.every((payload) => payload.items.every((item) => Object.keys(item.fields).length > 0))).toBe(true);
    expect(persistedPayloads.every((payload) => !JSON.stringify(payload).includes(SECRET) && !('callbackSecret' in payload))).toBe(true);
    for (const request of submissions) await submitBulkContentJob(request);
    expect(jobs.size).toBe(expectedBackendJobs);
    expect(bull).toHaveLength(expectedBackendJobs);
    expect(queueAdd).toHaveBeenCalledTimes(expectedBackendJobs);
    expect(status(itemCount)).toEqual({ status: 'queued', manifestId: MANIFEST, backendJobs: expectedBackendJobs, completedBackendJobs: 0, storedCount: 0, total: itemCount });

    const process = worker(); for (const job of bull) await process({ data: job.data, opts: { attempts: 1 }, attemptsMade: 0 });
    expect(status(itemCount)).toEqual({ status: 'completed', manifestId: MANIFEST, backendJobs: expectedBackendJobs, completedBackendJobs: expectedBackendJobs, storedCount: itemCount, total: itemCount });
    expect(outbox).toHaveLength(expectedBackendJobs); expect(outbox.every((row) => row.event === 'bulk_content_translation.completed' && row.status === 'delivered' && row.payload.failed_count === 0)).toBe(true);
    for (const row of outbox) {
      const job = jobs.get(row.jobId); expect(job).toBeDefined();
      expect(row.payload.submission_id).toBe(job?.submissionId);
      expect(row.payload.event_sequence).toBe(1);
      expect(row.payload.results.map((result) => result.ref)).toEqual(job?.payload.items.map((item) => item.ref));
    }
    expect(outbox.reduce((n, row) => n + row.payload.results.length, 0)).toBe(itemCount); expect(deliverWebhook).toHaveBeenCalledTimes(expectedBackendJobs); expect(loggerError).not.toHaveBeenCalled();
    const results: PersistedResult[] = [...jobs.values()].flatMap((job) => JSON.parse(job.translation || 'null'));
    expect(results).toHaveLength(itemCount); expect(new Set(results.map((row) => row.ref)).size).toBe(itemCount); expect(results.every((row) => row.status === 'completed' && !row.error)).toBe(true);
    expect(results.map((row) => row.ref).sort()).toEqual(items.map((item) => item.ref).sort());
    for (const row of results) {
      const source = items.find((item) => item.ref === row.ref); expect(source).toBeDefined();
      expect(row.fields).toEqual(deterministic(source?.fields || {}, source?.targetLang || 'es'));
      expect(row.fields).toBeDefined();
      expect(Object.values(row.fields || {}).every((value) => value.length > 0)).toBe(true);
    }
    // One provider call/item is current behavior and remains a production scalability risk.
    expect(translateStructured).toHaveBeenCalledTimes(itemCount); expect(maxActiveCalls).toBeGreaterThan(0); expect(maxActiveCalls).toBeLessThanOrEqual(8);
    expect(prismaMock.creditTransaction.create).toHaveBeenCalledTimes(expectedBackendJobs);
    expect(prismaMock.creditTransaction.create.mock.calls.reduce((n, [call]) => n + Number(call.data.amount), 0)).toBe(-expectedChars);

    const callsBeforeReload = translateStructured.mock.calls.length;
    expect(status(itemCount).storedCount).toBe(itemCount);
    for (const job of bull) expect(await process({ data: job.data, opts: { attempts: 1 }, attemptsMade: 0 })).toEqual({ translation: jobs.get(job.jobId)?.translation });
    expect(translateStructured).toHaveBeenCalledTimes(callsBeforeReload); expect(queueAdd).toHaveBeenCalledTimes(expectedBackendJobs); expect(prismaMock.creditTransaction.create).toHaveBeenCalledTimes(expectedBackendJobs); expect(clicks).toBe(1);
  }, 30000);
});
