#core/CLI.py
import logging

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(name)s: %(message)s"
)

import subprocess
import typer
from core.config import DEV_COMPOSE

logger = logging.getLogger(__name__)

app = typer.Typer()


def docker_commands(*args, mode=None):
    try:
        if mode == "test":
            if not DEV_COMPOSE.exists():
                typer.echo(f"No se encontró {DEV_COMPOSE}")
                raise typer.Exit(code=1)

            subprocess.run(
                ["docker", "compose", "-f", str(DEV_COMPOSE), *args],
                check=True
            )
        else:
            subprocess.run(
                ["docker", "compose", *args],
                check=True
            )
    except FileNotFoundError:
        typer.echo("Docker no está instalado o no está en el PATH")
        raise typer.Exit(code=127)

    except KeyboardInterrupt:
        # Ctrl+C: el usuario interrumpió manualmente → salida limpia
        typer.echo("Interrumpido por el usuario")
        raise typer.Exit(code=130)   # 130 = 128 + SIGINT(2), convención Unix

    except subprocess.CalledProcessError as e:
        code = e.returncode
        if code < 0:
            # El proceso murió por una señal (p. ej. -2 = SIGINT, -9 = SIGKILL)
            code = 128 + abs(code)
        code = min(code, 255)
        typer.echo(f"docker terminó con error (código {code})")
        raise typer.Exit(code=code)


#local docker 
@app.command()
def local_up():
    docker_commands("up", "-d", "--build", mode="test")
    
#produccion docker 
@app.command()
def up():
    docker_commands("up", "-d", "--build")


@app.command()
def down():
    docker_commands("down")

@app.command()
def stop():
    docker_commands("stop")

@app.command()
def start():
    docker_commands("start")

@app.command()
def logs():
    docker_commands("logs", "-f")

if __name__ == "__main__":
    app()

#__________________
#app/main.py
import logging
import asyncio

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(name)s: %(message)s",
)
from app.start import Class_Start
from contextlib import asynccontextmanager
from fastapi import FastAPI

logger = logging.getLogger(__name__)

@asynccontextmanager
async def lifespan(app: FastAPI):
    inicio = Class_Start()
    logger.info("corriendo programa")

    try:
        try:
            await inicio.start_up()
        except asyncio.CancelledError:
            logger.info("start_up cancelado")
            raise
        except Exception:
            logger.exception("error en start_up")
            raise

        yield
    finally:
        await inicio.shut_down()
        

app = FastAPI(
    lifespan=lifespan
)

@app.get("/health")
async def health():
    return {"status": "ok"}


#__________________
#app/start.py
import logging
import asyncio
import os
from core.Sql_Data import DB #instancia
from temp_workflows.create_temporal_client import CLASS_REGISTER #Instancia
from app.socket_run import Class_Socket #clase
from core.redis_ import REDIS_CONNECT #instancia
from temp_workflows.Workflow_ import Class_workflow #clase

class Class_Start:
    def __init__(self):
        self.Class_Register = CLASS_REGISTER
        self.socket_class = Class_Socket()
        self.redis = REDIS_CONNECT
        self.workflow_class = Class_workflow()
        self.event = asyncio.Event()
        self.logger = logging.getLogger(__name__)
        self._supervisor_task: asyncio.Task | None = None
        

    async def start_up(self):
        try:
            await DB.create_tables()
        except Exception:
            self.logger.exception("error creando tablas en SQL")
            raise

        try:
            await self.redis.connect()
        except Exception:
            self.logger.exception("Redis no está disponible")
            raise

        self._supervisor_task = asyncio.create_task(self._run_supervisor())

    
    async def _run_supervisor(self):
        try:
            async with asyncio.TaskGroup() as tg:
                tg.create_task(self.Class_Register.register_client(self.event)) 

                tg.create_task(self.socket_class.socket_server_start(self.redis))

                tg.create_task(self.workflow_class.workflow_disparador(self.redis, self.event))
        
        except* Exception as eg:
            for exc in eg.exceptions:
                self.logger.error("error inesperado", exc_info=(type(exc), exc, exc.__traceback__))
            raise RuntimeError("supervisor falló de forma irrecuperable") from eg.exceptions[0]


    async def shut_down(self):
        if self._supervisor_task is not None:
            self._supervisor_task.cancel()
            await asyncio.gather(self._supervisor_task, return_exceptions=True)
        
        try:
            await self.socket_class.stop()
        except Exception:
            self.logger.exception("error cerrando servidor de sockets")
        
        try:
            await DB.close()
        except Exception:
            self.logger.exception("error cerrando conexión SQL")

        try:
            await self.redis.close_redis()
        except Exception:
            self.logger.exception("error cerrando conexión Redis")



#__________________
#app/socket_run.py
import socket
import asyncio
import ipaddress
from collections import defaultdict
import hmac
import re
from functools import partial
import logging
import redis.exceptions as redis_exceptions
from core.dominio_utils import dominio_seguro
from core.config import(
SERVER_HOST, 
SERVER_PORT, 
SOCKET_TIMEOUT, 
MAX_MESSAGE_SIZE,
DOMINIO_PROHIBIDO,
TIMEOUT_CIERRE_CONEXIONES,
SEMAFORO,
SOCKET_AUTH_TOKEN,
MAX_CONEXIONES_POR_IP
)


class Class_Socket:
    def __init__(self):
        self.semaforo = None
        self.logger = logging.getLogger(__name__)
        self.dominios_prohibidos = {d.lower().strip() for d in DOMINIO_PROHIBIDO}
        self.conexiones_activas: set[asyncio.Task] = set()
        self.server: asyncio.AbstractServer | None = None
        self.conexiones_por_ip: dict[str, int] = defaultdict(int)
        self.TOKEN_LEN = len(SOCKET_AUTH_TOKEN.encode())   # ajustar dinámicamente (fix 2.1)
        self._token_bytes = SOCKET_AUTH_TOKEN.encode()
        

    async def socket_server_start(self, _redis_queue):
        self.semaforo = asyncio.Semaphore(SEMAFORO)

        try:
            self.server = await asyncio.start_server(
                partial(self.on_new_connection, _redis=_redis_queue),
                SERVER_HOST,
                SERVER_PORT
            )

            async with self.server:
                await self.server.serve_forever()
                
        except PermissionError:
            self.logger.warning("Sin permisos para abrir el puerto")
            raise

        except socket.gaierror:
            self.logger.warning("Host inválido")
            raise

        except OSError:
            self.logger.exception("Error del sistema")
            raise

        except asyncio.CancelledError:
            self.logger.info("servidor cancelado: %s/%s", SERVER_HOST, SERVER_PORT)
            raise

        except Exception:
            self.logger.exception("Error inesperado")
            raise


    async def on_new_connection(self, reader, writer, _redis):
        peername = writer.get_extra_info("peername")
        ip_origen = peername[0] if peername else "desconocido"
        
        if self.semaforo.locked():
            writer.close()
            await writer.wait_closed()
            return

        if self.conexiones_por_ip[ip_origen] >= MAX_CONEXIONES_POR_IP:
            self.logger.warning("límite de conexiones alcanzado para %s", ip_origen)
            writer.close()
            await writer.wait_closed()
            return

        self.conexiones_por_ip[ip_origen] += 1
            
        task = asyncio.current_task()

        self.conexiones_activas.add(task)
        
        try:
            async with self.semaforo:
                await self.handle_client(reader, writer, _redis)
            
        finally:
            self.conexiones_activas.discard(task)
            self.conexiones_por_ip[ip_origen] -= 1
            if self.conexiones_por_ip[ip_origen] <= 0:
                del self.conexiones_por_ip[ip_origen]
            

    async def stop(self):
        if self.server is not None:
            self.server.close()
            await self.server.wait_closed()
            self.server = None
            await self.cerrar_conexiones_activas()

    
    async def cerrar_conexiones_activas(self):
        if not self.conexiones_activas:
            return

        self.logger.info("esperando cierre de %d conexiones activas", len(self.conexiones_activas))
        
        pendientes = list(self.conexiones_activas)

        _, aun_pendientes = await asyncio.wait(
            pendientes,
            timeout=TIMEOUT_CIERRE_CONEXIONES
        )

        if aun_pendientes:
            self.logger.warning("cancelando %d conexiones que no cerraron a tiempo", len(aun_pendientes))
            for tarea in aun_pendientes:
                tarea.cancel()

            await asyncio.gather(*aun_pendientes, return_exceptions=True)

        
        
    async def handle_client(self, reader, writer, _redis):
        peername = writer.get_extra_info("peername")

        self.logger.info("informacion del host: %s", peername)

        try:
            token = await asyncio.wait_for(reader.readexactly(self.TOKEN_LEN), timeout=SOCKET_TIMEOUT)
            if not hmac.compare_digest(token, self._token_bytes):
                self.logger.warning("token inválido desde %s", peername)
                await self.enviar_respuesta(writer, "no autorizado")
                return
            
            header = await asyncio.wait_for(
                reader.readexactly(8),
                timeout=SOCKET_TIMEOUT
            )

            try:
                longitud_header = int(header.decode().strip())
            except (ValueError, UnicodeDecodeError):
                self.logger.warning("header inválido, no se pudo interpretar como longitud: %r", header)
                await self.enviar_respuesta(writer, "header invalido")
                return

            if longitud_header <= 0 or longitud_header > MAX_MESSAGE_SIZE:
                self.logger.warning("longitud bytes del mensaje incorrecta: %d, longitud sugerida: %d", longitud_header ,MAX_MESSAGE_SIZE)
                await self.enviar_respuesta(writer, "longitud incorrecta")
                return
            
            data = await asyncio.wait_for(
                reader.readexactly(longitud_header),
                timeout=SOCKET_TIMEOUT
            )
                    
            try:
                dominio_data = data.decode().strip()
            except UnicodeDecodeError:
                self.logger.warning("datos recibidos no son UTF-8 valido")
                await self.enviar_respuesta(writer, "encoding invalido")
                return

            if not self.verificar_dominio(dominio_data):
                self.logger.warning("dominio prohibido rechazado: %s", dominio_data)
                await self.enviar_respuesta(writer, "dominio prohibido")
                return
                    
            try:
                await _redis.get_client().xadd("cola:dominios", {"dominio": dominio_data}, maxlen=100_000, approximate=True)
            except redis_exceptions.ConnectionError as e:
                self.logger.warning("pool de redis agotado o conexion caida: %s", e)
                await self.enviar_respuesta(writer, "servicio no disponible, intenta mas tarde")
                return
                    
            await self.enviar_respuesta(writer, f"dominio recibido correctamente: {dominio_data}")
            
        except ConnectionAbortedError as e:
            self.logger.warning("Conexión abortada: %s", e)
    
        except BrokenPipeError as e:
            self.logger.warning("Broken pipe: %s", e)
    
        except (asyncio.TimeoutError, TimeoutError):
            self.logger.warning("tiempo de espera agotado en reader del host")
            
        except asyncio.IncompleteReadError as e:
            self.logger.warning("conexión cerrada antes de recibir los 8 bytes esperados: %s", e)
             
        except ConnectionResetError as e:
            self.logger.warning("conexión reiniciada por el cliente: %s", e)
            
        except asyncio.CancelledError:
            self.logger.info("conexion cancelada")
            raise
            
        except Exception:
            self.logger.exception("error inesperado leyendo del socket")

        finally:
            try:
                writer.close()
                await writer.wait_closed()
            except Exception:
                self.logger.exception("error al cerrar conexion")

    
    async def enviar_respuesta(self, writer, mensaje: str):
        try:
            writer.write(f"{mensaje}\n".encode())
            await writer.drain()
        except (ConnectionResetError, BrokenPipeError) as e:
            self.logger.warning("no se pudo enviar respuesta al cliente: %s", e)

    
    def verificar_dominio(self, data_dom: str) -> bool:
        dominio = data_dom.lower().strip()

        if self.is_ip(dominio):
            self.logger.warning("se envio ip en lugar de dominio: %s", dominio)
            return False

        if not self.dominio_valido(dominio):
            self.logger.warning("el formato del dominio es invalido: %s", dominio)
            return False

        if dominio in self.dominios_prohibidos:
            return False

        for prohibido in self.dominios_prohibidos:
            if dominio.endswith(f".{prohibido}"):
                return False

        return True

    
    def is_ip(self, dominio: str) -> bool:
        try:
            ipaddress.ip_address(dominio)
            return True
        except ValueError:
            return False

    
    def dominio_valido(self, dominio: str) -> bool:
        if len(dominio) > 253: 
            return False

        try:
            dominio_seguro(dominio)
            return True
        except ValueError:
            return False
            

#__________________
#core/redis_.py
import asyncio
import redis.asyncio as redis
from redis.asyncio import BlockingConnectionPool
import logging
from core.config import REDIS_HOST, REDIS_PORT, REDIS_DB, REDIS_PASSWORD, REDIS_MAX_CONNECTIONS



class Class_Redis:
    def __init__(self):
        self.pool = None
        self.r_connect = None
        self.logger = logging.getLogger(__name__)
        self._reconnect_lock = asyncio.Lock()

    async def connect(self):
        self.pool = BlockingConnectionPool(
            host=REDIS_HOST,
            password=REDIS_PASSWORD,
            port=REDIS_PORT,
            db=REDIS_DB,
            decode_responses=True,
            max_connections=REDIS_MAX_CONNECTIONS,
            timeout=10
        )
        self.r_connect = redis.Redis(connection_pool=self.pool)
        await self.r_connect.ping()   # verifica conexión aquí


    def get_client(self):
        """Siempre devuelve la instancia vigente, incluso tras una reconexión."""
        if self.r_connect is None:
            raise RuntimeError("Redis no está conectado todavía")
        return self.r_connect
        
    async def close_redis(self):
        if self.r_connect is not None:
            try:
                await self.r_connect.aclose()
            except Exception:
                self.logger.exception("error al cerrar conexion del pool redis")

        if self.pool is not None:
            try:
                await self.pool.disconnect()
            except Exception:
                self.logger.exception("error al cerrar el pool de conexiones redis")

    
    async def reconnect(self):
        async with self._reconnect_lock:
            if self.r_connect is not None:
                try:
                    await self.r_connect.ping()
                    return
                except Exception:
                    pass
                    
            self.logger.warning("intentando reconectar a Redis...")
            await self.close_redis()
            await self.connect()
            self.logger.info("reconexión a Redis exitosa")

REDIS_CONNECT = Class_Redis()


#__________________
#temp_workflows/create_temporal_client.py
from temporalio.client import Client
from temporalio.worker import Worker
from temporalio.service import RPCError
from core.config import TASK_QUEUE, TEMPORAL_NAMESPACE
from core.temporal_client_factory import TemporalClientFactory
from temporalio.api.workflowservice.v1 import DescribeTaskQueueRequest
from temp_workflows.ReconWorkflow_ import ReconWorkflow
from temp_workflows.activity import sample_activity
import logging
import asyncio


class Class_Register:
    def __init__(self):
        self.client = None
        self.logger = logging.getLogger(__name__)
        self.factory = TemporalClientFactory()


    async def Init_client(self):
        self.client = await self.factory.connect()

    async def _close_client(self):
        await self.factory.close(self.client)
        self.client = None
        

    async def _try_init_client(self) -> bool:
        try:
            await self.Init_client()
            return True
        except Exception:
            return False
            
            
    async def register_client(self, event):
        worker_task = None
            
        while True:
            if self.client is None:
                if not await self._try_init_client():
                    await asyncio.sleep(5)
                    continue
                
            reconectar = False
                
            try:
                if worker_task is None or worker_task.done():
                    worker_task = asyncio.create_task(self.register_worker())
                
                await asyncio.sleep(5)

                if worker_task.done():
                    exc = worker_task.exception()
                    if exc:
                        self.logger.error("worker falló al iniciar: %s", exc)
                    else:
                        self.logger.warning("worker terminó inesperadamente sin error")
                        reconectar = True
                
                else:
                    resp = await self.client.workflow_service.describe_task_queue(
                        DescribeTaskQueueRequest(namespace=TEMPORAL_NAMESPACE, task_queue={"name": TASK_QUEUE}))
                    
                    if resp.pollers:
                        self.logger.info("worker confirmado activo")
                        event.set()
                        await worker_task  # se queda aquí mientras el worker viva
                    
            except RPCError as e:
                self.logger.error(e)
                reconectar = True

            except ValueError as e:
                self.logger.error("configuración inválida, no se puede continuar: %s", e)
                worker_task.cancel()
                await asyncio.gather(worker_task, return_exceptions=True)
                await self._close_client()
                raise
                    
            except asyncio.CancelledError:
                self.logger.info("cliente temporal cancelado")
                worker_task.cancel()
                await asyncio.gather(worker_task, return_exceptions=True)
                await self._close_client()
                raise

            except Exception:
                self.logger.exception("error inesperado")
                reconectar = True

            if reconectar:
                worker_task.cancel()
                await asyncio.gather(worker_task, return_exceptions=True)
                await self._close_client()
                event.clear()
                await asyncio.sleep(5)
                if not await self._try_init_client():
                    await asyncio.sleep(5)
                    continue


    async def register_worker(self):
        try:
            n_worker = Worker(
                self.client,
                task_queue=TASK_QUEUE,
                workflows=[ReconWorkflow],
                activities=[sample_activity]
            )

            await n_worker.run()
                    

        except ValueError as e:
            self.logger.error("Configuración inválida del worker (falta decorador @activity.defn o @workflow.defn?): %s", e)
            self.logger.info("worker.run dejo de correr")
            raise
            
        except RPCError as e:
            self.logger.error("Worker perdió conexión con el servidor de Temporal: %s", e)
            self.logger.info("worker.run dejo de correr, se reintentara")
            raise

        except Exception as e:
            self.logger.exception("error inesperado: %s", e)
            self.logger.info("worker.run dejo de correr, se reintentara")
            raise

CLASS_REGISTER = Class_Register()


#__________________
#core/temporal_client_factory.py
import logging
from temporalio.client import Client
from temporalio.service import RPCError
from core.config import TEMPORAL_HOST, TEMPORAL_NAMESPACE

class TemporalClientFactory:
    def __init__(self):
        self.logger = logging.getLogger(__name__)

    async def connect(self) -> Client:
        try:
            client = await Client.connect(TEMPORAL_HOST, namespace=TEMPORAL_NAMESPACE)
            self.logger.info("Cliente conectado a Temporal en %s", TEMPORAL_HOST)
            return client
        except RPCError as e:
            self.logger.error("Temporal rechazó la conexión: %s", e)
            raise
        except ValueError as e:
            self.logger.error("TEMPORAL_HOST mal configurado: %s", e)
            raise
        except (ConnectionRefusedError, OSError) as e:
            self.logger.error("No se pudo alcanzar el host de Temporal: %s", e)
            raise
        except TimeoutError as e:
            self.logger.error("Timeout al conectar con Temporal: %s", e)
            raise
        except Exception:
            self.logger.exception("error inesperado al conectar con el cliente")
            raise

    async def close(self, client: Client | None):
        if client is None:
            return
        try:
            if hasattr(client, "close"):
                await client.close()
        except Exception as e:
            self.logger.warning("error al cerrar cliente temporal (ignorado): %s", e)
            
#__________________
#temp_workflows/Workflow_.py
import logging
import asyncio
import uuid
from core.redis_ import REDIS_CONNECT
from datetime import timedelta
from temp_workflows.ReconWorkflow_ import ReconWorkflow
from core.Sql_Data import DB
from temporalio.client import Client
from temporalio.service import RPCError
from temporalio.exceptions import WorkflowFailureError, ActivityError
from core.config import TEMPORAL_HOST, TASK_QUEUE, TEMPORAL_NAMESPACE, WORKFLOW_CONCURRENCY, LIMITE_GLOBAL_REINTENTOS
from core.temporal_client_factory import TemporalClientFactory
from redis.exceptions import ResponseError
import redis.exceptions as redis_exceptions


class Class_workflow:
    def __init__(self):
        self.client = None
        self.logger = logging.getLogger(__name__)
        self.factory = TemporalClientFactory()
        self._client_lock = asyncio.Lock()

    
    async def asegurar_grupo(self, redis, stream="cola:dominios", grupo="workers"):
        try:
            await redis.xgroup_create(name=stream, groupname=grupo, id="0", mkstream=True)
            self.logger.info("grupo de consumidores '%s' creado", grupo)
        except ResponseError as e:
            if "BUSYGROUP" in str(e):
                self.logger.info("grupo de consumidores '%s' ya existe", grupo)
            else:
                self.logger.error("error creando grupo de consumidores: %s", e)
                raise

    async def Init_client(self):
        async with self._client_lock:
            if self.client is not None:
                return   # otro worker ya reconectó mientras esperábamos el lock
            self.client = await self.factory.connect()

    async def _close_client(self):
        async with self._client_lock:
            await self.factory.close(self.client)
            self.client = None
    
    async def _reclaim_loop(self, redis_wrapper, stream="cola:dominios", grupo="workers",consumer="worker-1", min_idle_ms=60_000, intervalo=30):
        cursor = "0-0"

        while True:
            try:
                await asyncio.sleep(intervalo)

                cursor, mensajes_reclamados, _ = await redis_wrapper.get_client().xautoclaim(
                    name=stream,
                    groupname=grupo,
                    consumername=consumer,
                    min_idle_time=min_idle_ms,
                    start_id=cursor,
                    count=50
                )

                if mensajes_reclamados:
                    self.logger.warning("reclamados %d mensajes huérfanos del stream %s",len(mensajes_reclamados), stream)

                    for message_id, data in mensajes_reclamados:
                        dominio_data = data.get("dominio")
                        if not dominio_data:
                            await redis_wrapper.get_client().xack(stream, grupo, message_id)
                            continue

                        intentos_previos = int(data.get("intentos", 0))
                        nuevos_intentos = intentos_previos + 1
                        
                        if nuevos_intentos >= LIMITE_GLOBAL_REINTENTOS:
                            self.logger.warning("dominio %s agotó reintentos tras quedar huérfano", dominio_data)
                            await DB.insert_dead_letter(
                                f"reintentos agotados (huérfano) para {dominio_data}", nuevos_intentos, dominio_data
                            )
                        
                        else:
                            await redis_wrapper.get_client().xadd(
                                stream, {"dominio": dominio_data, "intentos": nuevos_intentos},
                                maxlen=100_000, approximate=True
                            )
                        
                        await redis_wrapper.get_client().xack(stream, grupo, message_id)
                    
            
            except asyncio.CancelledError:
                self.logger.info("reclaim_loop cancelado")
                raise

            except redis_exceptions.ConnectionError:
                self.logger.warning("redis caído en reclaim_loop, reconectando")
                try:
                    await redis_wrapper.reconnect()
                except Exception:
                    self.logger.exception("no se pudo reconectar a redis")
                
            except Exception:
                self.logger.exception("error en reclaim_loop, reintentando en %ds", intervalo)
                

    
    async def workflow_disparador(self, redis_wrapper, event):
        await event.wait()
        await self.asegurar_grupo(redis_wrapper.get_client())

        reclaim_task = asyncio.create_task(self._reclaim_loop(redis_wrapper))

        try:
            async with asyncio.TaskGroup() as tg:
                for i in range(WORKFLOW_CONCURRENCY):
                    tg.create_task(self._worker_loop(redis_wrapper, consumer=f"worker-{i}"))

        finally:
            reclaim_task.cancel()
            await asyncio.gather(reclaim_task, return_exceptions=True)
            

    async def _worker_loop(self, redis_wrapper, consumer):
        while True:
            if self.client is None:
                try:
                    await self.Init_client()
                except Exception:
                    self.logger.warning("no se pudo conectar, reintentando en 5s")
                    await asyncio.sleep(5)
                    continue

            try:
                messages = await redis_wrapper.get_client().xreadgroup(
                    groupname="workers",
                    consumername=consumer,
                    streams={"cola:dominios": ">"},
                    count=10,
                    block=2000
                )
            
            except redis_exceptions.ConnectionError as e:
                self.logger.error("conexión con redis perdida, reconectando: %s", e)
                try:
                    await redis_wrapper.reconnect()
                except Exception:
                    self.logger.exception("no se pudo reconectar a redis")
                await asyncio.sleep(5)
                continue
            
            except Exception as e:
                self.logger.error("error al leer de redis: %s", e)
                await asyncio.sleep(5)
                continue

            if not messages:
                continue
                
            stream_name, entries = messages[0]
            
            if not entries:
                continue

            tasks = [self._procesar_dominio(redis_wrapper, message_id, data["dominio"], int(data.get("intentos", 0))) for message_id, data in entries]

            result = await asyncio.gather(*tasks, return_exceptions=True)

            for e in result:
                if isinstance(e, Exception):
                    self.logger.error("error inesperado procesando un dominio del lote: %s", e)
                         
    
    async def _procesar_dominio(self, redis_wrapper, message_id, dominio_data, intentos_previos):
        if self.client is None:
            try:
                await self.Init_client()
            except Exception:
                self.logger.warning("no se pudo conectar a Temporal, se reintentará vía reclaim_loop")
                await asyncio.sleep(5)
                return

        workflow_id = f"recon:{dominio_data}"
        
        try:
            handle = await self.client.start_workflow(
                ReconWorkflow.run,
                dominio_data,
                id=workflow_id,
                #id=f"mainworkflow-{uuid.uuid4()}",
                task_queue=TASK_QUEUE,
                execution_timeout=timedelta(minutes=30),
            )

            await handle.result()
            await redis_wrapper.get_client().xack("cola:dominios", "workers", message_id)
            self.logger.info("workflow completado para %s", dominio_data)

            
        except RPCError as e:
            self.logger.error("no se pudo comunicar con el servidor temporal: %s", e)
            await self._close_client()

            nuevos_intentos = intentos_previos + 1
            if nuevos_intentos >= LIMITE_GLOBAL_REINTENTOS:
                await DB.insert_dead_letter(str(e), nuevos_intentos, dominio_data)
            else:
                await redis_wrapper.get_client().xadd(
                    "cola:dominios", {"dominio": dominio_data, "intentos": nuevos_intentos},maxlen=100_000, approximate=True)
            
            await redis_wrapper.get_client().xack("cola:dominios", "workers", message_id)
            await asyncio.sleep(5)


        except WorkflowFailureError as e:
            self.logger.error("el workflow para %s falló definitivamente: %s", dominio_data, e.cause)
            try:
                await DB.insert_dead_letter(str(e.cause), 1, dominio_data)
            except Exception:
                self.logger.warning("no se pudo guardar dead letter")
            finally:
                await redis_wrapper.get_client().xack("cola:dominios", "workers", message_id)

        except asyncio.CancelledError as e:
            self.logger.info("procesamiento de dominio cancelado")
            try:
                await asyncio.shield(DB.insert_dead_letter(str(e), 0, "cancelled"))
            except Exception:
                self.logger.exception("error al guardar dato en SQL")
            raise
                
        except Exception as e:
            self.logger.exception("error inesperado ejecutando workflow para %s: %s", dominio_data, e)
            try:
                await DB.insert_dead_letter(str(e), 1, dominio_data)
            except Exception:
                self.logger.warning("no se pudo guardar dead letter")
            finally:
                await redis_wrapper.get_client().xack("cola:dominios", "workers", message_id)                       
        
#__________________
#temp_workflows/ReconWorkflow_.py
import uuid
from temp_workflows.activity import sample_activity
from temporalio import workflow
from datetime import timedelta
from temporalio.common import RetryPolicy
from temporalio import activity
from core.config import TASK_QUEUE


@workflow.defn
class ReconWorkflow:
    @workflow.run
    async def run(self, dominio: str) -> dict:
        return await workflow.execute_activity(
            sample_activity,
            dominio,
            start_to_close_timeout=timedelta(minutes=45),
            heartbeat_timeout=timedelta(seconds=60),
            retry_policy=RetryPolicy(
                maximum_attempts=3,
                non_retryable_error_types=["ValueError"]),
        )

#__________________
#temp_workflows/activity.py
from temporalio import activity
from graph.graph_start import graph_ainvoke
import logging

logger = logging.getLogger(__name__)

@activity.defn
async def sample_activity(dominio: str) -> dict:
    async def _hb():
        while True:
            activity.heartbeat()
            await asyncio.sleep(10)
    hb = asyncio.create_task(_hb())
    try:
        return await graph_ainvoke(dominio)

    except Exception:
        logger.exception("error con graph_start")
        raise
        
    finally:
        hb.cancel()
        await asyncio.gather(hb, return_exceptions=True)

#__________________
#graph/graph_start.py
from graph.graph_init_compile import grafo_compilado
import logging
import asyncio


logger = logging.getLogger(__name__)

async def graph_ainvoke(dominio_data):

    platilla = {
        "dominio": dominio_data,
        "subdominios": {},
        "dominio_ips": {},
        "resultado_nmap": {},
        "resultado_nuclei": {},
        "resultado_filtro": {}
    }

    
    try:
        resultado = await grafo_compilado.ainvoke(platilla)

        return resultado
        
    except KeyError as e:
        logger.error("campo faltante en el estado, revisa el orden de los edges: %s", e)
        raise

    except asyncio.TimeoutError:
        logger.error("el pipeline para %s excedió el tiempo máximo", dominio_data)
        raise

    except asyncio.CancelledError:
        logger.info("ejecución del grafo cancelada")
        raise

    except Exception as e:
        logger.exception("error inesperado ejecutando el pipeline para %s: %s", dominio_data, e)
        raise
        

#__________________
#graph/graph_init_compile.py
from graph.graph_config import config_grafo

grafo_compilado = config_grafo()

#__________________
#graph/graph_config.py
from langgraph.graph import StateGraph, END
from typing import TypedDict
import logging
from tools._amass import analizer_amass
from tools._resolver_dns import DNS_resolve
from tools._nmap import analizer_nmap
from tools._nuclei import analizer_nuclei
from tools._filtro import analizer_filtro

resolver = DNS_resolve()
logger = logging.getLogger(__name__)


class State(TypedDict):
    dominio: str
    subdominios: dict[str, dict[str, list[str]]]
    dominio_ips: dict[str, list[str]]
    resultado_nmap: dict[str, list[str]]
    resultado_nuclei: dict[str, dict]
    resultado_filtro: dict[str, list[str]]


class PartialState(TypedDict, total=False):
    subdominios: dict[str, list[str]]
    dominio_ips: dict[str, list[str]]
    resultado_nmap: dict[str, list[str]]
    resultado_nuclei: dict[str, dict]
    resultado_filtro: dict[str, list[str]]


async def amass(state: State) -> PartialState:
    dominio = state.get("dominio", "desconocido")
    try:
        resultado = await analizer_amass(dominio)
        return {"subdominios": resultado}
    except Exception as e:
        logger.exception("fallo crítico en amass para %s: %s", dominio, e)
        raise

async def resolver_dns(state: State) -> PartialState:
    dominio = state.get("dominio", "desconocido")
    try:
        resultado = await resolver.analizer_resolver_dns(state["subdominios"])
        return {"dominio_ips": resultado}
    except Exception as e:
        logger.exception("error critico en resolver DNS para %s: %s", dominio, e)
        raise

async def nmap(state: State) -> PartialState:
    dominio = state.get("dominio", "desconocido")
    try:
        resultado = await analizer_nmap(state["dominio_ips"])
        return {"resultado_nmap": resultado}
    except Exception as e:
        logger.exception("error en analisis de nmap para %s: %s", dominio, e)
        raise
        
async def nuclei(state: State) -> PartialState:
    dominio = state.get("dominio", "desconocido")
    try:
        resultado = await analizer_nuclei(state["dominio_ips"])
        return {"resultado_nuclei": resultado}
    except Exception as e:
        logger.exception("error en analisis de nuclei para %s: %s", dominio, e)
        raise

async def filtro(state: State) -> PartialState:
    dominio = state.get("dominio", "desconocido")
    try:
        resultados_nmap = state.get("resultado_nmap")
        resultados_nuclei = state.get("resultado_nuclei")

        dominios = state.get("dominio_ips")

        if resultados_nmap is None or resultados_nuclei is None or dominios is None:
            logger.warning("filtro llamado antes de tener los 2 resultados para: %s", dominio)
            raise ValueError(f"faltan resultados de nmap/nuclei para completar filtro en {dominio}")
            
        resultado = await analizer_filtro(resultados_nmap, resultados_nuclei, dominios, dominio)
        return {"resultado_filtro": resultado}
        
    except Exception as e:
        logger.exception("error en filtro para %s: %s", dominio, e)
        raise


def config_grafo():
    grafo = StateGraph(State)

    grafo.add_node("amass", amass)
    grafo.add_node("resolver_dns", resolver_dns)
    grafo.add_node("nmap", nmap)
    grafo.add_node("nuclei", nuclei)
    grafo.add_node("filtro", filtro)

    grafo.set_entry_point("amass")

    grafo.add_edge("amass", "resolver_dns")
    grafo.add_edge("resolver_dns", "nmap")
    grafo.add_edge("resolver_dns", "nuclei")
    grafo.add_edge("nmap", "filtro")
    grafo.add_edge("nuclei", "filtro")
    grafo.add_edge("filtro", END)

    try:
        return grafo.compile()
    except Exception:
        logger.critical("no se pudo compilar el grafo")
        raise

#__________________
#tools/_amass.py
import asyncio
import json
import uuid
import logging
from pathlib import Path
from core.config import PATH_AMASS
from core.dominio_utils import dominio_seguro

logger = logging.getLogger(__name__)

async def analizer_amass(dominio_url: str) -> dict:
    dominio_url = dominio_seguro(dominio_url)
    tmp = Path("/tmp")
    tmp.mkdir(exist_ok=True)
    
    salida_json = tmp / f"amass_{dominio_url}_{uuid.uuid4().hex[:8]}.json"
    
    try:
        resultado = await asyncio.create_subprocess_exec(
            PATH_AMASS, "enum",
            "-d", dominio_url,
            "-timeout", "20",
            "-json", salida_json,
            stdout=asyncio.subprocess.PIPE,
            stderr=asyncio.subprocess.PIPE
        )

        try:
            stdout, stderr = await asyncio.wait_for(
                resultado.communicate(),
                timeout=1300
            )
        except asyncio.TimeoutError:
            logger.error("amass excedio el tiempo limite de python")
            resultado.kill()
            await resultado.wait()
            raise
        
        if resultado.returncode != 0:
            raise RuntimeError(f"error en amass {dominio_url}: {stderr.decode()}")

        subdominios_con_ip = {}
        subdominios_sin_ip = {}
        try:
            with open(salida_json) as f:
                for linea in f:
                    if linea.strip():
                        dato = json.loads(linea)
                        nombre = dato["name"]
                        direccion = dato.get("addresses", [])

                        if direccion:
                            ip = [addr["ip"] for addr in direccion]
                            subdominios_con_ip[nombre] = ip

                        else:
                            subdominios_sin_ip[nombre] = []
                        
        except FileNotFoundError:
            logger.error("amass no encontro nada o fallo al escribir el archivo")
            raise
        
        finally:
            try:
                salida_json.unlink()
            except FileNotFoundError:
                pass
            except OSError as e:
                logger.warning("no se pudo borrar el archivo temporal %s: %s", salida_json, e)
            
    except FileNotFoundError:
        logger.error("amass no está instalado o no se encuentra en el PATH")
        raise
    except PermissionError:
        logger.error("sin permisos para ejecutar amass")
        raise
    except OSError as e:
        logger.error("error del sistema al lanzar amass: %s", e)
        raise
    except ValueError:
        logger.error("parametros del comando incorrecta")
        raise

    return {
        "con_ip": subdominios_con_ip,
        "sin_ip": subdominios_sin_ip
    }
        

#__________________
#tools/_resolver_dns.py
import logging
import dns.asyncresolver
import dns.resolver
import dns.exception
import asyncio

class DNS_resolve:
    def __init__(self):
        self.logger = logging.getLogger(__name__)
        self.semaforo = asyncio.Semaphore(20)
    
    
    async def analizer_resolver_dns(self, datos_DNS: dict) -> dict:
        con_ip = datos_DNS.get("con_ip", {})
        sin_ip = datos_DNS.get("sin_ip", {})
        dict_final = {}
        lista_dominios_sin_ip = []

        if not con_ip and not sin_ip:
            self.logger.error("no se encontraron datos en los resultados de amass")
            raise ValueError("resultados de amass sin datos")

        for tipo, dict_dominios in datos_DNS.items():
            
            for dominio, ips in dict_dominios.items():
                if ips:
                    dict_final[dominio] = ips
                else:
                    lista_dominios_sin_ip.append(dominio)
        
        tarea = [self.resolver_dominio(dominio) for dominio in lista_dominios_sin_ip]

        resultado = await asyncio.gather(*tarea)

        for tupla in resultado:
            if not tupla:
                continue
            
            sub, ips = tupla
            dict_final[sub] = ips

        return dict_final


    async def resolver_dominio(self, dominio: str) -> tuple:
        resolver = dns.asyncresolver.Resolver()
        resolver.timeout = 5
        resolver.lifetime = 10

        async with self.semaforo:
            try:
                respuesta = await resolver.resolve(dominio, "A")

                lista = [str(ip) for ip in respuesta]

                return dominio, lista
            
            except dns.resolver.NXDOMAIN:
                self.logger.info("%s no existe (NXDOMAIN)", dominio)
                return ()
        
            except dns.resolver.NoAnswer:
                self.logger.info("%s no tiene registro A", dominio)
                return ()
        
            except dns.resolver.LifetimeTimeout:
                self.logger.warning("timeout resolviendo %s", dominio)
                return ()
        
            except dns.exception.DNSException as e:
                self.logger.warning("error DNS para %s: %s", dominio, e)
                return ()


#__________________
#tools/_nmap.py
import logging
import asyncio
from pathlib import Path
import tempfile
from collections import defaultdict
import json
import uuid
import xmltodict
from core.Sql_Data import DB
from core.dominio_utils import dominio_seguro
from core.config import PATH_NMAP, NMAP_TIMEOUT

logger = logging.getLogger(__name__)
folder_temp = Path(tempfile.gettempdir())
intentos = 2

async def analizer_nmap(dominios_y_ip: dict) -> dict:
    if not dominios_y_ip:
        raise ValueError("datos de entrada para nmap vacíos")
        
    sem = asyncio.Semaphore(10)

    dominios = [d for d, ips in dominios_y_ip.items() if ips]

    #dominios = [dominio_seguro(d) for d, ips in dominios_y_ip.items() if ips]
    
    dominios_validos = []

    for valido in dominios:
        try:
            dominios_validos.append(dominio_seguro(valido))
            
        except ValueError:
            logger.warning("dominio invalido descartado: %s",valido)
            
    if not dominios_validos:
        return {}
            

    tasks = [scanner_ip(d, dominios_y_ip[d], sem) for d in dominios_validos]
    
    resultados = await asyncio.gather(*tasks, return_exceptions=True)

    dict_resultados = {}
    for dominio, resultado in zip(dominios_validos, resultados):
        if isinstance(resultado, Exception):
            logger.error("nmap falló para %s: %s", dominio, resultado)
            await DB.insert_dead_letter(f"nmap falló para {dominio}", 1, dominio)
            continue
        dict_resultados[dominio] = resultado

    if not dict_resultados:
        raise RuntimeError("nmap falló para todos los dominios")

    return dict_resultados

    
async def scanner_ip(dominio: str, ips: list[str], sem: asyncio.Semaphore) -> dict:
    if not dominio or not ips: 
        raise ValueError("datos de entrada invalidos: dominio=%s y ip=%s para nmap vacíos", dominio, ips)

    ultimo_error = None

    async with sem:
        for intento in range(1, intentos + 1):
            file_nmap = folder_temp / f"nmap_{dominio}_{uuid.uuid4().hex[:8]}.xml"
            proc = None
    
            try:
                proc = await asyncio.create_subprocess_exec(
                    PATH_NMAP, *ips, "-oX", str(file_nmap),
                    stdout=asyncio.subprocess.PIPE,
                    stderr=asyncio.subprocess.PIPE
                )
                try:
                    stdout, stderr = await asyncio.wait_for(
                        proc.communicate(),
                        timeout=NMAP_TIMEOUT
                    )
                except (asyncio.TimeoutError, TimeoutError):
                    proc.kill()
                    await proc.wait()
                    raise

                if proc.returncode != 0:
                    error_msg = stderr.decode() if stderr else "Código de retorno no cero"
                    raise RuntimeError(f"nmap falló: {error_msg}")

                if stderr:
                    logger.warning("nmap emitió advertencias: %s", stderr.decode())
            

                data = parsear_salida(file_nmap)   # ya existente, borra el archivo al terminar/
                hosts = (data.get("nmap") or {}).get("host") or []
                
                if isinstance(hosts, dict):
                    hosts = [hosts]

                resultado_por_ip = {}

                for host in hosts:
                    address = host.get("address") or {}
                    if isinstance(address, dict):
                        address = [address]
                    ipv4 = next((a for a in address if a.get("@addrtype") == "ipv4"), None)
                    ip = (ipv4 or {}).get("@addr")
                    if ip:
                        resultado_por_ip[ip] = host
                
                return resultado_por_ip  # éxito: sale de la función, no reintenta más
                
                
            except (asyncio.TimeoutError, TimeoutError) as e:
                ultimo_error = e
                logger.warning("timeout escaneando %s (intento %d/%d)", dominio, intento, intentos)
                if proc is not None:
                    proc.kill()
                    await proc.wait()

            except (RuntimeError, FileNotFoundError, ValueError) as e:
                ultimo_error = e
                logger.warning("fallo escaneando %s (intento %d/%d): %s",dominio, intento, intentos, e)

            if intento < intentos:
                await asyncio.sleep(2 * intento)  # backoff simple: 2s, 4s...

    # se agotaron los intentos
    raise RuntimeError(f"nmap falló tras {intentos} intentos para {dominio}: {ultimo_error}")
                

def parsear_salida(ruta):
    if not ruta.exists():
        logger.error("archivo no existe: %s", ruta)
        raise FileNotFoundError(f"Archivo no encontrado: {ruta}")

    try:
        with open(ruta, "r") as f:
            data = xmltodict.parse(f.read())

        if not data or data == {}:
            raise ValueError(f"archivo xml vacio: {ruta}")

        if not data or 'nmap' not in data:
            raise ValueError(f"XML sin estructura válida: {ruta}")
        
        return data
        
    finally:
        try:
            if ruta.exists():
                ruta.unlink()
                logger.info("Archivo temporal eliminado: %s", ruta)

            else:
                logger.info("el archivo que intenta eliminar no existe: %s", ruta)
        
        except Exception as e:
            logger.error("no se pudo eliminar el archivo %s: %s", ruta, e)       
    

#__________________
#tools/_nuclei.py
import logging
import asyncio
import json
import uuid
import tempfile
from pathlib import Path
from collections import defaultdict
from core.Sql_Data import DB, Nuclei
from core.dominio_utils import dominio_seguro
from core.config import PATH_NUCLEI, NUCLEI_TIMEOUT

logger = logging.getLogger(__name__)

def _dividir_en_lotes(lista, tamano):
    for i in range(0, len(lista), tamano):
        yield lista[i:i + tamano]
        

async def analizer_nuclei(dominio_ips: dict) -> dict:
    if not dominio_ips:
        logger.warning("entrada con datos vacios")
        raise ValueError("datos recibidos por resolver_dns vacios")

    dominios = list(dominio_ips.keys())
    sem = asyncio.Semaphore(3)                 # nº de procesos nuclei simultáneos
    lotes = list(_dividir_en_lotes(dominios, 25))  # 25 dominios por proceso

    resultados_lotes = await asyncio.gather(
        *(_escanear_lote(lote, sem) for lote in lotes),
        return_exceptions=True
    )

    dict_final = {}
    for lote, resultado in zip(lotes, resultados_lotes):
        fallidos = 0
        if isinstance(resultado, Exception):
            fallidos += 1
            if fallidos == len(lotes):
                raise RuntimeError("nuclei falló para todos los lotes")
                
            logger.error("lote: %s de nuclei falló: %s",lote, resultado)
            for dominio in lote:
                await DB.insert_dead_letter(f"nuclei falló para {dominio} (lote)", 1, dominio)
            continue
            
        for dominio, hallazgos in resultado.items():
            dict_final[dominio] = hallazgos

    for dominio, hallazgos in dict_final.items():
        try:
            await DB.Upsert(Nuclei, hallazgos, dominio)
        except Exception as e:
            logger.error("error al guardar datos de %s: %s", dominio, e)

    return dict_final


async def _escanear_lote(dominios_lote: list[str], sem: asyncio.Semaphore) -> dict:
    """Un solo proceso nuclei para varios dominios a la vez, usando -l."""
    folder_file = tempfile.gettempdir()
    lote_id = uuid.uuid4().hex[:8]

    dominios_validos = []

    for d in dominios_lote:
        try:
            dominios_validos.append(dominio_seguro(d))
        except ValueError:
            logger.warning("dominio invalido descartado: %s", d)

    if not dominios_validos:
        return {}   # nada válido que escanear en este lote
            
    urls = [
        d if d.startswith(("http://", "https://")) else f"https://{d}"
        for d in dominios_validos
    ]
    url_a_dominio = dict(zip(urls, dominios_validos))

    lista_urls_path = Path(f"{folder_file}/nuclei_urls_{lote_id}.txt")
    salida_path = Path(f"{folder_file}/nuclei_out_{lote_id}.json")

    lista_urls_path.write_text("\n".join(urls))

    async with sem:
        try:
            proc = await asyncio.create_subprocess_exec(
                PATH_NUCLEI, "-l", str(lista_urls_path), "-json-export", str(salida_path),
                stdout=asyncio.subprocess.PIPE,
                stderr=asyncio.subprocess.PIPE
            )
            try:
                stdout, stderr = await asyncio.wait_for(proc.communicate(), timeout=NUCLEI_TIMEOUT)
            except (asyncio.TimeoutError, TimeoutError):
                proc.kill()
                await proc.wait()
                raise

            if proc.returncode != 0:
                raise RuntimeError(f"nuclei falló en lote: {stderr.decode()}")

            if not salida_path.exists():
                return {}  # lote sin hallazgos es válido, no es error

            contenido = salida_path.read_text().strip()
            if not contenido:
                return {}

            items = json.loads(contenido)
            if isinstance(items, dict):
                items = [items]

            resultado_por_dominio = defaultdict(list)
            for item in items:
                url_encontrada = item.get("matched-at") or item.get("host") or ""
                dominio = next(
                    (d for u, d in url_a_dominio.items() if url_encontrada.startswith(u)),
                    None
                )
                if dominio:
                    resultado_por_dominio[dominio].append(item)

            return dict(resultado_por_dominio)


        finally:
            for p in (lista_urls_path, salida_path):
                try:
                    p.unlink()
                except FileNotFoundError:
                    pass
                except OSError as e:
                    logger.warning("no se pudo borrar %s: %s", p, e)
                    
#__________________
#tools/_filtro.py
from urllib.parse import urlparse
import logging
from core.Sql_Data import Salida_final, DB
from core.config import DOWNLOADS_DIR
import json
import aiofiles

logger = logging.getLogger(__name__)

TAGS_LOGIN = {"panel", "login", "admin", "cms"}
EXTENSIONES_LOGIN = (".php", ".asp", ".aspx", ".jsp", ".cgi")
SEGMENTOS_LOGIN = ["login", "admin", "xmlrpc", "signin", "wp-admin", "cpanel", "portal", "manage"]

async def analizer_filtro(resultado_nmap: dict, resultado_nuclei: dict, dominios: dict, dominio_principal: str) -> dict:
    formato_dict = {}
    
    for d in dominios.keys():

        formato_dict[d] = {
            "hosts": [],
            "url_login": [],
            "descriptions": [],
        }

        
        puertos = []      

        try:
            nmap_por_ip = resultado_nmap.get(d) or {}

            if not nmap_por_ip:
                logger.info("datos nmap de %s vacios", d)

            for ip, datos_ip in nmap_por_ip.items():          # <-- NUEVO: recorre TODAS las IPs
                
                nmap_data = ((datos_ip.get("nmap") or {}).get("host") or {})
                if not nmap_data:
                    continue

                puertos_ip = (nmap_data.get("ports") or {}).get("port") or []
                if isinstance(puertos_ip, dict):
                    puertos_ip = [puertos_ip]
                elif not isinstance(puertos_ip, list):
                    logger.warning("formato de puertos inesperado en %s (ip %s): %r", d, ip, puertos_ip)
                    puertos_ip = []
                    
                data_puertos = []
                for port in puertos_ip:
                    if not port:
                        continue
                        
                    estado = (port.get("state") or {}).get("@state")
                    if estado != "open":
                        continue
                    service = port.get("service") or {}
                    data_puertos.append({
                        "puerto": port.get("@portid"),
                        "protocolo": port.get("@protocol"),
                        "servicio": service.get("@name"),
                        "tecnologia": service.get("@product"),
                    })
                
                formato_dict[d]["hosts"].append({"ip": ip, "data": data_puertos})
        
        except (AttributeError, TypeError, KeyError, IndexError) as e:
            logger.warning("error procesando IP de %s: %s", d, e)
            
        except Exception as e:
            logger.exception("error inesperado procesando nmap de %s: %s", d, e)

        
        try: 
            nuclei_data = resultado_nuclei.get(d) or []
            if isinstance(nuclei_data, dict):
                nuclei_data = [nuclei_data]
            elif not isinstance(nuclei_data, list):
                logger.warning("formato de nuclei inesperado en %s: %r", d, nuclei_data)
                nuclei_data = []

            descriptions = []
            for item in nuclei_data:
                try:
                    desc = (item.get("info") or {}).get("description")
                except AttributeError:
                    logger.warning("item de nuclei inválido en %s", d)
                    continue
                
                if desc:
                    descriptions.append(desc)
                    formato_dict[d]["descriptions"] = list(dict.fromkeys(descriptions))

            
            for item in nuclei_data:
                try:
                    url = item.get("matched-at")
                    tags_raw = (item.get("info") or {}).get("tags", "")
                    tags = set(t.strip() for t in tags_raw.split(",")) if tags_raw else set()
                except AttributeError:
                    logger.warning("item de nuclei inválido en %s", d)
                    continue
                    
                if not url:
                    continue
                    
                path = urlparse(url).path.lower()
                
                coincide_por_tag = bool(tags & TAGS_LOGIN)
                
                coincide_por_patron = (
                    any(ext in path for ext in EXTENSIONES_LOGIN) or
                    any(seg in path for seg in SEGMENTOS_LOGIN)
                )
                
                if coincide_por_tag or coincide_por_patron:
                    formato_dict[d]["url_login"].append(url)

            formato_dict[d]["url_login"] = list(dict.fromkeys(formato_dict[d]["url_login"]))
        
            try:
                await DB.Upsert(Salida_final, formato_dict[d], d)
                logger.info("datos de %s guardado correctamente", d)
                    
            except Exception as e:
                logger.error("error al guardar datos de %s: %s", d, e)
                
            
        except (AttributeError, TypeError, KeyError, IndexError) as e:
            logger.warning("error procesando %s: %s", d, e)

        except Exception as e:
            logger.exception("error inesperado, %s, %s", d, e)
        

    nombre_archivo = f"{dominio_principal}.json"
    ruta_salida = DOWNLOADS_DIR / nombre_archivo

    async with aiofiles.open(ruta_salida, "w", encoding="utf-8") as f:
        await f.write(json.dumps(formato_dict, ensure_ascii=False, indent=4))
        
    
    return formato_dict
            
#__________________
#core/dominio_utils.py
import re

PATRON_DOMINIO = re.compile(
    r'^(?!-)[A-Za-z0-9-]{1,63}(?<!-)(\.[A-Za-z0-9-]{1,63}(?<!-))*\.[A-Za-z]{2,63}$'
)

def dominio_seguro(dominio: str) -> str:
    """Revalida el dominio antes de usarlo en filesystem/subprocess.
    Lanza ValueError si no cumple el formato esperado."""
    if not dominio or len(dominio) > 253 or not PATRON_DOMINIO.match(dominio):
        raise ValueError(f"dominio no válido para uso en filesystem/subprocess: {dominio!r}")
    return dominio
    
#__________________
#core/Sql_Data.py
import logging
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
from sqlalchemy import JSON, DateTime, Integer, String
from sqlalchemy.dialects.postgresql import insert as pg_insert
from datetime import datetime, timezone
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession
from sqlalchemy.exc import IntegrityError, OperationalError, ProgrammingError

from core.config import DB_URL, SQL_ECHO, MAX_WORKERS, MAX_OVERFLOWS

logger = logging.getLogger(__name__)

class Base(DeclarativeBase):
    pass


#tabla nmap
class Socket_nmap(Base):
    __tablename__ = "nmap_socket"

    id: Mapped[int] = mapped_column(primary_key=True)

    result: Mapped[dict] = mapped_column(JSON)

    domain: Mapped[str] = mapped_column(String, unique=True, index=True)


#tabla filtro final.
class Salida_final(Base):
    __tablename__ = "salida_final_dominio"

    id: Mapped[int] = mapped_column(primary_key=True)

    result: Mapped[dict] = mapped_column(JSON)

    domain: Mapped[str] = mapped_column(String, unique=True, index=True)
    

#tablas nuclei
class Nuclei(Base):
    __tablename__ = "nuclei"

    id: Mapped[int] = mapped_column(primary_key=True)

    result: Mapped[dict] = mapped_column(JSON)

    domain: Mapped[str] = mapped_column(String, unique=True, index=True)


#tabla dead letter
class DeadLetter(Base):
    __tablename__ = "dead_letter"

    id: Mapped[int] = mapped_column(primary_key=True)

    error: Mapped[str] = mapped_column(String)

    intentos: Mapped[int] = mapped_column(Integer, default=0)

    domain: Mapped[str] = mapped_column(String)

    date: Mapped[datetime] = mapped_column(
        DateTime(timezone=True), 
        default=lambda:datetime.now(timezone.utc)
    )


class Base_Sql:
    def __init__(self):
        logger.info("Inicializando conexión a base de datos...")
        self.engine = create_async_engine(
            DB_URL,
            echo=SQL_ECHO,
            pool_size=MAX_WORKERS,
            max_overflow=MAX_OVERFLOWS,
            pool_pre_ping=True
        )
        self._session_factory = async_sessionmaker(
            self.engine,
            expire_on_commit=False
        )
        logger.info("Conexión a base de datos inicializada")
        

    #esto es para que cuando se necesite la sesion solo se llame la funcion por si se modifica async_sessionmaker
    def get_session(self) -> AsyncSession:
        """Obtiene una sesión asíncrona."""
        return self._session_factory()

        

    async def create_tables(self):
        try:
            async with self.engine.begin() as conn:
                await conn.run_sync(
                    Base.metadata.create_all
                )
                logger.info("Tablas creadas exitosamente")
        except OperationalError as e:
            logger.warning("no se pudo conectar a la base de datos %s", e)
            raise
            
        except ProgrammingError as e:
            logger.warning("error de permisos o esquema al crear tablas: %s", e)
            raise

        except Exception as e:
            logger.exception("Error inesperado creando tablas: %s", e)
            raise

    
    async def close(self):
        await self.engine.dispose()



    async def insert_dead_letter(self, error: str, intentos: int, domain: str):
        try:
            async with self._session_factory() as session:
                async with session.begin():
                    session.add(DeadLetter(
                        error=error,
                        intentos=intentos,
                        domain=domain
                    ))
                    logger.info("Dead letter insertada para %s", domain)
                
        except Exception as e:
            logger.error("Error insertando dead letter para %s: %s", domain, e)
            raise
            
    
    async def Upsert(self, model, Result, Domain):
        try:
            async with self._session_factory() as session:
                async with session.begin():
                    stmt = pg_insert(model).values(result=Result, domain=Domain)
                    stmt = stmt.on_conflict_do_update(
                        index_elements=["domain"],
                        set_={"result": Result}
                    )

                    await session.execute(stmt)
                logger.info("Datos upserted para dominio: %s", Domain)

        except IntegrityError as e:
            logger.error(f"Error de integridad al insertar {Domain}: {e}")
            raise
        except OperationalError as e:
            logger.error(f"Error operacional en BD para {Domain}: {e}")
            raise
        except Exception as e:
            logger.exception(f"Error inesperado en Upsert para {Domain}: {e}")
            raise

DB = Base_Sql()


#__________________
#core/config.py
from pathlib import Path
import os
from dotenv import load_dotenv

load_dotenv(Path(__file__).resolve().parent.parent / ".env")   # raíz del repo, no core/

def get_env(name: str, default: str | None = None) -> str:

    value = os.getenv(name, default)

    if value is None or value == "":
        raise RuntimeError(f"Falta la variable de entorno requerida: {name}")
    
    return value

SOCKET_AUTH_TOKEN = get_env("SOCKET_AUTH_TOKEN")

DOWNLOADS_DIR = Path(get_env("DOWNLOADS_DIR", "/downloads"))

DOWNLOADS_DIR.mkdir(parents=True, exist_ok=True)

DEV_COMPOSE = Path(get_env("DEV_COMPOSE_PATH", "docker-compose.dev.yml"))

TEMPORAL_HOST = get_env("TEMPORAL_HOST")

TEMPORAL_NAMESPACE = get_env("TEMPORAL_NAMESPACE")

TASK_QUEUE = get_env("TASK_QUEUE")

POSTGRES_USER = get_env("POSTGRES_USER")

POSTGRES_PASSWORD = get_env("POSTGRES_PASSWORD")

REDIS_PASSWORD = get_env("REDIS_PASSWORD")

SQL_ECHO = get_env("SQL_ECHO", "False").lower() == "true"

MAX_WORKERS = int(get_env("MAX_WORKERS"))

WORKFLOW_CONCURRENCY = int(get_env("WORKFLOW_CONCURRENCY"))

MAX_OVERFLOWS = int(get_env("MAX_OVERFLOWS"))

SERVER_HOST = get_env("SERVER_HOST")

SERVER_PORT = int(get_env("SERVER_PORT"))

SOCKET_TIMEOUT = int(get_env("SOCKET_TIMEOUT"))

TIMEOUT_CIERRE_CONEXIONES = float(get_env("TIMEOUT_CIERRE_CONEXIONES"))

SEMAFORO = int(get_env("SEMAFORO"))

NMAP_TIMEOUT = int(get_env("NMAP_TIMEOUT"))

NUCLEI_TIMEOUT = int(get_env("NUCLEI_TIMEOUT"))

MAX_MESSAGE_SIZE = int(get_env("MAX_MESSAGE_SIZE"))

REDIS_HOST = get_env("REDIS_HOST")

REDIS_PORT = int(get_env("REDIS_PORT"))

REDIS_DB = int(get_env("REDIS_DB"))

REDIS_MAX_CONNECTIONS = int(get_env("REDIS_MAX_CONNECTIONS"))

PATH_AMASS = get_env("PATH_AMASS")

PATH_NMAP = get_env("PATH_NMAP")

PATH_NUCLEI = get_env("PATH_NUCLEI")

DOMINIO_PROHIBIDO = [
    "localhost",
    "example.com",
    "example.org",
    "example.net",
]

DB_URL = f"postgresql+asyncpg://{POSTGRES_USER}:{POSTGRES_PASSWORD}@postgres:5432/proyect_framework"

MAX_CONEXIONES_POR_IP = int(get_env("MAX_CONEXIONES_POR_IP"))

LIMITE_GLOBAL_REINTENTOS = int(get_env("LIMITE_GLOBAL_REINTENTOS", "5"))
#__________________
#.env

ALLOWED_INGEST_IP=<ip real de tu equipo/servicio autorizado> #cambiar esto __________________

SOCKET_AUTH_TOKEN="9f1c2a7e4b3d8f0a6c5e2d1b9a8f7e6d5c4b3a2918f7e6d5c4b3a2918f7e6d5"

DEV_COMPOSE_PATH="docker-compose.dev.yml"

TEMPORAL_HOST=temporal:7233

TEMPORAL_NAMESPACE="default"

TASK_QUEUE="framework_2"

POSTGRES_USER=app_user

POSTGRES_PASSWORD="jeisonpostgres"

REDIS_PASSWORD="jeisonredis"

SQL_ECHO=False

MAX_WORKERS=10

WORKFLOW_CONCURRENCY=5

MAX_OVERFLOWS=20

SERVER_HOST=0.0.0.0

SERVER_PORT=9000

SOCKET_TIMEOUT=120

TIMEOUT_CIERRE_CONEXIONES=0.5

SEMAFORO=100

NMAP_TIMEOUT=120

NUCLEI_TIMEOUT=120

MAX_MESSAGE_SIZE=8192

REDIS_HOST=redis

REDIS_PORT=6379

REDIS_DB=0

REDIS_MAX_CONNECTIONS=100

#PATH_AMASS=/usr/bin/amass

#PATH_NMAP=/usr/bin/nmap

#PATH_NUCLEI=/usr/bin/nuclei

PATH_AMASS = "amass"

PATH_NUCLEI = "nuclei"

PATH_NMAP = "nmap"

MAX_CONEXIONES_POR_IP=10

LIMITE_GLOBAL_REINTENTOS=5
#__________________
# .gitignore
# .gitignore
downloads/*
!downloads/.gitkeep
.env
!.env.example
nginx/certs/
nginx/.htpasswd

#__________________
# .dockerignore
**/.env
.env.*
!.env.example
.git
__pycache__
downloads/
nginx/certs/
    
#__________________
#Dockerfile
# syntax=docker/dockerfile:1

ARG AMASS_VERSION=5.1.1
ARG NUCLEI_VERSION=3.6.2

# ===================== ETAPA 1: builder =====================
FROM python:3.12-slim AS builder

ARG AMASS_VERSION
ARG NUCLEI_VERSION

RUN DEBIAN_FRONTEND=noninteractive apt-get update && \
    apt-get install -y --no-install-recommends \
        curl \
        unzip \
        ca-certificates \
    && rm -rf /var/lib/apt/lists/*

RUN mkdir -p /tools && \
    curl -fsSL "https://[Log in to view URL]" \
    -o /tmp/amass.zip && \
    curl -fsSL "https://[Log in to view URL]" \
    -o /tmp/amass_checksums.txt && \
    grep "amass_linux_amd64.zip" /tmp/amass_checksums.txt | sed 's|amass_linux_amd64.zip|/tmp/amass.zip|' | sha256sum --check --status && \
    unzip /tmp/amass.zip -d /tmp/amass && \
    find /tmp/amass -type f -name amass -exec cp {} /tools/amass \; && \
    rm -rf /tmp/amass.zip /tmp/amass_checksums.txt /tmp/amass && \
    chmod 755 /tools/amass


RUN curl -fsSL "https://[Log in to view URL]" \
    -o /tmp/nuclei.zip && \
    curl -fsSL "https://[Log in to view URL]" \
    -o /tmp/nuclei_checksums.txt && \
    grep "nuclei_${NUCLEI_VERSION}_linux_amd64.zip" /tmp/nuclei_checksums.txt | sed "s|nuclei_${NUCLEI_VERSION}_linux_amd64.zip|/tmp/nuclei.zip|" | sha256sum --check --status && \
    unzip /tmp/nuclei.zip -d /tmp/nuclei && \
    mv /tmp/nuclei/nuclei /tools/nuclei && \
    rm -rf /tmp/nuclei.zip /tmp/nuclei_checksums.txt /tmp/nuclei && \
    chmod 755 /tools/nuclei

WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir --upgrade pip && \
    pip install --no-cache-dir --user -r requirements.txt

# ===================== ETAPA 2: runtime =====================
FROM python:3.12-slim AS runtime

LABEL org.opencontainers.image.authors="equipo@empresa.com"
LABEL org.opencontainers.image.version="1.0"

ENV PYTHONDONTWRITEBYTECODE=1 \
    PYTHONUNBUFFERED=1 \
    PATH="/tools:/home/app/.local/bin:${PATH}"

RUN DEBIAN_FRONTEND=noninteractive apt-get update && \
    apt-get install -y --no-install-recommends \
        nmap \
        curl \
        ca-certificates \
    && rm -rf /var/lib/apt/lists/*

RUN addgroup --system app && adduser --system --group app

WORKDIR /app

# Copiamos solo lo ya listo desde la etapa builder — sin curl, unzip, tar.gz, ni caché de apt
COPY --from=builder /tools /tools
COPY --from=builder /root/.local /home/app/.local
COPY --chown=app:app . .

RUN chown -R app:app /home/app/.local

USER app

EXPOSE 8000
EXPOSE 9000

HEALTHCHECK --interval=30s --timeout=5s --retries=3 \
    CMD curl -f http://localhost:8000/health || exit 1

CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]


#__________________
#requirements.txt
fastapi==0.136.3
uvicorn[standard]==0.34.0
typer==0.15.1
sqlalchemy[asyncio]==2.0.49
asyncpg==0.30.0
redis==5.2.1
temporalio==1.33.0
dnspython==2.7.0
xmltodict==0.14.2
aiofiles==24.1.0
python-dotenv==1.0.1
langgraph==1.2.4
pytest==8.3.3
pytest-asyncio==0.24.0
httpx==0.27.2

#_______________________________
#tests/test_health.py
import pytest
from httpx import AsyncClient, ASGITransport
from app.main import app

@pytest.mark.asyncio
async def test_health():
    transport = ASGITransport(app=app)
    async with AsyncClient(transport=transport, base_url="http://test") as client:
        resp = await client.get("/health")
        assert resp.status_code == 200
        assert resp.json() == {"status": "ok"}
#__________________
#downloads
#carpeta donde van los archivos JSON
#__________________
#.github/workflows/ci-cd.yml
name: CI/CD

on:
  push:
    branches: [main]
  pull_request:
    branches: [main]

env:
  IMAGE_NAME: ${{ github.repository }}

jobs:

  # ===================== CI: Build + Test =====================
  build-and-test:
    runs-on: ubuntu-latest
    steps:
      - name: Checkout código
        uses: actions/checkout@v4

      - name: Configurar Docker Buildx
        uses: docker/setup-buildx-action@v3

      - name: Construir imagen (solo validación, sin publicar)
        uses: docker/build-push-action@v6
        with:
          context: .
          push: false
          tags: ${{ env.IMAGE_NAME }}:test
          cache-from: type=gha
          cache-to: type=gha,mode=max

      - name: Crear .env para pruebas
        run: |
          cat > .env <<EOF
          DEV_COMPOSE_PATH=docker-compose.dev.yml
          TEMPORAL_HOST=temporal:7233
          TEMPORAL_NAMESPACE=default
          TASK_QUEUE=framework_2
          POSTGRES_USER=ci_user
          POSTGRES_PASSWORD=ci_password_temporal
          REDIS_PASSWORD=ci_redis_temporal
          SQL_ECHO=False
          MAX_WORKERS=10
          MAX_OVERFLOWS=20
          SERVER_HOST=0.0.0.0
          SERVER_PORT=9000
          SOCKET_TIMEOUT=120
          TIMEOUT_CIERRE_CONEXIONES=0.5
          SEMAFORO=100
          NMAP_TIMEOUT=120
          NUCLEI_TIMEOUT=120
          MAX_MESSAGE_SIZE=8192
          REDIS_HOST=redis
          REDIS_PORT=6379
          REDIS_DB=0
          REDIS_MAX_CONNECTIONS=100
          #PATH_AMASS=/usr/bin/amass
          #PATH_NMAP=/usr/bin/nmap
          #PATH_NUCLEI=/usr/bin/nuclei

          PATH_AMASS = "amass"
          PATH_NUCLEI = "nuclei"
          PATH_NMAP = "nmap"
          DOWNLOADS_DIR=/downloads
          WORKFLOW_CONCURRENCY=5
          SOCKET_AUTH_TOKEN=9f1c2a7e4b3d8f0a6c5e2d1b9a8f7e6d5c4b3a2918f7e6d5c4b3a2918f7e6d5
          ALLOWED_INGEST_IP=0.0.0.0/0 #SOLO EN CI NUNCA EN PRODUCCIÓN!!!
          MAX_CONEXIONES_POR_IP=10
          EOF

      - name: Levantar stack para pruebas
        run: docker compose -f docker-compose.dev.yml up -d --build
        env:
          POSTGRES_USER: ci_user
          POSTGRES_PASSWORD: ci_password_temporal
          REDIS_PASSWORD: ci_redis_temporal

      - name: Esperar a que la app esté healthy
        run: |
          for i in {1..15}; do
            status=$(docker inspect --format='{{.State.Health.Status}}' $(docker compose -f docker-compose.dev.yml ps -q app) 2>/dev/null || echo "starting")
            if [ "$status" == "healthy" ]; then
              echo "App healthy"
              exit 0
            fi
            echo "Esperando... ($status)"
            sleep 5
          done
          echo "La app no llegó a healthy a tiempo"
          docker compose -f docker-compose.dev.yml logs
          exit 1

      - name: Ejecutar tests
        run: docker compose -f docker-compose.dev.yml exec -T app pytest

      - name: Apagar stack
        if: always()
        run: docker compose -f docker-compose.dev.yml down -v

  # ===================== CD: Build + Push + Deploy (solo en main) =====================
  build-push-deploy:
    needs: build-and-test
    if: github.ref == 'refs/heads/main' && github.event_name == 'push'
    runs-on: ubuntu-latest
    steps:
      - name: Checkout código
        uses: actions/checkout@v4

      - name: Login al registry
        uses: docker/login-action@v3
        with:
          registry: ghcr.io
          username: ${{ github.actor }}
          password: ${{ secrets.GITHUB_TOKEN }}

      - name: Configurar Docker Buildx
        uses: docker/setup-buildx-action@v3

      - name: Construir y publicar imagen
        uses: docker/build-push-action@v6
        with:
          context: .
          push: true
          tags: |
            ghcr.io/${{ env.IMAGE_NAME }}:latest
            ghcr.io/${{ env.IMAGE_NAME }}:${{ github.sha }}
          cache-from: type=gha
          cache-to: type=gha,mode=max

      - name: Desplegar en el servidor por SSH
        uses: appleboy/ssh-action@v1.0.3
        with:
          host: ${{ secrets.DEPLOY_HOST }}
          username: ${{ secrets.DEPLOY_USER }}
          key: ${{ secrets.DEPLOY_SSH_KEY }}
          script: |
            cd /ruta/a/tu/proyecto #aqui colocar ruta
            git pull origin main
            docker compose pull
            docker compose up -d --build
            docker image prune -f
#__________________
# docker-compose.yml
#este es para producción 
services:

  # ===================== REVERSE PROXY (único punto expuesto) =====================
  nginx:
    image: nginx:1.27-alpine
    entrypoint: ["/bin/sh", "/entrypoint.sh"]
    volumes:
      - ./nginx/templates/nginx.conf.template:/etc/nginx/templates/nginx.conf.template:ro
      - ./nginx/entrypoint.sh:/entrypoint.sh:ro
      - ./nginx/.htpasswd:/etc/nginx/.htpasswd:ro
      - ./nginx/certs:/etc/nginx/certs:ro
    environment:
      ALLOWED_INGEST_IP: ${ALLOWED_INGEST_IP}
    restart: unless-stopped
    ports:
      - "80:80"
      - "9000:9000"
    networks:
      - public
      - internal
    depends_on:
      app:
        condition: service_healthy
    logging:
      driver: "json-file"
      options:
        max-size: "10m"
        max-file: "3"

  # ===================== APP (imagen ya construida y publicada en ghcr.io) =====================
  app:
    image: ghcr.io/TU_USUARIO/TU_REPOSITORIO:latest #esto se cambia _______________________________$$$
    restart: unless-stopped
    expose:
      - "8000"
      - "9000"
    env_file:
      - .env
    volumes:
    - downloads_data:/downloads
    environment:
      REDIS_HOST: redis
      TEMPORAL_HOST: temporal:7233
    networks:
      - internal
      - egress
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_healthy
      temporal:
        condition: service_healthy
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
      interval: 30s
      timeout: 5s
      retries: 3

  # ===================== POSTGRES (sin puertos publicados) =====================
  postgres:
    image: postgres:16-alpine
    restart: unless-stopped
    env_file:
      - .env
    environment:
      POSTGRES_USER: ${POSTGRES_USER}
      POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
      POSTGRES_DB: proyect_framework
    volumes:
      - postgres_data:/var/lib/postgresql/data
      - ./postgres-init:/docker-entrypoint-initdb.d:ro
    networks:
      - internal
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER}"]
      interval: 10s
      timeout: 5s
      retries: 5

  # ===================== REDIS (con password, sin puertos publicados) =====================
  redis:
    image: redis:7-alpine
    restart: unless-stopped
    command: ["redis-server", "--requirepass", "${REDIS_PASSWORD}"]
    env_file:
      - .env
    volumes:
      - redis_data:/data
    networks:
      - internal
    healthcheck:
      test: ["CMD", "redis-cli", "-a", "${REDIS_PASSWORD}", "ping"]
      interval: 10s
      timeout: 5s
      retries: 5

  # ===================== TEMPORAL (sin puertos publicados) =====================
  temporal:
    image: temporalio/auto-setup:1.24
    restart: unless-stopped
    environment:
      DB: postgres12
      DB_PORT: 5432
      POSTGRES_USER: ${POSTGRES_USER}
      POSTGRES_PWD: ${POSTGRES_PASSWORD}
      POSTGRES_SEEDS: postgres
      DBNAME: temporal
      VISIBILITY_DBNAME: temporal_visibility
    networks:
      - internal
    depends_on:
      postgres:
        condition: service_healthy
    healthcheck:
      test: ["CMD", "tctl", "--address", "temporal:7233", "cluster", "health"]
      interval: 15s
      timeout: 5s
      retries: 5

  # ===================== TEMPORAL UI (solo vía nginx, con autenticación) =====================
  temporal-ui:
    image: temporalio/ui:2.31.2
    restart: unless-stopped
    environment:
      TEMPORAL_ADDRESS: temporal:7233
    networks:
      - internal
    depends_on:
      temporal:
        condition: service_healthy

networks:
  public:
    driver: bridge
  internal:
    driver: bridge
    internal: true
  egress:
    driver: bridge

volumes:
  postgres_data:
  redis_data:
  downloads_data:


#_______________________________
# nginx/entrypoint.sh
#!/bin/sh
#este archivo es para el compose de produccióna

set -e

case "$ALLOWED_INGEST_IP" in
  *"<"*|*">"*|"")
    echo "ERROR: ALLOWED_INGEST_IP no fue configurado con un valor real (${ALLOWED_INGEST_IP})" >&2
    exit 1
    ;;
esac

envsubst '${ALLOWED_INGEST_IP}' < /etc/nginx/templates/nginx.conf.template > /etc/nginx/nginx.conf
exec nginx -g "daemon off;"

#__________________
# docker-compose.dev.yml
#este es para probar local
services:

  # ===================== REVERSE PROXY =====================
  nginx:
    image: nginx:1.27-alpine
    restart: unless-stopped
    ports:
      - "80:80"
      - "9000:9000"
    volumes:
      - ./nginx/templates/nginx.dev.conf.template:/etc/nginx/templates/nginx.dev.conf.template:ro
      - ./nginx/entrypoint.sh:/entrypoint.sh:ro
      - ./nginx/.htpasswd:/etc/nginx/.htpasswd:ro
    entrypoint: ["/bin/sh", "-c", "envsubst < /etc/nginx/templates/nginx.dev.conf.template > /etc/nginx/nginx.conf && exec nginx -g 'daemon off;'"]
    networks:
      - public
      - internal
    depends_on:
      app:
        condition: service_healthy

  # ===================== APP (construida localmente, no descargada) =====================
  app:
    build:
      context: .
      dockerfile: Dockerfile
    restart: unless-stopped
    expose:
      - "8000"
      - "9000"
    env_file:
      - .env
    environment:
      REDIS_HOST: redis
      TEMPORAL_HOST: temporal:7233
    networks:
      - internal
      - egress
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_healthy
      temporal:
        condition: service_healthy
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
      interval: 30s
      timeout: 5s
      retries: 3
    volumes:
      - ./app:/app/app
      - ./core:/app/core
      - ./tools:/app/tools
      - ./temp_workflows:/app/temp_workflows
      - ./graph:/app/graph
      - downloads_data:/downloads   # monta tu código en vivo, para no reconstruir en cada cambio

  # ===================== POSTGRES =====================
  postgres:
    image: postgres:16-alpine
    restart: unless-stopped
    env_file:
      - .env
    environment:
      POSTGRES_USER: ${POSTGRES_USER}
      POSTGRES_PASSWORD: ${POSTGRES_PASSWORD}
      POSTGRES_DB: proyect_framework
    volumes:
      - postgres_data_dev:/var/lib/postgresql/data
      - ./postgres-init:/docker-entrypoint-initdb.d:ro
    networks:
      - internal
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER}"]
      interval: 10s
      timeout: 5s
      retries: 5

  # ===================== REDIS =====================
  redis:
    image: redis:7-alpine
    restart: unless-stopped
    command: ["redis-server", "--requirepass", "${REDIS_PASSWORD}"]
    env_file:
      - .env
    volumes:
      - redis_data_dev:/data
    networks:
      - internal
    healthcheck:
      test: ["CMD", "redis-cli", "-a", "${REDIS_PASSWORD}", "ping"]
      interval: 10s
      timeout: 5s
      retries: 5

  # ===================== TEMPORAL =====================
  temporal:
    image: temporalio/auto-setup:1.24
    restart: unless-stopped
    environment:
      DB: postgres12
      DB_PORT: 5432
      POSTGRES_USER: ${POSTGRES_USER}
      POSTGRES_PWD: ${POSTGRES_PASSWORD}
      POSTGRES_SEEDS: postgres
      DBNAME: temporal
      VISIBILITY_DBNAME: temporal_visibility
    networks:
      - internal
    depends_on:
      postgres:
        condition: service_healthy
    healthcheck:
      test: ["CMD", "tctl", "--address", "temporal:7233", "cluster", "health"]
      interval: 15s
      timeout: 5s
      retries: 5

  # ===================== TEMPORAL UI =====================
  temporal-ui:
    image: temporalio/ui:2.31.2
    restart: unless-stopped
    environment:
      TEMPORAL_ADDRESS: temporal:7233
    networks:
      - internal
    depends_on:
      temporal:
        condition: service_healthy

networks:
  public:
    driver: bridge
  internal:
    driver: bridge
  egress:
    driver: bridge


volumes:
  postgres_data_dev:
  redis_data_dev:
  downloads_data:



#__________________
#nginx/templates/nginx.conf.template

#esto genera los certificados de ssl_certificate y 
#ssl_certificate_key. comando para terminal

#openssl req -x509 -nodes -newkey rsa:2048 \
  #-keyout nginx/certs/server.key -out nginx/certs/server.crt \
  #-days 825 -subj "/CN=ingest.miempresa.internal"

user nginx;
worker_processes auto;

events {
    worker_connections 1024;
}


stream {
    server {
        listen 9000 ssl;
        ssl_certificate     /etc/nginx/certs/server.crt;
        ssl_certificate_key /etc/nginx/certs/server.key;

        allow ${ALLOWED_INGEST_IP};
        deny all;

        proxy_pass app:9000;
        proxy_timeout 130s;
        proxy_connect_timeout 10s;
    }
}

http {
    include       /etc/nginx/mime.types;
    default_type  application/octet-stream;

    sendfile        on;
    keepalive_timeout 65;

    access_log /var/log/nginx/access.log;
    error_log  /var/log/nginx/error.log warn;

    server_tokens off;
    client_max_body_size 10M;

    # ===================== App principal (acceso por IP, sin TLS) =====================
    server {
        listen 80 default_server;
        server_name _;

        add_header X-Frame-Options "SAMEORIGIN" always;
        add_header X-Content-Type-Options "nosniff" always;
        add_header Referrer-Policy "strict-origin-when-cross-origin" always;

        location / {
            proxy_pass http://app:8000;
            proxy_http_version 1.1;

            proxy_set_header Host $host;
            proxy_set_header X-Real-IP $remote_addr;
            proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
            proxy_set_header X-Forwarded-Proto $scheme;

            proxy_set_header Upgrade $http_upgrade;
            proxy_set_header Connection "upgrade";

            proxy_read_timeout 60s;
            proxy_connect_timeout 10s;
        }

        location /health {
            proxy_pass http://app:8000/health;
            access_log off;
        }

        # Temporal UI, protegido con autenticación básica, bajo una ruta distinta
        location /temporal/ {
            auth_basic "Acceso restringido";
            auth_basic_user_file /etc/nginx/.htpasswd;

            rewrite ^/temporal/(.*)$ /$1 break;
            proxy_pass http://temporal-ui:8080;
            proxy_http_version 1.1;

            proxy_set_header Host $host;
            proxy_set_header X-Real-IP $remote_addr;
            proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
            proxy_set_header X-Forwarded-Proto $scheme;

            proxy_set_header Upgrade $http_upgrade;
            proxy_set_header Connection "upgrade";
        }
    }
}

#esto sudo apt-get install apache2-utils
#htpasswd -c ./nginx/.htpasswd admin
#crea Cómo generar el nginx/.htpasswd 
#(comando, no archivo de contenido fijo) del archivo nginx

#_______________________________
#nginx/templates/nginx.dev.conf.template

user nginx;
worker_processes auto;

events {
    worker_connections 1024;
}

stream {
    server {
        listen 9000;
        allow 127.0.0.1;   # cambiar a ip local real
        deny all;
        proxy_pass app:9000;
        proxy_timeout 130s;
        proxy_connect_timeout 10s;
    }
}


http {
    include       /etc/nginx/mime.types;
    default_type  application/octet-stream;

    sendfile        on;
    keepalive_timeout 65;

    access_log /var/log/nginx/access.log;
    error_log  /var/log/nginx/error.log warn;

    server_tokens off;
    client_max_body_size 10M;

    # ===================== App principal (acceso por IP, sin TLS) =====================
    server {
        listen 80 default_server;
        server_name _;

        add_header X-Frame-Options "SAMEORIGIN" always;
        add_header X-Content-Type-Options "nosniff" always;
        add_header Referrer-Policy "strict-origin-when-cross-origin" always;

        location / {
            proxy_pass http://app:8000;
            proxy_http_version 1.1;

            proxy_set_header Host $host;
            proxy_set_header X-Real-IP $remote_addr;
            proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
            proxy_set_header X-Forwarded-Proto $scheme;

            proxy_set_header Upgrade $http_upgrade;
            proxy_set_header Connection "upgrade";

            proxy_read_timeout 60s;
            proxy_connect_timeout 10s;
        }

        location /health {
            proxy_pass http://app:8000/health;
            access_log off;
        }

        # Temporal UI, protegido con autenticación básica, bajo una ruta distinta
        location /temporal/ {
            auth_basic "Acceso restringido";
            auth_basic_user_file /etc/nginx/.htpasswd;

            rewrite ^/temporal/(.*)$ /$1 break;
            proxy_pass http://temporal-ui:8080;
            proxy_http_version 1.1;

            proxy_set_header Host $host;
            proxy_set_header X-Real-IP $remote_addr;
            proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
            proxy_set_header X-Forwarded-Proto $scheme;

            proxy_set_header Upgrade $http_upgrade;
            proxy_set_header Connection "upgrade";
        }
    }
}


#__________________
#postgres-init:/docker-entrypoint-initdb
#!/bin/bash
set -e

psql -v ON_ERROR_STOP=1 --username "$POSTGRES_USER" <<-EOSQL
    CREATE DATABASE temporal;
    CREATE DATABASE temporal_visibility;

EOSQL






Embed on website

To embed this project on your website, copy the following code and paste it into your website's HTML: