import { Logger, OnModuleDestroy, OnModuleInit } from '@nestjs/common';
import { JwtService } from '@nestjs/jwt';
import { OnGatewayConnection, SubscribeMessage, WebSocketGateway, WebSocketServer } from '@nestjs/websockets';
import { io as ioClient, Socket } from 'socket.io-client';
import { Server, Socket as ServerSocket } from 'socket.io';
import { PrismaService } from 'src/prisma/prisma.service';
import { envs } from 'src/config';

function getMicroUrl(): string {
  return (envs.botMarangatuBaseUrl ?? 'http://localhost:8500').replace(/\/$/, '');
}

// El microservicio manda un snapshot al recibir `worker_public.subscribe`, pero no
// vuelve a mandarlo solo cuando cambia el estado del worker. Sin este refresco el
// panel congelaba lo que hubiera al arrancar el backend (p.ej. seguía diciendo
// "pausado — Login falló por credenciales inválidas" con el worker ya desactivado).
const REFRESCO_SNAPSHOT_MS = 60_000;
// Piso entre refrescos pedidos por un cliente, para que abrir varias pestañas no
// bombardee al microservicio.
const MIN_MS_ENTRE_REFRESCOS = 5_000;

const TIPOS_REGISTRO = ['COMPRA', 'VENTA'] as const;
type TipoRegistro = (typeof TIPOS_REGISTRO)[number];

function normalizeRuc(ruc?: string | null): string | null {
  const base = ruc?.split('-')[0]?.trim();
  return base && /^\d+$/.test(base) ? base : null;
}

// cors: origin reflejado (no el wildcard '*') + credentials. La barrera real de
// seguridad es el JWT que se valida en el handshake (handleConnection).
@WebSocketGateway({ namespace: '/marangatu', cors: { origin: true, credentials: true } })
export class MarangatuGateway implements OnModuleInit, OnModuleDestroy, OnGatewayConnection {
  @WebSocketServer()
  server: Server;

  private readonly logger = new Logger(MarangatuGateway.name);
  private active = false;
  // Una conexión al microservicio por `${empresaRuc}:${tipoRegistro}`.
  private readonly conexiones = new Map<string, Socket>();
  // Último snapshot conocido por `${empresaRuc}:${tipo}`. Se reenvía SOLO a los sockets
  // de esa misma empresa cuando se conectan, para que el panel refleje el estado actual.
  private readonly lastSnapshotByTipo = new Map<string, any>();
  private refrescoTimer: NodeJS.Timeout | null = null;
  private ultimoRefresco = 0;

  constructor(
    private readonly prisma: PrismaService,
    private readonly jwtService: JwtService,
  ) {}

  private roomFor(ruc: string): string {
    return `empresa:${ruc}`;
  }

  private extractToken(client: ServerSocket): string | null {
    const auth = client.handshake?.auth as { token?: string } | undefined;
    const fromAuth = auth?.token;
    if (fromAuth) return fromAuth.replace(/^Bearer\s+/i, '');
    const header = client.handshake?.headers?.authorization;
    if (header) return header.replace(/^Bearer\s+/i, '');
    const q = client.handshake?.query?.token;
    if (typeof q === 'string' && q) return q;
    return null;
  }

  // Handshake autenticado: valida el JWT, resuelve el RUC de la empresa activa del
  // usuario y une el socket al room de esa empresa. Sin token válido → se desconecta.
  async handleConnection(client: ServerSocket) {
    const token = this.extractToken(client);
    if (!token) {
      client.disconnect(true);
      return;
    }
    let empresaRuc: string | null = null;
    try {
      const payload = await this.jwtService.verifyAsync<{ sub: string }>(token, {
        secret: envs.jwtSecret,
      });
      const user = await this.prisma.usuario.findUnique({
        where: { id: payload.sub },
        select: { empresa_id: true, empresa_id_actual: true },
      });
      const empresaId = user?.empresa_id_actual || user?.empresa_id;
      if (empresaId) {
        const empresa = await this.prisma.empresas.findUnique({
          where: { id: empresaId },
          select: { ruc: true },
        });
        empresaRuc = normalizeRuc(empresa?.ruc);
      }
    } catch {
      empresaRuc = null;
    }

    if (!empresaRuc) {
      client.disconnect(true);
      return;
    }

    (client.data as { empresaRuc?: string }).empresaRuc = empresaRuc;
    await client.join(this.roomFor(empresaRuc));

    // Reenviar solo los snapshots de ESTA empresa.
    for (const snapshot of this.lastSnapshotByTipo.values()) {
      if (snapshot?.empresaRuc === empresaRuc) client.emit('marangatu_event', snapshot);
    }
  }

  onModuleDestroy() {
    if (this.refrescoTimer) clearInterval(this.refrescoTimer);
    this.refrescoTimer = null;
    for (const socket of this.conexiones.values()) socket.disconnect();
    this.conexiones.clear();
  }

  /**
   * Vuelve a suscribirse para que el microservicio mande un snapshot fresco. Es la
   * única forma de enterarse de un cambio de estado (el bot no lo avisa solo).
   */
  private async refrescarSnapshots(motivo: string) {
    if (!this.active) return;
    const ahora = Date.now();
    if (ahora - this.ultimoRefresco < MIN_MS_ENTRE_REFRESCOS) return;
    this.ultimoRefresco = ahora;
    this.logger.debug(`Refrescando snapshots Marangatu (${motivo})`);
    await this.subscribeAll();
  }

  /** El panel lo pide al abrirse y con el botón de refrescar. */
  @SubscribeMessage('solicitar_snapshot')
  async handleSolicitarSnapshot() {
    await this.refrescarSnapshots('pedido del panel');
  }

  async onModuleInit() {
    const count = await this.prisma.marangatu_empresa_config.count();
    if (count > 0) {
      this.logger.log(`${count} empresa(s) con config Marangatu — iniciando conexión Socket.IO`);
      this.active = true;
      await this.reconciliarConexiones();
      this.iniciarRefrescoPeriodico();
    } else {
      this.logger.log('Ninguna empresa con config Marangatu — Socket.IO en espera');
    }
  }

  /** Llamar desde MarangatuConfigService al guardar la primera config */
  activateWs() {
    if (this.active) return;
    this.active = true;
    this.logger.log('Primera config Marangatu guardada — activando Socket.IO');
    void this.reconciliarConexiones().then(() => this.iniciarRefrescoPeriodico());
  }

  /** Agregar suscripción cuando se registra una nueva empresa */
  async resubscribeAll() {
    if (!this.active) return;
    await this.reconciliarConexiones();
  }

  deactivateWs() {
    this.active = false;
    for (const socket of this.conexiones.values()) socket.disconnect();
    this.conexiones.clear();
    if (this.refrescoTimer) clearInterval(this.refrescoTimer);
    this.refrescoTimer = null;
    this.logger.log('Socket.IO Marangatu desactivado');
  }

  private iniciarRefrescoPeriodico() {
    if (this.refrescoTimer) return;
    this.refrescoTimer = setInterval(() => {
      void this.refrescarSnapshots('refresco periódico');
    }, REFRESCO_SNAPSHOT_MS);
  }

  /**
   * Una conexión por empresa y tipo: el microservicio hace `leave` de los rooms
   * públicos anteriores en cada `worker_public.subscribe`, así que un solo socket
   * sólo recibe los eventos de la ÚLTIMA suscripción. Con un socket por par, cada
   * uno queda en su room y sabemos de qué empresa y tipo es lo que llega (no hace
   * falta adivinar por workflowId).
   */
  private async reconciliarConexiones() {
    if (!this.active) return;

    const configs = await this.prisma.marangatu_empresa_config.findMany({
      include: { empresa: { select: { ruc: true } } },
    });

    const esperadas = new Set<string>();
    for (const cfg of configs) {
      // El bot exige empresaRuc SOLO dígitos, sin DV (regex ^\d+$).
      const ruc = normalizeRuc(cfg.empresa?.ruc);
      if (!ruc) continue;
      for (const tipoRegistro of TIPOS_REGISTRO) {
        const clave = `${ruc}:${tipoRegistro}`;
        esperadas.add(clave);
        if (!this.conexiones.has(clave)) this.conectarPar(ruc, tipoRegistro, clave);
      }
    }

    // Empresas que ya no tienen config: cerrar su conexión.
    for (const [clave, socket] of this.conexiones) {
      if (!esperadas.has(clave)) {
        socket.disconnect();
        this.conexiones.delete(clave);
      }
    }
  }

  private conectarPar(empresaRuc: string, tipoRegistro: TipoRegistro, clave: string) {
    const socket = ioClient(`${getMicroUrl()}/cdc-worker-public`, {
      transports: ['websocket'],
      reconnection: true,
      reconnectionDelay: 10_000,
      reconnectionAttempts: Infinity,
    });
    this.conexiones.set(clave, socket);

    socket.on('connect', () => {
      this.logger.log(`Conectado al microservicio Marangatu (${clave})`);
      socket.emit('worker_public.subscribe', { empresaRuc, tipoRegistro });
    });

    socket.on('disconnect', (reason) => {
      this.logger.warn(`Desconectado del microservicio Marangatu (${clave}): ${reason}`);
    });

    socket.on('connect_error', (err) => {
      this.logger.error(`Error de conexión Marangatu (${clave}): ${err.message}`);
    });

    // Eventos en tiempo real del worker. El bot manda `data` = PublicProgressEvent:
    //   { at, workflowId, stage, stageLabel, level, message, jobStatus, meta: {...} }
    socket.on('worker_public.event', (payload) => {
      try {
        const data = payload?.data ?? payload;
        const meta = data?.meta ?? {};
        // La empresa y el tipo salen de ESTA conexión, no del payload: el evento en
        // vivo no trae empresaRuc y antes se descartaba si el workflowId era nuevo.
        this.server.to(this.roomFor(empresaRuc)).emit('marangatu_event', {
          tipo: meta.tipoRegistro ?? data.tipoRegistro ?? data.tipo ?? tipoRegistro,
          // "Estado" = estado del job (QUEUED/active/FAILED/completed); "Etapa" = stage legible.
          estado: data.jobStatus ?? data.stage ?? data.estado,
          etapa: data.stageLabel ?? data.stage ?? data.etapa,
          mensaje: data.message ?? data.mensaje ?? data.stageLabel ?? data.stage,
          nivel: data.level ?? data.nivel ?? 'INFO',
          totalRecuperados: meta.total ?? meta.processed ?? data.totalRecuperados,
          workflowId: data.workflowId,
          at: data.at,
          raw: data,
        });
      } catch {
        // ignore
      }
    });

    // Snapshot del estado actual, que el bot manda SOLO como respuesta a
    // `worker_public.subscribe` (por eso hay que volver a pedirlo para refrescar).
    // Campos ANIDADOS:
    //   data.progress.{ workflowId, jobStatus, currentStage, currentStageLabel, nextCapture, nextEnrich, events }
    //   data.worker.{ operationalStatus, operationalReason, credencialEstado, ... }
    socket.on('worker_public.snapshot', (payload) => {
      try {
        if (payload?.status === 'error') {
          this.logger.warn(`Snapshot Marangatu (${clave}) con error: ${payload?.message ?? 'desconocido'}`);
          return;
        }
        const data = payload?.data ?? payload;
        const progress = data?.progress ?? {};
        const worker = data?.worker ?? {};
        // totalRecuperados no viene suelto en el snapshot: lo derivamos del último evento con meta.total.
        const events = Array.isArray(progress.events) ? progress.events : [];
        const lastConTotal = [...events].reverse().find((e) => e?.meta?.total != null);
        const evento = {
          tipo: data.tipoRegistro ?? data.tipo ?? tipoRegistro,
          empresaRuc,
          estado: progress.jobStatus ?? progress.currentStage ?? null,
          etapa: progress.currentStageLabel ?? progress.currentStage ?? null,
          totalRecuperados: lastConTotal?.meta?.total ?? null,
          // nextCapture/nextEnrich son objetos { status, reason, at, remainingSec }. El motivo del
          // pausado (p.ej. "Login Marangatu falló por credenciales inválidas") llega en .reason.
          proximaAutoCaptura: progress.nextCapture ?? null,
          proximoEnrich: progress.nextEnrich ?? null,
          workflowId: progress.workflowId ?? null,
          // Info extra del worker (por si la UI la quiere mostrar): estado operativo y credenciales.
          operationalStatus: worker.operationalStatus ?? null,
          operationalReason: worker.operationalReason ?? null,
          credencialEstado: worker.credencialEstado ?? null,
          isSnapshot: true,
          // Momento en que llegó: el panel lo usa para avisar si el estado quedó viejo.
          at: new Date().toISOString(),
          raw: data,
        };
        // Guardar el último snapshot por empresa+tipo para reenviarlo al conectar (handleConnection).
        this.lastSnapshotByTipo.set(clave, evento);
        this.server.to(this.roomFor(empresaRuc)).emit('marangatu_event', evento);
      } catch {
        // ignore
      }
    });
  }

  /** Vuelve a suscribir cada conexión para que el bot mande un snapshot fresco. */
  private async subscribeAll() {
    for (const [clave, socket] of this.conexiones) {
      if (!socket.connected) continue;
      const [empresaRuc, tipoRegistro] = clave.split(':');
      socket.emit('worker_public.subscribe', { empresaRuc, tipoRegistro });
    }
  }
}
