const mockProcessors = new Map<string, (job: any) => Promise<unknown>>();
const mockQueueProcess = jest.fn((name: string, ...rest: unknown[]) => {
  const processor = rest[rest.length - 1] as (job: any) => Promise<unknown>;
  mockProcessors.set(name, processor);
});
const mockTranslateStructured = jest.fn();
const mockIsRetryableError = jest.fn();
const mockOutbox = {
  findMany: jest.fn().mockResolvedValue([]),
  updateMany: jest.fn().mockResolvedValue({ count: 0 }),
  create: jest.fn().mockResolvedValue({}),
};

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

process.env.WORKER_ID = 'bulk-content-concurrency-test-worker';
process.env.WORKER_GEMINI_CONCURRENCY = '2';
process.env.WORKER_GEMINI_RETRY_BASE_DELAY_MS = '1';
process.env.WORKER_GEMINI_RETRY_MAX_DELAY_MS = '1';

import { prismaMock } from '../setup';
import { registerProcessors } from '../../worker';
import { Tone } from '../../types';
import type { BulkContentItemInput, BulkContentTranslationJobData } from '../../types';

function processor(): (job: any) => Promise<any> {
  registerProcessors();
  const value = mockProcessors.get('bulk-content');
  if (!value) throw new Error('bulk-content processor was not registered');
  return value;
}

function makeJobData(jobId: string, items: BulkContentItemInput[]): BulkContentTranslationJobData {
  return {
    jobId,
    clientJobId: `${jobId}-client`,
    userId: 'user-1',
    items,
    sourceLang: 'en',
    tone: Tone.NEUTRAL,
    callbackUrl: 'https://site.example.test/callback',
    callbackSecret: 'x'.repeat(16),
  };
}

function translatedFields(targetLang: string, fields: Record<string, string>) {
  return Object.fromEntries(Object.keys(fields).map((key) => [key, `${targetLang}:${key}`]));
}

describe('bulk-content worker concurrency', () => {
  beforeEach(() => {
    mockProcessors.clear();
    mockQueueProcess.mockClear();
    mockTranslateStructured.mockReset();
    mockIsRetryableError.mockReset();
    mockOutbox.findMany.mockResolvedValue([]);
    mockOutbox.updateMany.mockResolvedValue({ count: 0 });
    mockOutbox.create.mockClear();
    Object.assign(prismaMock, { webhookOutbox: mockOutbox });
    (prismaMock.$transaction as jest.Mock).mockImplementation(
      async (callback: (tx: typeof prismaMock) => Promise<unknown>) => callback(prismaMock)
    );
    prismaMock.translationJob.updateMany.mockResolvedValue({ count: 1 } as never);
    prismaMock.creditTransaction.findFirst.mockResolvedValue({ balance_after: 1000 } as never);
    prismaMock.creditTransaction.create.mockResolvedValue({} as never);
    prismaMock.user.update.mockResolvedValue({} as never);
    prismaMock.translationJob.update.mockResolvedValue({} as never);
    mockIsRetryableError.mockReturnValue(false);
  });

  it('translates all items concurrently while preserving result order and exact billing', async () => {
    let activeCalls = 0;
    let maxActiveCalls = 0;
    mockTranslateStructured.mockImplementation(async (fields: Record<string, string>, _source: string, target: string) => {
      activeCalls += 1;
      maxActiveCalls = Math.max(maxActiveCalls, activeCalls);
      await new Promise((resolve) => setTimeout(resolve, target === 'es' ? 20 : 2));
      activeCalls -= 1;
      const output = translatedFields(target, fields);
      return { fields: output, translatedFields: output };
    });

    const items: BulkContentItemInput[] = [
      { ref: 'one', targetLang: 'es', title: 'A', content: 'one' },
      { ref: 'two', targetLang: 'fr', title: 'BB', excerpt: '<em>two</em>', content: 'three' },
      { ref: 'three', targetLang: 'de', title: 'CCC', content: 'six' },
      { ref: 'four', targetLang: 'it', title: 'DDDD', content: 'seven' },
    ];

    const result = await processor()({
      data: makeJobData('job-order', items),
      opts: { attempts: 1 },
      attemptsMade: 0,
    });

    expect(mockTranslateStructured).toHaveBeenCalledTimes(items.length);
    expect(maxActiveCalls).toBeLessThanOrEqual(2);
    const payload = mockOutbox.create.mock.calls[0][0].data.payload;
    expect(payload.results.map((item: { ref: string }) => item.ref)).toEqual(['one', 'two', 'three', 'four']);
    expect(payload.results.every((item: { status: string }) => item.status === 'completed')).toBe(true);
    expect(payload.total_characters_used).toBe(29);
    expect(result.totalCharactersUsed).toBe(29);
    expect(prismaMock.creditTransaction.create.mock.calls[0][0].data.amount).toBe(-29);
  });

  it('settles a failed item without aborting later items or changing result order', async () => {
    mockTranslateStructured.mockImplementation(async (fields: Record<string, string>, _source: string, target: string) => {
      if (target === 'fr') {
        throw new Error('non-retryable item failure');
      }
      const output = translatedFields(target, fields);
      return { fields: output, translatedFields: output };
    });

    const items: BulkContentItemInput[] = [
      { ref: 'first', targetLang: 'es', title: 'one', content: 'body' },
      { ref: 'failed', targetLang: 'fr', title: 'two', content: 'body' },
      { ref: 'last', targetLang: 'de', title: 'three', content: 'body' },
    ];

    const result = await processor()({
      data: makeJobData('job-failure', items),
      opts: { attempts: 1 },
      attemptsMade: 0,
    });

    const payload = mockOutbox.create.mock.calls[0][0].data.payload;
    expect(payload.results).toMatchObject([
      { ref: 'first', status: 'completed' },
      { ref: 'failed', status: 'failed', error: 'non-retryable item failure' },
      { ref: 'last', status: 'completed' },
    ]);
    expect(payload.failed_count).toBe(1);
    expect(payload.total_characters_used).toBe(16);
    expect(result.totalFailedCount).toBe(1);
    expect(mockTranslateStructured).toHaveBeenCalledTimes(3);
  });

  it('retries retryable provider errors with bounded attempts', async () => {
    mockIsRetryableError.mockReturnValue(true);
    mockTranslateStructured
      .mockRejectedValueOnce(new Error('rate limit'))
      .mockImplementationOnce(async (fields: Record<string, string>) => {
        const output = translatedFields('es', fields);
        return { fields: output, translatedFields: output };
      });

    await processor()({
      data: makeJobData('job-retry', [{ ref: 'retry', targetLang: 'es', title: 'one', content: 'body' }]),
      opts: { attempts: 1 },
      attemptsMade: 0,
    });

    expect(mockTranslateStructured).toHaveBeenCalledTimes(2);
    expect(mockIsRetryableError).toHaveBeenCalledWith(expect.any(Error));
  });

  it('keeps the provider call ceiling process-wide across concurrent jobs', async () => {
    let activeCalls = 0;
    let maxActiveCalls = 0;
    mockTranslateStructured.mockImplementation(async (fields: Record<string, string>, _source: string, target: string) => {
      activeCalls += 1;
      maxActiveCalls = Math.max(maxActiveCalls, activeCalls);
      await new Promise((resolve) => setTimeout(resolve, 8));
      activeCalls -= 1;
      const output = translatedFields(target, fields);
      return { fields: output, translatedFields: output };
    });

    const items = Array.from({ length: 4 }, (_, index) => ({
      ref: `item-${index}`,
      targetLang: 'es',
      title: `Title ${index}`,
      content: `Body ${index}`,
    }));
    const process = processor();

    await Promise.all([
      process({ data: makeJobData('job-a', items), opts: { attempts: 1 }, attemptsMade: 0 }),
      process({ data: makeJobData('job-b', items), opts: { attempts: 1 }, attemptsMade: 0 }),
    ]);

    expect(maxActiveCalls).toBeLessThanOrEqual(2);
    expect(mockTranslateStructured).toHaveBeenCalledTimes(8);
  });
});
