#!/usr/bin/env python3
"""Servicio de paralelizacion de disponibilidades (solo stdlib, sin dependencias).

Recibe POST /parallel con los mismos parametros de disponibilidad.php,
forma jobs de (MAX 2 ACTIVIDADES x 1 FECHA) y lanza UN worker PHP-CLI
(run_dispo_worker.php) por job, TODOS EN PARALELO (tope MAX_PARALLEL_WORKERS),
fusiona EXACTO por uuid (tipos concatenados en orden de fecha, Rpta/cancel
de la ultima fecha, como el secuencial) y devuelve
{"respuesta": [...], "Rpta": "...", "idsRate": [...]} para que
disponibilidad.php arme el envelope y haga sus UPDATEs como siempre.

Arranque ejemplo:
    nohup python3 /ruta/actividades-completa/dispo_parallel.py >/tmp/dispo_parallel.log 2>&1 &
Salud:
    curl http://127.0.0.1:8766/health
"""
import json
import os
import subprocess
import urllib.parse
from concurrent.futures import ThreadPoolExecutor, as_completed
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer

# ----------------- CONFIG -----------------
HOST = "127.0.0.1"
PORT = 8766
SERVICE_TOKEN = "d1sp0-p4r4ll3l-t0k3n-CAMBIAR-EN-PRODUCCION"
PHP_BIN = "php"
WORKER_PATH = os.environ.get(
    "DISPO_WORKER_PATH",
    os.path.join(os.path.dirname(os.path.abspath(__file__)), "run_dispo_worker.php"),
)
ACTS_POR_HILO = 2        # actividades por worker (1 sola fecha por worker),
# igual para cualquier tarifa (el worker siempre usa datosInternos();
# la tarifa solo selecciona credenciales en firmaAPI).
# NOTA 2026-09-10: medido 57.8s por worker de (2 acts x rango completo),
# ~7s por (act x fecha). Con mas acts x fecha se roza el muro de 45s.
MAX_PARALLEL_WORKERS = 48  # tope de workers simultaneos (0 = sin tope).
# 2026-09-14: 12 acts x 16 dias = 96 jobs; con 25 eran 4 oleadas (>muro 45s).
# Con 48 son 2 oleadas (~20-30s). Si los tiempos EMPEORAN al subir workers
# => throttling de Musement: bajar a 32/25.
WORKER_TIMEOUT = 100     # segundos por grupo
BODY_MAX_BYTES = 65536
# ------------------------------------------


def fechas_rango(desde, hasta):
    """Lista de YYYY-MM-DD inclusivas entre desde y hasta (1 elemento si iguales/vacio)."""
    import datetime as _dt
    try:
        d0 = _dt.date.fromisoformat(str(desde)[:10])
        d1 = _dt.date.fromisoformat(str(hasta)[:10])
    except Exception:
        return [str(desde)[:10]]
    if d1 < d0:
        d0, d1 = d1, d0
    out, d = [], d0
    while d <= d1:
        out.append(d.isoformat())
        d += _dt.timedelta(days=1)
    return out


def run_worker(params, fecha, chunk):
    """Ejecuta el worker PHP para UN GRUPO (<=5 acts) y UNA FECHA.
    Devuelve (dict|None, segundos)."""
    import time as _time
    t0 = _time.time()
    q = dict(params)
    q["activitiesList"] = ",".join(chunk)
    q["fecha_unica"] = fecha
    qs = urllib.parse.urlencode({k: ("" if v is None else str(v)) for k, v in q.items()})
    try:
        p = subprocess.run(
            [PHP_BIN, WORKER_PATH, qs],
            capture_output=True, text=True, timeout=WORKER_TIMEOUT,
        )
    except Exception:
        return None, _time.time() - t0
    if p.returncode != 0:
        return None, _time.time() - t0
    try:
        data = json.loads(p.stdout.strip())
    except Exception:
        return None, _time.time() - t0
    if not isinstance(data, dict):
        return None, _time.time() - t0
    return data, _time.time() - t0


def handle_parallel(body):
    try:
        req = json.loads(body)
    except Exception:
        return None, "JSON invalido"
    acts = req.get("activitiesList", [])
    if isinstance(acts, str):
        acts = [a.strip() for a in acts.split(",") if a.strip()]
    acts = [a for a in acts if a]
    if not acts:
        return None, "activitiesList vacio"
    base = {k: req.get(k, "") for k in (
        "sociedad", "tarifa", "codigo_destino", "fecha_desde", "fecha_hasta",
        "pasajeros", "edades", "codigo_idioma", "lastId", "bbddreal",
    )}
    results = {}
    # Camino unico para cualquier tarifa: jobs de (acts x 1 fecha).
    fechas = fechas_rango(base.get("fecha_desde", ""), base.get("fecha_hasta", ""))
    act_chunks = [acts[i:i + ACTS_POR_HILO] for i in range(0, len(acts), ACTS_POR_HILO)]
    jobs = [("fecha", di, f, ci, ch) for di, f in enumerate(fechas) for ci, ch in enumerate(act_chunks)]
    maxw = len(jobs) if not MAX_PARALLEL_WORKERS else min(len(jobs), MAX_PARALLEL_WORKERS)
    import time as _time
    import copy as _copy
    t0 = _time.time()
    with ThreadPoolExecutor(max_workers=maxw) as ex:
        futs = {ex.submit(run_worker, base, f, ch): (modo, di, ci) for modo, di, f, ci, ch in jobs}
        for f in as_completed(futs):
            try:
                results[futs[f]] = f.result()
            except Exception:
                results[futs[f]] = (None, 0.0)
    respuesta, rpta_parts, ids_rate = [], [], []
    # regroup por fecha con fusion exacta (ambas tarifas devuelven
    # envelopes por fecha con modalidades; ver below).
    por_uuid = {u: {} for u in acts}
    for (modo, di, ci), (w, _el) in results.items():
        if not w:
            continue
        for pieza in (w.get("resultados") or []):
            u = pieza.get("uuid")
            if u in por_uuid:
                por_uuid[u][di] = pieza
    # Fusion EXACTA replicando datosInternos secuencial por uuid y modalidad:
    #  - por modalidad (orden de 1a aparicion): concat de listas-por-dia
    #    en orden de fecha (solo dias con tipos en esa modalidad)
    #  - Rpta/idsRate: nivel fecha como el secuencial (Rpta = ultima fecha)
    #  - entrada de modalidad: la de la ultima fecha disponible (su politica),
    #    con tipos sustituidos; modalidad sin tipos => se omite;
    #    uuid sin tipos => se descarta.
    for u in acts:  # orden determinista = orden de entrada
        piezas = por_uuid.get(u, {})
        if not piezas:
            continue
        orden_mods, tipos_pm, env_pm = [], {}, {}
        rpta_final, ids_final = "", []
        for di in sorted(piezas.keys()):
            p = piezas[di]
            interna = p.get("interna_datos") if isinstance(p.get("interna_datos"), dict) else {}
            dia_con_tipos = False
            for mod in (interna.get("modalidades") or []):
                if not isinstance(mod, dict):
                    continue
                name = mod.get("modalidad", "")
                try:
                    td = mod["BreakDown"][0]["tipos"] or []
                except Exception:
                    td = []
                if name not in tipos_pm:
                    tipos_pm[name] = []
                    orden_mods.append(name)
                if td:
                    tipos_pm[name].extend(td)
                    dia_con_tipos = True
                env_pm[name] = mod  # ultima fecha disponible gana
            if dia_con_tipos:
                for ir in (p.get("idsRate") or []):
                    ids_final.append(ir)
            r = p.get("Rpta")
            if isinstance(r, str) and r:
                rpta_final = r  # ultima fecha disponible gana
        orden_mods = [n for n in orden_mods if tipos_pm.get(n)]
        if not orden_mods:
            continue
        try:
            base_env = None
            for di in sorted(piezas.keys(), reverse=True):
                ie = piezas[di].get("interna_datos")
                if isinstance(ie, dict) and len(ie) > 0:
                    base_env = ie
                    break
            if base_env is None:
                continue
            env = _copy.deepcopy(base_env)
            nuevas_mods = []
            for name in orden_mods:
                m = _copy.deepcopy(env_pm[name])
                m["BreakDown"][0]["tipos"] = tipos_pm[name]
                nuevas_mods.append(m)
            env["modalidades"] = nuevas_mods
        except Exception:
            continue
        respuesta.append(env)
        if rpta_final:
            rpta_parts.append(rpta_final)
        # un elemento por uuid (lista de rateIds por dia), como el secuencial
        ids_rate.append(ids_final)
    chunk_times = [el for (_w, el) in results.values()]
    wvers = set()
    for (_w, _el) in results.values():
        try:
            if isinstance(_w, dict) and _w.get("worker_v"):
                wvers.add(str(_w.get("worker_v")))
        except Exception:
            pass
    total = _time.time() - t0
    try:
        import datetime as _dt
        import tempfile as _tf
        with open(_tf.gettempdir() + "/dispo_parallel.log", "a") as _lf:
            _lf.write("%s | port=%d | lastId=%s | tarifa=%s | nacts=%d ndias=%d njobs=%d maxw=%d | total=%.1fs maxjob=%.1fs | wver=%s\n" % (
                _dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
                PORT,
                base.get("lastId", ""), base.get("tarifa", ""), len(acts), len(fechas), len(jobs), maxw,
                total, max(chunk_times) if chunk_times else 0.0, ",".join(sorted(wvers)) if wvers else "legacy"))
    except Exception:
        pass
    return {"respuesta": respuesta, "Rpta": "".join(rpta_parts), "idsRate": ids_rate}, None


class Handler(BaseHTTPRequestHandler):
    server_version = "DispoParallel/1.0"

    def _send(self, code, obj):
        body = json.dumps(obj).encode("utf-8")
        self.send_response(code)
        self.send_header("Content-Type", "application/json")
        self.send_header("Content-Length", str(len(body)))
        self.end_headers()
        self.wfile.write(body)

    def do_GET(self):
        if self.path == "/health":
            self._send(200, {"ok": True})
        else:
            self._send(404, {"ok": False})

    def do_POST(self):
        if self.path != "/parallel":
            self._send(404, {"ok": False})
            return
        if self.headers.get("X-Dispo-Token", "") != SERVICE_TOKEN:
            self._send(403, {"ok": False, "error": "token"})
            return
        try:
            length = int(self.headers.get("Content-Length", "0"))
        except ValueError:
            length = 0
        if length <= 0 or length > BODY_MAX_BYTES:
            self._send(400, {"ok": False, "error": "body"})
            return
        out, err = handle_parallel(self.rfile.read(length).decode("utf-8", "replace"))
        if out is None:
            self._send(422, {"ok": False, "error": err or "fail"})
        else:
            self._send(200, out)

    def log_message(self, *args):
        pass


if __name__ == "__main__":
    srv = ThreadingHTTPServer((HOST, PORT), Handler)
    print("dispo_parallel en %s:%d actsxhilo=%d maxw=%s" % (HOST, PORT, ACTS_POR_HILO, MAX_PARALLEL_WORKERS or "all"), flush=True)
    srv.serve_forever()
