import { randomUUID } from 'node:crypto';
import { Worker } from 'node:worker_threads';
import type {
  PerceptionBoundary,
  PerceptionEventProposal,
  PerceptionFrame,
  PerceptionSpatialFacts,
} from './contracts.js';
import type { RuntimeDiagnosticSeverity } from '../../../framework/plugins/index.js';
import type {
  PerceptionThreadLoaderConfig,
  PerceptionThreadModuleSpec,
} from './perception-module-loader.js';

const WORKER_BOOTSTRAP_URL = new URL(
  '../../../framework/modules/runtime-module-worker-bootstrap.mjs',
  import.meta.url,
);
const WORKER_HOST_URL = new URL('./perception-reducer-thread-worker.ts', import.meta.url).href;
const MAX_REDUCER_WORKERS = 10;
const START_TIMEOUT_MS = 20_000;
const STOP_RESET_TIMEOUT_MS = 1_000;

export type PerceptionReducerResetReason =
  Exclude<PerceptionBoundary['kind'], 'events_requested'> | 'replay';
export type PerceptionReducerOperation = 'update' | 'commit' | 'flush' | 'reset';

export interface PerceptionReducerOperationInput {
  readonly frame?: Readonly<PerceptionFrame>;
  readonly analysis?: unknown;
  readonly boundary?: Readonly<PerceptionBoundary>;
  readonly reason?: PerceptionReducerResetReason;
}

export interface PerceptionReducerThreadPoolOptions {
  readonly modules: readonly PerceptionThreadModuleSpec[];
  readonly loader: PerceptionThreadLoaderConfig;
  readonly spatialFacts: PerceptionSpatialFacts;
  readonly maxWorkers?: number;
  readonly maxOldGenerationSizeMb?: number;
  readonly diagnosticsEnabled?: boolean;
  readonly diagnostic?: (
    name: string,
    attributes?: Readonly<Record<string, unknown>>,
    severity?: RuntimeDiagnosticSeverity,
  ) => void;
}

export interface PerceptionReducerThreadPoolSnapshot {
  readonly state: 'new' | 'starting' | 'running' | 'stopping' | 'stopped';
  readonly maxWorkerCount: number;
  readonly workerCount: number;
  readonly busyWorkerCount: number;
  readonly queuedRequestCount: number;
  readonly inFlightRequestCount: number;
  readonly activeModuleCount: number;
  readonly unavailableModuleIds: readonly string[];
  readonly restartedWorkerCount: number;
  readonly workerHeapUsedBytes: number;
  readonly peakWorkerHeapUsedBytes: number;
}

interface PendingRequest {
  readonly requestId: string;
  readonly moduleId: string;
  readonly operation: PerceptionReducerOperation;
  readonly input: PerceptionReducerOperationInput;
  readonly signal: AbortSignal;
  readonly resolve: (events: readonly PerceptionEventProposal[]) => void;
  readonly reject: (error: Error) => void;
  readonly onAbort: () => void;
  settled: boolean;
}

interface ReducerLane {
  readonly id: number;
  readonly modules: readonly PerceptionThreadModuleSpec[];
  readonly queue: PendingRequest[];
  readonly activeModuleIds: Set<string>;
  worker?: Worker;
  state: 'starting' | 'idle' | 'busy' | 'restarting' | 'stopped';
  current?: PendingRequest;
  readyResolve?: () => void;
  readyReject?: (error: Error) => void;
  heapUsedBytes?: number;
  generation: number;
  consecutiveWorkerFailures: number;
  restartTask?: Promise<void>;
}

interface WorkerMessageBase {
  readonly type: string;
  readonly poolId: string;
  readonly laneId: number;
  readonly requestId?: string;
  readonly moduleId?: string;
  readonly heapUsedBytes?: number;
}

interface WorkerReadyMessage extends WorkerMessageBase {
  readonly type: 'ready';
  readonly moduleIds: readonly string[];
  readonly startupFailures: readonly {
    readonly moduleId: string;
    readonly code: string;
  }[];
  readonly debugSnapshots: Readonly<Record<string, unknown>>;
}

interface WorkerResultMessage extends WorkerMessageBase {
  readonly type: 'result';
  readonly requestId: string;
  readonly moduleId: string;
  readonly events: readonly PerceptionEventProposal[];
  readonly debugSnapshot: unknown;
}

interface WorkerErrorMessage extends WorkerMessageBase {
  readonly type: 'request_error' | 'error';
  readonly code: string;
  readonly debugSnapshot?: unknown;
}

function asError(value: unknown, fallback: string): Error {
  return value instanceof Error ? value : new Error(typeof value === 'string' ? value : fallback);
}

function abortError(signal: AbortSignal): Error {
  return asError(signal.reason, 'perception_reducer_request_aborted');
}

export class PerceptionReducerThreadPool {
  private readonly poolId = randomUUID();
  private readonly maxWorkers: number;
  private readonly maxOldGenerationSizeMb: number;
  private readonly lanes: ReducerLane[];
  private readonly laneByModule = new Map<string, ReducerLane>();
  private readonly moduleFailures = new Map<string, Error>();
  private readonly debugSnapshots = new Map<string, unknown>();
  private readonly releasedModuleIds = new Set<string>();
  private state: PerceptionReducerThreadPoolSnapshot['state'] = 'new';
  private startTask?: Promise<void>;
  private stopTask?: Promise<void>;
  private peakWorkerHeapUsedBytes = 0;
  private restartedWorkerCount = 0;
  private diagnosticCallback: PerceptionReducerThreadPoolOptions['diagnostic'];

  constructor(private readonly options: PerceptionReducerThreadPoolOptions) {
    this.diagnosticCallback = options.diagnostic;
    const ids = options.modules.map(({ moduleId }) => moduleId);
    if (ids.length === 0) throw new Error('perception_reducer_modules_empty');
    if (new Set(ids).size !== ids.length) throw new Error('perception_reducer_module_duplicate');
    this.maxWorkers = Math.max(1, Math.min(
      MAX_REDUCER_WORKERS,
      Math.floor(options.maxWorkers ?? MAX_REDUCER_WORKERS),
    ));
    this.maxOldGenerationSizeMb = Math.max(16, Math.min(
      256,
      Math.floor(options.maxOldGenerationSizeMb ?? 64),
    ));
    const laneCount = Math.min(this.maxWorkers, options.modules.length);
    const groups = Array.from({ length: laneCount }, () => [] as PerceptionThreadModuleSpec[]);
    options.modules.forEach((module, index) => groups[index % laneCount]!.push(module));
    this.lanes = groups.map((modules, index) => ({
      id: index + 1,
      modules,
      queue: [],
      activeModuleIds: new Set(),
      state: 'stopped',
      generation: 0,
      consecutiveWorkerFailures: 0,
    }));
    for (const lane of this.lanes) {
      for (const module of lane.modules) this.laneByModule.set(module.moduleId, lane);
    }
  }

  setDiagnostic(
    diagnostic: PerceptionReducerThreadPoolOptions['diagnostic'],
  ): void {
    if (diagnostic) this.diagnosticCallback = diagnostic;
  }

  start(): Promise<void> {
    if (this.state === 'running') return Promise.resolve();
    if (this.startTask) return this.startTask;
    if (this.stopTask) return this.stopTask.then(() => this.start());
    let task!: Promise<void>;
    task = this.startInternal().finally(() => {
      if (this.startTask === task) this.startTask = undefined;
    });
    this.startTask = task;
    return task;
  }

  private async startInternal(): Promise<void> {
    this.state = 'starting';
    this.moduleFailures.clear();
    this.debugSnapshots.clear();
    this.releasedModuleIds.clear();
    for (const lane of this.lanes) {
      lane.consecutiveWorkerFailures = 0;
      lane.heapUsedBytes = undefined;
      lane.activeModuleIds.clear();
    }
    const starts = await Promise.allSettled(this.lanes.map((lane) => this.spawn(lane)));
    if (this.state !== 'starting') throw new Error('perception_reducer_pool_stopped');
    for (const [index, result] of starts.entries()) {
      if (result.status === 'fulfilled') continue;
      const lane = this.lanes[index]!;
      lane.state = 'stopped';
      lane.activeModuleIds.clear();
      const failure = asError(result.reason, 'perception_reducer_worker_start_failed');
      for (const { moduleId } of lane.modules) this.moduleFailures.set(moduleId, failure);
      const worker = lane.worker;
      lane.worker = undefined;
      await worker?.terminate().catch(() => -1);
    }
    this.state = 'running';
    this.diagnostic('perception.reducer_pool.started', {
      worker_count: this.lanes.filter(({ worker }) => !!worker).length,
      max_worker_count: this.maxWorkers,
      active_module_count: this.activeModuleCount(),
      unavailable_module_ids: [...this.moduleFailures.keys()],
      worker_max_old_generation_mb: this.maxOldGenerationSizeMb,
    }, 'info');
    this.pumpAll();
  }

  assertAvailable(moduleId: string): void {
    const lane = this.laneByModule.get(moduleId);
    if (!lane || !lane.activeModuleIds.has(moduleId) || !lane.worker) {
      throw this.moduleFailures.get(moduleId)
        ?? new Error(`perception_reducer_unavailable:${moduleId}`);
    }
  }

  operate(
    moduleId: string,
    operation: PerceptionReducerOperation,
    input: PerceptionReducerOperationInput,
    signal: AbortSignal,
  ): Promise<readonly PerceptionEventProposal[]> {
    if (this.state !== 'running') {
      return Promise.reject(new Error('perception_reducer_pool_not_running'));
    }
    if (signal.aborted) return Promise.reject(abortError(signal));
    const lane = this.laneByModule.get(moduleId);
    if (!lane || !lane.activeModuleIds.has(moduleId) || !lane.worker) {
      return Promise.reject(this.moduleFailures.get(moduleId)
        ?? new Error(`perception_reducer_unavailable:${moduleId}`));
    }
    return new Promise((resolve, reject) => {
      const request: PendingRequest = {
        requestId: randomUUID(),
        moduleId,
        operation,
        input,
        signal,
        resolve,
        reject,
        settled: false,
        onAbort: () => this.abortRequest(lane, request),
      };
      signal.addEventListener('abort', request.onAbort, { once: true });
      lane.queue.push(request);
      this.pump(lane);
    });
  }

  async releaseModule(
    moduleId: string,
    reason: 'segment_ending' | 'shutdown',
  ): Promise<void> {
    if (this.releasedModuleIds.has(moduleId)) return;
    if (this.state === 'running') {
      const controller = new AbortController();
      const timer = setTimeout(
        () => controller.abort(new Error('perception_reducer_stop_timeout')),
        STOP_RESET_TIMEOUT_MS,
      );
      timer.unref?.();
      try {
        await this.operate(moduleId, 'reset', { reason }, controller.signal);
      } catch {
        // The worker is terminated below when the last runner releases it.
      } finally {
        clearTimeout(timer);
      }
    }
    this.releasedModuleIds.add(moduleId);
    if (this.releasedModuleIds.size === this.options.modules.length) await this.stop();
  }

  stop(): Promise<void> {
    if (this.state === 'stopped') return Promise.resolve();
    if (this.stopTask) return this.stopTask;
    let task!: Promise<void>;
    task = this.stopInternal().finally(() => {
      if (this.stopTask === task) this.stopTask = undefined;
    });
    this.stopTask = task;
    return task;
  }

  private async stopInternal(): Promise<void> {
    this.state = 'stopping';
    const error = new Error('perception_reducer_pool_stopped');
    await Promise.all(this.lanes.map(async (lane) => {
      lane.readyReject?.(error);
      if (lane.current) this.rejectRequest(lane.current, error);
      for (const request of [...lane.queue]) this.rejectRequest(request, error);
      lane.queue.length = 0;
      lane.current = undefined;
      lane.activeModuleIds.clear();
      lane.state = 'stopped';
      lane.generation += 1;
      const worker = lane.worker;
      lane.worker = undefined;
      await worker?.terminate().catch(() => -1);
    }));
    this.state = 'stopped';
    this.moduleFailures.clear();
    this.debugSnapshots.clear();
    this.diagnostic('perception.reducer_pool.stopped', {
      worker_count: 0,
      peak_worker_heap_used_bytes: this.peakWorkerHeapUsedBytes,
    }, 'info');
  }

  debugSnapshot(moduleId: string): unknown {
    return structuredClone(this.debugSnapshots.get(moduleId) ?? null);
  }

  snapshot(): PerceptionReducerThreadPoolSnapshot {
    const queuedRequestCount = this.lanes.reduce((sum, lane) => sum + lane.queue.length, 0);
    const busyWorkerCount = this.lanes.filter(({ state }) => state === 'busy').length;
    return {
      state: this.state,
      maxWorkerCount: this.maxWorkers,
      workerCount: this.lanes.filter(({ worker }) => !!worker).length,
      busyWorkerCount,
      queuedRequestCount,
      inFlightRequestCount: queuedRequestCount + busyWorkerCount,
      activeModuleCount: this.activeModuleCount(),
      unavailableModuleIds: [...this.moduleFailures.keys()].sort(),
      restartedWorkerCount: this.restartedWorkerCount,
      workerHeapUsedBytes: this.currentWorkerHeapUsedBytes(),
      peakWorkerHeapUsedBytes: this.peakWorkerHeapUsedBytes,
    };
  }

  private spawn(lane: ReducerLane): Promise<void> {
    lane.generation += 1;
    const generation = lane.generation;
    lane.state = 'starting';
    lane.activeModuleIds.clear();
    const worker = new Worker(WORKER_BOOTSTRAP_URL, {
      execArgv: [],
      resourceLimits: { maxOldGenerationSizeMb: this.maxOldGenerationSizeMb },
      stdout: true,
      stderr: true,
      workerData: {
        runtimeModuleBootstrap: {
          hostUrl: WORKER_HOST_URL,
          hooksUrl: this.options.loader.hooksUrl,
          sdkAliases: this.options.loader.sdkAliases,
          esmModuleRoots: [...new Set(lane.modules.map((module) => (
            new URL('.', module.sourceUrl).href
          )))],
        },
        poolId: this.poolId,
        laneId: lane.id,
        modules: lane.modules,
        spatialFacts: this.options.spatialFacts,
        diagnosticsEnabled: this.options.diagnosticsEnabled === true,
      },
    });
    lane.worker = worker;
    worker.stdout?.on('data', (chunk: Buffer) => this.diagnostic(
      'perception.reducer_worker.stdout',
      { lane_id: lane.id, module_ids: lane.modules.map(({ moduleId }) => moduleId), bytes: chunk.length },
    ));
    worker.stderr?.on('data', (chunk: Buffer) => this.diagnostic(
      'perception.reducer_worker.stderr',
      { lane_id: lane.id, module_ids: lane.modules.map(({ moduleId }) => moduleId), bytes: chunk.length },
    ));
    worker.on('message', (message: unknown) => {
      if (lane.worker === worker && lane.generation === generation) this.handleMessage(lane, message);
    });
    worker.once('error', (error) => {
      if (lane.worker === worker && lane.generation === generation) this.failLane(lane, error);
    });
    worker.once('exit', (code) => {
      if (lane.worker === worker && lane.generation === generation && lane.state !== 'stopped') {
        this.failLane(lane, new Error(`perception_reducer_worker_exited:${code}`));
      }
    });
    return new Promise<void>((resolve, reject) => {
      let timer!: ReturnType<typeof setTimeout>;
      const finishResolve = (): void => {
        clearTimeout(timer);
        lane.readyResolve = undefined;
        lane.readyReject = undefined;
        resolve();
      };
      const finishReject = (error: Error): void => {
        clearTimeout(timer);
        lane.readyResolve = undefined;
        lane.readyReject = undefined;
        reject(error);
      };
      lane.readyResolve = finishResolve;
      lane.readyReject = finishReject;
      timer = setTimeout(
        () => finishReject(new Error('perception_reducer_worker_start_timeout')),
        START_TIMEOUT_MS,
      );
      timer.unref?.();
    });
  }

  private handleMessage(lane: ReducerLane, value: unknown): void {
    if (!value || typeof value !== 'object' || Array.isArray(value)) {
      this.failLane(lane, new Error('perception_reducer_protocol_invalid'));
      return;
    }
    const message = value as WorkerReadyMessage | WorkerResultMessage | WorkerErrorMessage;
    if (message.poolId !== this.poolId || message.laneId !== lane.id) {
      this.failLane(lane, new Error('perception_reducer_identity_mismatch'));
      return;
    }
    if (this.options.diagnosticsEnabled === true && (
      !Number.isFinite(message.heapUsedBytes)
      || Number(message.heapUsedBytes) < 0
    )) {
      this.failLane(lane, new Error('perception_reducer_worker_memory_invalid'));
      return;
    }
    if (typeof message.heapUsedBytes === 'number') {
      lane.heapUsedBytes = message.heapUsedBytes;
      this.peakWorkerHeapUsedBytes = Math.max(
        this.peakWorkerHeapUsedBytes,
        this.currentWorkerHeapUsedBytes(),
      );
    }
    if (message.type === 'ready') {
      if (lane.state !== 'starting' || !Array.isArray(message.moduleIds)) {
        this.failLane(lane, new Error('perception_reducer_ready_invalid'));
        return;
      }
      const expected = new Set(lane.modules.map(({ moduleId }) => moduleId));
      if (message.moduleIds.some((moduleId) => !expected.has(moduleId))) {
        this.failLane(lane, new Error('perception_reducer_ready_invalid'));
        return;
      }
      lane.activeModuleIds.clear();
      for (const moduleId of message.moduleIds) {
        lane.activeModuleIds.add(moduleId);
        this.moduleFailures.delete(moduleId);
        this.debugSnapshots.set(moduleId, structuredClone(message.debugSnapshots[moduleId] ?? null));
      }
      for (const failure of message.startupFailures ?? []) {
        if (!expected.has(failure.moduleId)) continue;
        this.moduleFailures.set(failure.moduleId, new Error(failure.code));
        this.debugSnapshots.delete(failure.moduleId);
      }
      for (const moduleId of expected) {
        if (!lane.activeModuleIds.has(moduleId) && !this.moduleFailures.has(moduleId)) {
          this.moduleFailures.set(moduleId, new Error('perception_reducer_worker_module_missing'));
        }
      }
      lane.state = lane.activeModuleIds.size > 0 ? 'idle' : 'stopped';
      lane.readyResolve?.();
      if (lane.state === 'stopped') {
        const worker = lane.worker;
        lane.worker = undefined;
        void worker?.terminate().catch(() => -1);
      }
      this.pump(lane);
      return;
    }
    const request = lane.current;
    if (!request || message.requestId !== request.requestId || message.moduleId !== request.moduleId) {
      this.failLane(lane, new Error('perception_reducer_request_mismatch'));
      return;
    }
    lane.current = undefined;
    lane.state = 'idle';
    if ('debugSnapshot' in message) {
      this.debugSnapshots.set(request.moduleId, structuredClone(message.debugSnapshot ?? null));
    }
    if (message.type === 'result') {
      lane.consecutiveWorkerFailures = 0;
      this.resolveRequest(request, message.events);
    } else if (message.type === 'request_error') {
      this.rejectRequest(request, new Error(message.code));
    } else {
      this.rejectRequest(request, new Error(message.code ?? 'perception_reducer_worker_failed'));
      this.restartLane(lane, 'worker_error');
      return;
    }
    this.pump(lane);
  }

  private pumpAll(): void {
    for (const lane of this.lanes) this.pump(lane);
  }

  private pump(lane: ReducerLane): void {
    if (this.state !== 'running' || lane.state !== 'idle' || lane.current) return;
    while (lane.queue.length > 0) {
      const request = lane.queue.shift()!;
      if (request.settled || request.signal.aborted) {
        if (!request.settled) this.rejectRequest(request, abortError(request.signal));
        continue;
      }
      if (!lane.activeModuleIds.has(request.moduleId)) {
        this.rejectRequest(request, this.moduleFailures.get(request.moduleId)
          ?? new Error(`perception_reducer_unavailable:${request.moduleId}`));
        continue;
      }
      lane.current = request;
      lane.state = 'busy';
      try {
        lane.worker!.postMessage({
          type: 'operate',
          poolId: this.poolId,
          requestId: request.requestId,
          moduleId: request.moduleId,
          operation: request.operation,
          ...request.input,
        });
      } catch (error) {
        this.rejectRequest(request, asError(error, 'perception_reducer_submit_failed'));
        this.restartLane(lane, 'submit_failed');
      }
      return;
    }
  }

  private abortRequest(lane: ReducerLane, request: PendingRequest): void {
    if (request.settled) return;
    const queued = lane.queue.indexOf(request);
    if (queued >= 0) lane.queue.splice(queued, 1);
    const active = lane.current === request;
    this.rejectRequest(request, abortError(request.signal));
    if (active) this.restartLane(lane, 'request_aborted');
  }

  private resolveRequest(
    request: PendingRequest,
    events: readonly PerceptionEventProposal[],
  ): void {
    if (request.settled) return;
    request.settled = true;
    request.signal.removeEventListener('abort', request.onAbort);
    request.resolve(structuredClone(events));
  }

  private rejectRequest(request: PendingRequest, error: Error): void {
    if (request.settled) return;
    request.settled = true;
    request.signal.removeEventListener('abort', request.onAbort);
    request.reject(error);
  }

  private failLane(lane: ReducerLane, error: Error): void {
    lane.readyReject?.(error);
    if (lane.current) {
      const current = lane.current;
      lane.current = undefined;
      this.rejectRequest(current, error);
    }
    if (this.state !== 'running') return;
    this.restartLane(lane, 'worker_failed', error);
  }

  private restartLane(lane: ReducerLane, reason: string, failure?: Error): void {
    if (lane.restartTask || this.state !== 'running') return;
    lane.consecutiveWorkerFailures += 1;
    for (const moduleId of lane.activeModuleIds) this.debugSnapshots.delete(moduleId);
    lane.activeModuleIds.clear();
    const stateLost = failure ?? new Error(`perception_reducer_state_lost:${reason}`);
    for (const request of [...lane.queue]) this.rejectRequest(request, stateLost);
    lane.queue.length = 0;
    if (lane.consecutiveWorkerFailures >= 3) {
      const worker = lane.worker;
      lane.worker = undefined;
      lane.current = undefined;
      lane.state = 'stopped';
      lane.generation += 1;
      void worker?.terminate().catch(() => -1);
      for (const { moduleId } of lane.modules) {
        this.moduleFailures.set(
          moduleId,
          new Error('perception_reducer_disabled_after_worker_failures'),
        );
      }
      this.diagnostic('perception.module.disabled', {
        module_ids: lane.modules.map(({ moduleId }) => moduleId),
        reason: 'reducer_worker_failures',
        consecutive_failures: lane.consecutiveWorkerFailures,
      }, 'error');
      return;
    }
    const worker = lane.worker;
    lane.worker = undefined;
    lane.current = undefined;
    lane.state = 'restarting';
    lane.generation += 1;
    lane.restartTask = Promise.resolve(worker?.terminate()).catch(() => -1).then(async () => {
      if (this.state !== 'running') return;
      await this.spawn(lane);
      this.restartedWorkerCount += 1;
      this.diagnostic('perception.reducer_worker.restarted', {
        lane_id: lane.id,
        module_ids: lane.modules.map(({ moduleId }) => moduleId),
        reason,
      }, 'warn');
    }).catch((error) => {
      lane.state = 'stopped';
      for (const { moduleId } of lane.modules) {
        this.moduleFailures.set(
          moduleId,
          asError(error, 'perception_reducer_restart_failed'),
        );
      }
      this.diagnostic('perception.reducer_worker.replacement_failed', {
        lane_id: lane.id,
        module_ids: lane.modules.map(({ moduleId }) => moduleId),
        reason,
        error_code: asError(error, 'perception_reducer_restart_failed').message.split(':', 1)[0],
      }, 'error');
    }).finally(() => {
      lane.restartTask = undefined;
      this.pump(lane);
    });
  }

  private activeModuleCount(): number {
    return this.lanes.reduce((sum, lane) => sum + lane.activeModuleIds.size, 0);
  }

  private currentWorkerHeapUsedBytes(): number {
    return this.lanes.reduce((sum, lane) => sum + (lane.heapUsedBytes ?? 0), 0);
  }

  private diagnostic(
    name: string,
    attributes?: Readonly<Record<string, unknown>>,
    severity?: RuntimeDiagnosticSeverity,
  ): void {
    try {
      this.diagnosticCallback?.(name, attributes, severity);
    } catch {
      // Diagnostics never alter reducer execution or lifecycle.
    }
  }
}
