#!/usr/bin/env python3
"""
daemon-swl.py — Daemon lite de ejecución de tareas SWL.

Portado del patrón daemon de Multica (cmd/daemon/main.go).
Detecta CLIs disponibles en PATH (claude, codex, gemini) y ejecuta
tareas pendientes desde .planning/task-queue.json mediante polling.

Ciclo de vida de tarea:
  pending → claimed → running → success
                                failed (con reintentos automáticos)

Uso:
  python scripts/daemon-swl.py               # modo daemon (loop infinito)
  python scripts/daemon-swl.py --once        # un ciclo y salir
  python scripts/daemon-swl.py --status      # estado del daemon y cola
  python scripts/daemon-swl.py --stop        # detener daemon existente

Archivos:
  .planning/task-queue.json   — cola de tareas (compartida con task-service.js)
  .planning/daemon.pid        — PID del daemon activo
  .planning/daemon.log        — log de actividad

Variables de entorno:
  SWL_DAEMON_POLL_INTERVAL  — intervalo de polling en segundos (default: 5)
  SWL_DAEMON_MAX_WORKERS    — tareas concurrentes por ciclo (default: 2)
  SWL_DAEMON_CLI_TIMEOUT    — timeout por tarea en segundos (default: 300)
"""

from __future__ import annotations

import json
import logging
import os
import shutil
import signal
import subprocess
import sys
import tempfile
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import Optional

# ---------------------------------------------------------------------------
# Configuración
# ---------------------------------------------------------------------------

POLL_INTERVAL = int(os.environ.get('SWL_DAEMON_POLL_INTERVAL', '5'))
MAX_WORKERS   = int(os.environ.get('SWL_DAEMON_MAX_WORKERS',   '2'))
CLI_TIMEOUT   = int(os.environ.get('SWL_DAEMON_CLI_TIMEOUT',   '300'))

TASK_FILE = '.planning/task-queue.json'
PID_FILE  = '.planning/daemon.pid'
LOG_FILE  = '.planning/daemon.log'

# CLIs soportados en orden de preferencia
SUPPORTED_CLIS = ['claude', 'codex', 'gemini']

# ---------------------------------------------------------------------------
# Logging
# ---------------------------------------------------------------------------

def configurar_logging(cwd: Path) -> None:
    log_path = cwd / LOG_FILE
    log_path.parent.mkdir(parents=True, exist_ok=True)
    logging.basicConfig(
        level=logging.INFO,
        format='%(asctime)s [%(levelname)s] %(message)s',
        handlers=[
            logging.FileHandler(str(log_path), encoding='utf-8'),
            logging.StreamHandler(sys.stdout),
        ],
    )

# ---------------------------------------------------------------------------
# Detección de CLI
# ---------------------------------------------------------------------------

def detectar_cli() -> Optional[str]:
    """
    Detecta el primer CLI soportado disponible en PATH.
    Retorna el nombre del binario (ej: 'claude') o None si ninguno está disponible.
    """
    for cli in SUPPORTED_CLIS:
        ruta = shutil.which(cli)
        if ruta:
            logging.info(f'CLI detectado: {cli} ({ruta})')
            return cli
    return None

# ---------------------------------------------------------------------------
# Cola de tareas (I/O JSON atómico)
# ---------------------------------------------------------------------------

def leer_cola(cwd: Path) -> dict:
    """Lee la cola de tareas. Retorna estructura vacía si no existe."""
    try:
        return json.loads((cwd / TASK_FILE).read_text(encoding='utf-8'))
    except (FileNotFoundError, json.JSONDecodeError):
        return {'version': 1, 'tareas': []}


def escribir_cola_atomica(cwd: Path, cola: dict) -> None:
    """Escritura atómica usando archivo temporal + os.replace()."""
    ruta = cwd / TASK_FILE
    ruta.parent.mkdir(parents=True, exist_ok=True)
    with tempfile.NamedTemporaryFile(
        mode='w',
        encoding='utf-8',
        dir=str(ruta.parent),
        delete=False,
        suffix='.tmp',
    ) as tmp:
        json.dump(cola, tmp, ensure_ascii=False, indent=2)
        tmp_path = tmp.name
    os.replace(tmp_path, str(ruta))


def reclamar_tarea(cwd: Path, worker_id: str) -> Optional[dict]:
    """Reclama la siguiente tarea 'pending' de mayor prioridad (menor número)."""
    cola = leer_cola(cwd)
    pendientes = sorted(
        [t for t in cola['tareas'] if t.get('estado') == 'pending'],
        key=lambda t: (t.get('prioridad', 3), t.get('creadaEn', '')),
    )
    if not pendientes:
        return None

    tarea = pendientes[0]
    tarea['estado']        = 'claimed'
    tarea['workerId']      = worker_id
    tarea['actualizadaEn'] = datetime.now(timezone.utc).isoformat()
    escribir_cola_atomica(cwd, cola)
    return dict(tarea)


def actualizar_estado(cwd: Path, task_id: str, estado: str, **kwargs) -> None:
    """Actualiza el estado de una tarea en la cola."""
    cola = leer_cola(cwd)
    for tarea in cola['tareas']:
        if tarea.get('id') == task_id:
            tarea['estado']        = estado
            tarea['actualizadaEn'] = datetime.now(timezone.utc).isoformat()
            for k, v in kwargs.items():
                tarea[k] = v
            if estado in ('success', 'failed'):
                tarea['completadaEn'] = datetime.now(timezone.utc).isoformat()
            break
    escribir_cola_atomica(cwd, cola)

# ---------------------------------------------------------------------------
# Ejecución de tarea
# ---------------------------------------------------------------------------

def ejecutar_tarea(cwd: Path, tarea: dict, cli: str) -> tuple[bool, str]:
    """
    Ejecuta una tarea usando el CLI disponible.

    Tipos soportados:
      agent  — invoca un agente SWL via CLI con prompt + subagent_type
      bash   — ejecuta un comando shell directamente
      prompt — envía un prompt libre al CLI

    Retorna (exito: bool, output: str).
    """
    tipo    = tarea.get('tipo', 'generic')
    payload = tarea.get('payload', {})

    if tipo == 'agent':
        prompt        = payload.get('prompt', '')
        subagent_type = payload.get('subagent_type', 'general-purpose')
        args = [cli, '--print', f'Actúa como {subagent_type}. {prompt}']

    elif tipo == 'bash':
        comando = payload.get('comando', '').strip()
        if not comando:
            return False, 'Tarea bash sin campo "comando" en payload'
        args = ['bash', '-c', comando]

    elif tipo == 'prompt':
        prompt = payload.get('prompt', '')
        if not prompt:
            return False, 'Tarea prompt sin campo "prompt" en payload'
        args = [cli, '--print', prompt]

    else:
        return False, f'Tipo de tarea no soportado: {tipo}'

    try:
        resultado = subprocess.run(
            args,
            capture_output=True,
            text=True,
            timeout=CLI_TIMEOUT,
            cwd=str(cwd),
        )
        exito  = resultado.returncode == 0
        output = (resultado.stdout or resultado.stderr or '')[:2000]
        return exito, output

    except subprocess.TimeoutExpired:
        return False, f'Timeout después de {CLI_TIMEOUT}s'
    except Exception as exc:
        return False, f'Error de ejecución: {exc}'

# ---------------------------------------------------------------------------
# Ciclo de polling
# ---------------------------------------------------------------------------

def ciclo(cwd: Path, cli: str) -> int:
    """
    Ejecuta un ciclo de polling: reclama y procesa hasta MAX_WORKERS tareas.
    Retorna el número de tareas procesadas en este ciclo.
    """
    worker_id  = f'daemon-{os.getpid()}'
    procesadas = 0

    for _ in range(MAX_WORKERS):
        tarea = reclamar_tarea(cwd, worker_id)
        if tarea is None:
            break

        task_id     = tarea['id']
        descripcion = tarea.get('descripcion', '')
        logging.info(f'Ejecutando [{task_id}]: {descripcion}')

        actualizar_estado(cwd, task_id, 'running')
        exito, output = ejecutar_tarea(cwd, tarea, cli)

        if exito:
            actualizar_estado(cwd, task_id, 'success', output=output)
            logging.info(f'[{task_id}] completada exitosamente')
        else:
            reintentos     = tarea.get('reintentos', 0) + 1
            max_reintentos = tarea.get('maxReintentos', 3)

            if reintentos < max_reintentos:
                actualizar_estado(
                    cwd, task_id, 'pending',
                    workerId=None,
                    reintentos=reintentos,
                    error=f'Reintento {reintentos}/{max_reintentos}: {output}',
                )
                logging.warning(f'[{task_id}] fallida (reintento {reintentos}/{max_reintentos})')
            else:
                actualizar_estado(cwd, task_id, 'failed', error=output, reintentos=reintentos)
                logging.error(f'[{task_id}] fallida definitivamente: {output[:200]}')

        procesadas += 1

    return procesadas

# ---------------------------------------------------------------------------
# Gestión de PID
# ---------------------------------------------------------------------------

def escribir_pid(cwd: Path) -> None:
    pid_path = cwd / PID_FILE
    pid_path.parent.mkdir(parents=True, exist_ok=True)
    pid_path.write_text(str(os.getpid()))


def leer_pid(cwd: Path) -> Optional[int]:
    try:
        return int((cwd / PID_FILE).read_text())
    except (FileNotFoundError, ValueError):
        return None


def limpiar_pid(cwd: Path) -> None:
    (cwd / PID_FILE).unlink(missing_ok=True)

# ---------------------------------------------------------------------------
# Comandos de gestión (--status / --stop)
# ---------------------------------------------------------------------------

def cmd_status(cwd: Path) -> None:
    """Muestra estado del daemon y estadísticas de la cola."""
    pid = leer_pid(cwd)
    print(f'Daemon: {"activo (PID " + str(pid) + ")" if pid else "no activo"}')

    cola   = leer_cola(cwd)
    tareas = cola.get('tareas', [])
    por_estado: dict[str, int] = {}
    for t in tareas:
        est = t.get('estado', 'unknown')
        por_estado[est] = por_estado.get(est, 0) + 1

    print(f'\nCola: {len(tareas)} tarea(s) total')
    for est in ('pending', 'claimed', 'running', 'success', 'failed'):
        cnt = por_estado.get(est, 0)
        if cnt:
            print(f'  {est:<10} {cnt}')


def cmd_stop(cwd: Path) -> None:
    """Detiene el daemon enviando SIGTERM al PID registrado."""
    pid = leer_pid(cwd)
    if not pid:
        print('No hay daemon activo.')
        return
    try:
        os.kill(pid, signal.SIGTERM)
        print(f'SIGTERM enviado al daemon (PID {pid})')
    except ProcessLookupError:
        print(f'PID {pid} ya no existe — limpiando registro')
        limpiar_pid(cwd)

# ---------------------------------------------------------------------------
# Entrypoint
# ---------------------------------------------------------------------------

def main() -> None:
    cwd  = Path.cwd()
    args = sys.argv[1:]

    configurar_logging(cwd)

    if '--status' in args:
        cmd_status(cwd)
        return

    if '--stop' in args:
        cmd_stop(cwd)
        return

    cli = detectar_cli()
    if cli is None:
        logging.error(
            'No se encontró ningún CLI soportado en PATH. '
            'Candidatos: ' + ', '.join(SUPPORTED_CLIS)
        )
        sys.exit(1)

    una_sola_vez = '--once' in args

    if not una_sola_vez:
        # Verificar que no hay otro daemon corriendo
        pid_existente = leer_pid(cwd)
        if pid_existente:
            try:
                os.kill(pid_existente, 0)
                logging.error(
                    f'Ya existe un daemon activo (PID {pid_existente}). '
                    f'Usa --stop para detenerlo.'
                )
                sys.exit(1)
            except ProcessLookupError:
                logging.warning(f'PID obsoleto {pid_existente} — limpiando')
                limpiar_pid(cwd)

        escribir_pid(cwd)
        logging.info(
            f'Daemon SWL iniciado — PID: {os.getpid()}, '
            f'CLI: {cli}, '
            f'Intervalo: {POLL_INTERVAL}s, '
            f'Max workers/ciclo: {MAX_WORKERS}'
        )

        def handle_signal(signum, frame):
            logging.info(f'Señal {signum} recibida — deteniendo daemon')
            limpiar_pid(cwd)
            sys.exit(0)

        signal.signal(signal.SIGTERM, handle_signal)
        signal.signal(signal.SIGINT,  handle_signal)

    try:
        if una_sola_vez:
            procesadas = ciclo(cwd, cli)
            logging.info(f'Ciclo único: {procesadas} tarea(s) procesadas')
        else:
            while True:
                try:
                    procesadas = ciclo(cwd, cli)
                    if procesadas:
                        logging.debug(f'Ciclo completado: {procesadas} tarea(s)')
                except Exception as exc:
                    logging.error(f'Error en ciclo: {exc}')
                time.sleep(POLL_INTERVAL)
    finally:
        if not una_sola_vez:
            limpiar_pid(cwd)


if __name__ == '__main__':
    main()
