#!/usr/bin/env python3
"""
Lead import queue processor.

This script is designed to be triggered every minute via cron. On each run it:
    * pulls up to MAX_LEADS_PER_RUN entries from the Redis queue,
    * removes them to avoid reprocessing,
    * processes them concurrently using threads,
    * logs success/error details for every lead.
"""

from __future__ import annotations

import json
import logging
import os
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
from typing import Any, Dict, List, Sequence

try:
    import redis
    from redis.exceptions import RedisError
except ImportError:  # pragma: no cover - redis is optional during development
    redis = None  # type: ignore[assignment]

    class RedisError(Exception):
        """Fallback RedisError used when redis-py is unavailable."""

        pass

MAX_LEADS_PER_RUN = 20
LOG_PATH = Path(__file__).resolve().parents[1] / "storage" / "logs" / "lead_import_queue.log"
LEAD_QUEUE_KEY = os.getenv("LEAD_IMPORT_QUEUE_KEY", "lead_queue")

POP_LEADS_SCRIPT = """
local limit = tonumber(ARGV[1])
if (not limit) or limit <= 0 then
    return {}
end
local entries = redis.call('LRANGE', KEYS[1], 0, limit - 1)
local consumed = #entries
if consumed > 0 then
    redis.call('LTRIM', KEYS[1], consumed, -1)
end
return entries
"""



def configure_logger() -> logging.Logger:
    """Configure a stream + file logger for the queue processor."""
    LOG_PATH.parent.mkdir(parents=True, exist_ok=True)
    logger = logging.getLogger("lead_queue_processor")
    logger.setLevel(logging.INFO)

    if not logger.handlers:
        formatter = logging.Formatter("%(asctime)s [%(levelname)s] %(message)s")
        stream_handler = logging.StreamHandler()
        stream_handler.setFormatter(formatter)
        file_handler = logging.FileHandler(LOG_PATH, encoding="utf-8")
        file_handler.setFormatter(formatter)
        logger.addHandler(stream_handler)
        logger.addHandler(file_handler)

    return logger


LOGGER = configure_logger()


def get_redis_client() -> "redis.Redis[Any]":
    """Instantiate a Redis client using env vars or the default localhost."""
    if redis is None:  # pragma: no cover - runtime guard
        raise RuntimeError(
            "O pacote 'redis' não está instalado. Execute `pip install redis` para usar a fila real."
        )

    redis_url = os.getenv("REDIS_URL")
    if redis_url:
        return redis.from_url(redis_url, decode_responses=True)

    host = os.getenv("REDIS_HOST", "127.0.0.1")
    port = int(os.getenv("REDIS_PORT", "6379"))
    db = int(os.getenv("REDIS_DB", "0"))
    password = os.getenv("REDIS_PASSWORD")
    return redis.Redis(host=host, port=port, db=db, password=password, decode_responses=True)


def parse_queue_entries(entries: Sequence[Any]) -> List[Dict[str, Any]]:
    """Convert Redis queue entries into dictionaries, ignoring malformed payloads."""
    leads: List[Dict[str, Any]] = []
    for raw in entries:
        if raw is None:
            continue
        if isinstance(raw, bytes):
            value = raw.decode("utf-8", errors="ignore")
        else:
            value = str(raw)
        try:
            payload = json.loads(value)
        except json.JSONDecodeError:
            LOGGER.warning("Entrada inválida removida da fila %s: %s", LEAD_QUEUE_KEY, value)
            continue

        if isinstance(payload, dict):
            leads.append(payload)
            continue

        LOGGER.warning(
            "Entrada ignorada na fila %s (esperado dict, recebido %s).",
            LEAD_QUEUE_KEY,
            type(payload).__name__,
        )
    return leads


def fetch_leads_from_queue(limit: int) -> List[Dict[str, Any]]:
    """Pop up to `limit` leads from Redis, removing them atomically from the list."""
    if limit <= 0:
        return []

    client = get_redis_client()
    try:
        entries = client.eval(POP_LEADS_SCRIPT, 1, LEAD_QUEUE_KEY, limit)
    except RedisError as exc:
        LOGGER.error("Falha ao consumir fila %s: %s", LEAD_QUEUE_KEY, exc)
        return []

    return parse_queue_entries(entries)


class LeadQueue:
    """Thread-safe in-memory queue abstraction."""

    def __init__(self, initial_data: List[Dict[str, Any]] | None = None) -> None:
        self._lock = threading.Lock()
        self._items: List[Dict[str, Any]] = list(initial_data or [])

    def get_batch(self, limit: int) -> List[Dict[str, Any]]:
        """Pop up to `limit` leads from the queue."""
        with self._lock:
            batch = self._items[:limit]
            self._items = self._items[limit:]
        return batch


def processar_lead(lead: Dict[str, Any]) -> None:
    """
    Simulate importing a single lead.

    Raises:
        ValueError: when required data is missing (simulated failure path).
    """
    lead_id = lead.get("id")
    print(f"Processando lead #{lead_id} - {lead.get('name')}")
    time.sleep(1)  # Simulates IO/CPU work

    if not lead.get("phone"):
        raise ValueError("Telefone obrigatório não informado.")

    # Insert real import logic here (API calls, DB inserts, etc.).


def process_batch(leads: List[Dict[str, Any]]) -> None:
    """Process a list of leads concurrently with worker threads."""
    if not leads:
        LOGGER.info("Nenhum lead pendente para processar nesta execução.")
        return

    max_workers = min(8, len(leads))
    LOGGER.info("Iniciando processamento de %d lead(s) com %d thread(s).", len(leads), max_workers)

    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        future_to_lead = {executor.submit(processar_lead, lead): lead for lead in leads}
        for future in as_completed(future_to_lead):
            lead = future_to_lead[future]
            lead_id = lead.get("id")
            try:
                future.result()
                LOGGER.info("Lead #%s processado com sucesso.", lead_id)
            except Exception as exc:  # noqa: BLE001 - we must keep processing remaining leads
                LOGGER.error("Lead #%s falhou: %s", lead_id, exc)


def main() -> None:
    """Entry point for cron execution."""
    leads = fetch_leads_from_queue(MAX_LEADS_PER_RUN)
    queue = LeadQueue(leads)
    batch = queue.get_batch(MAX_LEADS_PER_RUN)

    if not batch:
        LOGGER.info("Fila %s vazia - nada a fazer.", LEAD_QUEUE_KEY)
        return

    LOGGER.info(
        "Execução iniciada - %d lead(s) consumidos da fila %s.",
        len(batch),
        LEAD_QUEUE_KEY,
    )
    process_batch(batch)
    LOGGER.info("Execução concluída.")


if __name__ == "__main__":
    main()
