import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Logger } from '@nestjs/common';
import type { Job } from 'bullmq';
import { LimiteTokensExcedidoError } from 'src/ai-usage/ai-usage.service';
import { PrismaService } from 'src/prisma/prisma.service';
import { IMAGENES_IA_QUEUE, type ImagenesIaJob } from './busquedas.service';
import { LOTE_NORMALIZACION } from './pipeline/normalizador-queries';
import type { ProductoParaBuscar } from './pipeline/tipos';
import { PipelineProductoService } from './pipeline-producto.service';
import { StorageSugeridasService } from './storage-sugeridas.service';

const CONCURRENCIA_PRODUCTOS = 3;

/**
 * Worker de una búsqueda: normaliza por lotes, procesa productos con concurrencia 3 y
 * persiste sugerencias. Reanudable: en un reintento salta los productos que ya tienen
 * sugerencias de esta búsqueda (las queries salen de cache, así que cuesta casi nada).
 */
@Processor(IMAGENES_IA_QUEUE)
export class ImagenesIaProcessor extends WorkerHost {
  private readonly logger = new Logger(ImagenesIaProcessor.name);

  constructor(
    private readonly prisma: PrismaService,
    private readonly pipeline: PipelineProductoService,
    private readonly storage: StorageSugeridasService,
  ) {
    super();
  }

  async process(job: Job<ImagenesIaJob>): Promise<void> {
    const b = await this.prisma.producto_imagen_busquedas.findUnique({ where: { id: job.data.busqueda_id } });
    if (!b || b.estado === 'CANCELADA' || b.estado === 'COMPLETADA') return;
    const ctx = { empresaId: b.empresa_id, busquedaId: b.id, usuarioId: b.usuario_id };
    await this.prisma.producto_imagen_busquedas.update({
      where: { id: b.id },
      data: { estado: 'EN_PROCESO', iniciado_en: b.iniciado_en ?? new Date(), error: null, updated_at: new Date() },
    });

    try {
      const yaHechos = new Set(
        (
          await this.prisma.producto_imagenes_sugeridas.findMany({ where: { busqueda_id: b.id }, select: { producto_id: true }, distinct: ['producto_id'] })
        ).map((s) => s.producto_id),
      );
      const productos = await this.cargarProductos(b.producto_ids.filter((id) => !yaHechos.has(id)));
      let procesados = b.procesados;
      let conSugerencias = b.con_sugerencias;
      let queries = b.queries_busqueda;

      for (let i = 0; i < productos.length; i += LOTE_NORMALIZACION) {
        if (await this.cancelada(b.id)) return;
        const lote = productos.slice(i, i + LOTE_NORMALIZACION);
        const normalizados = await this.pipeline.normalizarLote(ctx, lote);

        await this.pool(lote, CONCURRENCIA_PRODUCTOS, async (p) => {
          if (await this.cancelada(b.id)) return;
          const norm = normalizados.get(p.id)!;
          try {
            const { sugerencias, queriesNuevas } = await this.pipeline.procesarProducto(ctx, p, norm);
            queries += queriesNuevas;
            for (const s of sugerencias) {
              const subida = await this.storage.subirSugerida(b.empresa_id, p.id, s.buffer);
              await this.prisma.producto_imagenes_sugeridas.create({
                data: {
                  empresa_id: b.empresa_id,
                  producto_id: p.id,
                  busqueda_id: b.id,
                  path: subida.path,
                  fuente_url: s.fuente_url.slice(0, 2000),
                  fuente_dominio: s.fuente_dominio?.slice(0, 255) || null,
                  fuente_titulo: s.fuente_titulo?.slice(0, 500) || null,
                  ancho: s.ancho,
                  alto: s.alto,
                  coincide: s.vision.coincide,
                  fondo_blanco: s.vision.fondo_blanco,
                  marca_agua: s.vision.marca_agua,
                  calidad: s.vision.calidad,
                  score: s.score,
                  nota: s.vision.nota ?? null,
                },
              });
            }
            if (sugerencias.length) conSugerencias++;
          } catch (e) {
            if (e instanceof LimiteTokensExcedidoError) throw e;
            this.logger.warn(`Producto ${p.id} sin sugerencias: ${(e as Error).message}`);
          }
          procesados++;
          await this.actualizarProgreso(b.id, { procesados, con_sugerencias: conSugerencias, queries_busqueda: queries });
        });
      }

      if (await this.cancelada(b.id)) return;
      await this.actualizarProgreso(b.id, {
        procesados,
        con_sugerencias: conSugerencias,
        queries_busqueda: queries,
        estado: 'COMPLETADA',
        finalizado_en: new Date(),
      });
    } catch (e) {
      const msg = e instanceof LimiteTokensExcedidoError ? e.message : `Error procesando la búsqueda: ${(e as Error).message}`;
      this.logger.error(`Búsqueda ${b.id} en ERROR: ${msg}`);
      await this.actualizarProgreso(b.id, { estado: 'ERROR', error: msg.slice(0, 2000), finalizado_en: new Date() });
      // Sin cuota no tiene sentido reintentar; cualquier otro error deja que BullMQ reintente (attempts: 2).
      if (!(e instanceof LimiteTokensExcedidoError)) throw e;
    }
  }

  private async cancelada(id: string): Promise<boolean> {
    const b = await this.prisma.producto_imagen_busquedas.findUnique({ where: { id }, select: { estado: true } });
    return b?.estado === 'CANCELADA';
  }

  /** Recalcula tokens/costo desde ai_consumo (fuente de verdad) y persiste el progreso. */
  private async actualizarProgreso(id: string, data: Record<string, unknown>) {
    const agg = await this.prisma.ai_consumo.aggregate({ where: { referencia_id: id }, _sum: { tokens_in: true, tokens_out: true, costo_usd: true } });
    await this.prisma.producto_imagen_busquedas.update({
      where: { id },
      data: { ...data, tokens_in: agg._sum.tokens_in ?? 0, tokens_out: agg._sum.tokens_out ?? 0, costo_usd: agg._sum.costo_usd ?? null, updated_at: new Date() },
    });
  }

  private async cargarProductos(ids: string[]): Promise<ProductoParaBuscar[]> {
    if (!ids.length) return [];
    const rows = await this.prisma.productos.findMany({
      where: { id: { in: ids } },
      select: { id: true, cod_producto: true, descripcion: true, marcas: { select: { descripcion: true } }, categoria: { select: { descripcion: true } } },
    });
    const orden = new Map(ids.map((id, i) => [id, i]));
    return rows
      .map((r) => ({ id: r.id, cod_producto: r.cod_producto, descripcion: r.descripcion, marca: r.marcas?.descripcion ?? null, categoria: r.categoria?.descripcion ?? null }))
      .sort((a, b) => orden.get(a.id)! - orden.get(b.id)!);
  }

  private async pool<T>(items: T[], n: number, fn: (t: T) => Promise<void>): Promise<void> {
    let idx = 0;
    await Promise.all(
      Array.from({ length: Math.min(n, items.length) }, async () => {
        while (idx < items.length) await fn(items[idx++]);
      }),
    );
  }
}
