#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__))
os._exit(1)
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
import re
from functools import partial
import logging
import redis.exceptions as redis_exceptions
from core.config import(
SERVER_HOST,
SERVER_PORT,
SOCKET_TIMEOUT,
MAX_MESSAGE_SIZE,
DOMINIO_PROHIBIDO,
TIMEOUT_CIERRE_CONEXIONES,
SEMAFORO,
SOCKET_AUTH_TOKEN
)
PATRON_DOMINIO = re.compile(
r'^(?!-)[A-Za-z0-9-]{1,63}(?<!-)(\.[A-Za-z0-9-]{1,63}(?<!-))*\.[A-Za-z]{2,63}$'
)
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.TOKEN_LEN = 64
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):
if self.semaforo.locked():
writer.close()
await writer.wait_closed()
return
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)
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 hmac.compare_digest(token, SOCKET_AUTH_TOKEN.encode()):
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.xadd("cola:dominios", {"dominio": dominio_data})
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
return bool(PATRON_DOMINIO.match(dominio))
#__________________
#core/redis_.py
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__)
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í
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):
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 TEMPORAL_HOST, TASK_QUEUE, TEMPORAL_NAMESPACE
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__)
async def Init_client(self):
try:
self.client = await Client.connect(TEMPORAL_HOST, namespace=TEMPORAL_NAMESPACE)
self.logger.info("Cliente conectado a Temporal en %s", TEMPORAL_HOST)
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 _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 _close_client(self):
if self.client is None:
return
try:
if hasattr(self.client, "close"):
await self.client.close()
except Exception as e:
self.logger.warning("error al cerrar cliente temporal (ignorado): %s", e)
finally:
self.client = None
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()
#__________________
#temp_workflows/Workflow_.py
import logging
import asyncio
import uuid
from core.redis_ import REDIS_CONNECT
from datetime import timedelta
from temporalio.common import RetryPolicy
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
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.max_workflow_attemp = 3
async def connect_client_workflow(self):
try:
self.client = await Client.connect(TEMPORAL_HOST, namespace=TEMPORAL_NAMESPACE)
self.logger.info("Cliente conectado a Temporal en %s", TEMPORAL_HOST)
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 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 _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.r_connect.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.r_connect.xack(stream, grupo, message_id)
continue
await redis_wrapper.r_connect.xadd(stream, {"dominio": dominio_data})
await redis_wrapper.r_connect.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.r_connect)
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.connect_client_workflow()
except Exception:
self.logger.warning("no se pudo conectar, reintentando en 5s")
await asyncio.sleep(5)
continue
try:
messages = await redis_wrapper.r_connect.xreadgroup(
groupname="workers",
consumername=consumer,
streams={"cola:dominios": ">"},
count=1,
block=0
)
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
stream_name, entries = messages[0]
if not entries:
self.logger.warning("stream sin entradas, reintentando")
continue
message_id, data = entries[0]
dominio_data = data["dominio"]
intentos = 0
while True:
if self.client is None:
try:
await self.connect_client_workflow()
except Exception:
await asyncio.sleep(5)
continue
error_except = None
try:
handle = await self.client.start_workflow(
ReconWorkflow.run,
dominio_data,
id=f"mainworkflow-{uuid.uuid4()}",
task_queue=TASK_QUEUE,
execution_timeout=timedelta(minutes=30),
retry_policy=RetryPolicy(
maximum_attempts=5,
initial_interval=timedelta(seconds=10))
)
resultado = await handle.result()
await redis_wrapper.r_connect.xack("cola:dominios","workers",message_id)
self.logger.info("workflow completado para %s", dominio_data)
self.logger.info("handle del workflow: %s", handle.id)
break
except RPCError as e:
self.logger.error("no se pudo comunicar con el servidor temporal: %s", e)
error_except = e
intentos += 1
if self.client is not None:
try:
await self.client.close()
except Exception as close_err:
self.logger.warning("no se pudo cerrar la conexión: %s", close_err)
finally:
self.client = None
except WorkflowFailureError as e:
causa = e.cause
if isinstance(causa, TimeoutError):
self.logger.error("el workflow para %s excedió el tiempo máximo (execution_timeout): %s",dominio_data, causa)
else:
self.logger.error("el workflow para %s falló: %s", dominio_data, causa)
intentos += 1
error_except = e
except asyncio.CancelledError as e:
self.logger.info("procesamiento de dominio cancelado")
try:
await asyncio.shield(
DB.insert_dead_letter(str(e), intentos, "guardando datos cancelled")
)
except Exception:
self.logger.exception("error al guardar dato en SQL")
raise
except Exception as e:
self.logger.exception("error al ejecutar workflow para %s: %s", dominio_data, e)
intentos += 1
error_except = e
if intentos >= self.max_workflow_attemp:
try:
await DB.insert_dead_letter(str(error_except) if error_except else "error desconocido", intentos, dominio_data)
except Exception:
self.logger.warning("no se pudo guardar datos de error")
self.logger.warning("cantidad de intentos superada: %d/%d", intentos, self.max_workflow_attemp)
await redis_wrapper.r_connect.xack("cola:dominios", "workers", message_id)
#await redis_wrapper.r_connect.xadd("cola:dominios", {"dominio": dominio_data})
break
if intentos < self.max_workflow_attemp:
await asyncio.sleep(5)
#__________________
#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_data: str) -> dict:
return await workflow.execute_activity(
sample_activity,
dominio_data,
schedule_to_close_timeout=timedelta(minutes=25),
schedule_to_start_timeout=timedelta(minutes=5),
start_to_close_timeout=timedelta(minutes=20),
heartbeat_timeout=timedelta(seconds=30),
retry_policy=RetryPolicy(maximum_attempts=5),
cancellation_type=activity.ActivityCancellationType.TRY_CANCEL,
task_queue=TASK_QUEUE,
activity_id=f"graph-{uuid.uuid4()}",
versioning_intent=None,
summary="ejecución del grafo que contiene las herramientas"
)
#__________________
#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_data: str):
try:
result = await graph_ainvoke(dominio_data)
return result
except Exception:
logger.exception("error con graph_start")
raise
#__________________
#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, 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 analisis para nuclei en %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
logger = logging.getLogger(__name__)
async def analizer_amass(dominio_url: str) -> dict:
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, Socket_nmap
from core.config import PATH_NMAP, NMAP_TIMEOUT
logger = logging.getLogger(__name__)
db = DB
async def analizer_nmap(dominios_y_ip: dict) -> dict:
dict_resultados = defaultdict(dict)
max_intentos = 3
if not dominios_y_ip:
raise ValueError("datos de entrada para nmap vacíos")
folder_temp = Path(tempfile.gettempdir())
for dominio, lista_ip in dominios_y_ip.items():
for ip in lista_ip:
file_nmap = folder_temp / f"nmap_{dominio}_{ip}_{uuid.uuid4().hex[:8]}.xml"
intento = 0
bucle = True
while bucle:
intento += 1
try:
resultado = await asyncio.create_subprocess_exec(
PATH_NMAP, ip, "-oX", file_nmap,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE
)
try:
stdout, stderr = await asyncio.wait_for(
resultado.communicate(),
timeout= NMAP_TIMEOUT
)
except (asyncio.TimeoutError, TimeoutError):
logger.error("nmap excedio el tiempo limite de python")
resultado.kill()
await resultado.wait()
raise
if resultado.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())
parser = parsear_salida(file_nmap)
logger.info("Dominio %s - IP %s escaneado correctamente", dominio, ip)
try:
await db.Upsert(Socket_nmap, parser, ip)
logger.info("guardado correctamente")
except Exception:
logger.error("Error al guardar datos para %s - %s", dominio, ip)
dict_resultados[dominio][ip] = parser
break
except FileNotFoundError:
logger.error("Nmap no está instalado o no está en PATH")
raise
except (asyncio.TimeoutError, TimeoutError):
logger.warning("Timeout escaneando %s, reintento %d/%d", ip, intento, max_intentos)
if intento >= max_intentos:
bucle = False
except (PermissionError, OSError, ValueError, TypeError) as e:
logger.error(f"Error no recuperable: {e}")
bucle = False
except asyncio.CancelledError:
logger.error("La operación fue cancelada")
try:
if file_nmap.exists():
file_nmap.unlink()
except Exception as cleanup_error:
logger.error("Error en limpieza: %s", cleanup_error)
raise
except Exception:
logger.exception("error inesperado")
bucle = False
if not bucle and ip not in dict_resultados.get(dominio, {}):
try:
await DB.insert_dead_letter(
f"nmap falló para {dominio} - {ip}", intento, dominio
)
except Exception:
logger.warning("no se pudo guardar dead letter de nmap para %s", dominio)
return dict(dict_resultados)
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.config import PATH_NUCLEI, NUCLEI_TIMEOUT
logger = logging.getLogger(__name__)
async def analizer_nuclei(dominio_ips: dict) -> dict:
dict_final = defaultdict(dict)
max_intentos = 3
if not dominio_ips:
logger.warning("entrada con datos vacios")
raise ValueError("datos recibidos por resolver_dns vacios")
folder_file = tempfile.gettempdir()
for dominio in dominio_ips.keys():
if dominio.startswith(("http://", "https://")):
url = dominio
else:
url = f"https://{dominio}"
intento = 0
nuclei_file = Path(f"{folder_file}/{dominio}_{uuid.uuid4().hex[:8]}.json")
bucle = True
while bucle:
intento += 1
try:
resultado = await asyncio.create_subprocess_exec(
PATH_NUCLEI, "-u", url, "-json-export", nuclei_file,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE
)
try:
stdout, stderr = await asyncio.wait_for(
resultado.communicate(),
timeout=NUCLEI_TIMEOUT
)
except (asyncio.TimeoutError, TimeoutError):
logger.error("nuclei excedio el tiempo limite de python")
resultado.kill()
await resultado.wait()
raise
if resultado.returncode != 0:
msge = stderr.decode() if stderr else "codigo de retorno no es 0"
raise RuntimeError(f"error en la salida de subprocess: {msge}")
if stderr:
logger.warning("nuclei emitió advertencias: %s", stderr.decode())
parser = await parsear_salida(nuclei_file, dominio)
try:
await DB.Upsert(Nuclei, parser, dominio)
logger.info("datos de %s guardado correctamente", dominio)
except Exception as e:
logger.error("error al guardar datos de %s: %s", dominio, e)
dict_final[dominio] = parser
break
except FileNotFoundError:
logger.error("Nuclei no está instalado o no está en PATH")
raise
except (asyncio.TimeoutError, TimeoutError):
logger.warning("Timeout escaneando %s, reintento %d/%d", dominio, intento, max_intentos)
if intento >= max_intentos:
bucle = False
except (PermissionError, OSError, ValueError, TypeError) as e:
logger.error(f"Error no recuperable: {e}")
bucle = False
except asyncio.CancelledError:
logger.error("La operación fue cancelada")
try:
if nuclei_file.exists():
nuclei_file.unlink()
except Exception as cleanup_error:
logger.error("Error en limpieza: %s", cleanup_error)
raise
except Exception:
logger.exception("error inesperado")
bucle = False
#revisar
if not bucle and dominio not in dict_final:
try:
await DB.insert_dead_letter(
f"nuclei falló para {dominio}", intento, dominio
)
except Exception:
logger.warning("no se pudo guardar dead letter de nuclei para %s", dominio)
return dict(dict_final)
async def parsear_salida(ruta, dominio):
if not ruta.exists():
logger.warning("no esxite archivo nuclei de %s del dominio: %s", ruta, dominio)
raise FileNotFoundError(f"archivo nuclei {ruta} no encontrado del dominio: {dominio}")
try:
with open(ruta, "r") as f:
contenido = f.read()
if not contenido.strip():
logger.warning("el archivo %s está vacío", ruta)
raise ValueError(f"archivo {ruta} sin datos para procesar")
data = json.loads(contenido)
if not data:
logger.warning("el archivo %s no contiene datos útiles tras parsear", ruta)
raise ValueError(f"archivo {ruta} sin datos útiles")
return data
finally:
try:
if ruta.exists():
ruta.unlink()
logger.info("archivo temporal %s eliminado", ruta)
except Exception as e:
logger.warning("no se pudo eliminar el archivo %s: %s", ruta, 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__)
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] = {
"ip": "",
"data": [],
"url_login": [],
"descriptions": [],
}
nmap_data = {}
puertos = []
try:
nmap_por_ip = resultado_nmap.get(d) or {}
primera_ip = next(iter(nmap_por_ip), None)
if primera_ip is None:
nmap_data, puertos = {}, []
else:
nmap_data = ((nmap_por_ip[primera_ip].get("nmap") or {}).get("host") or {})
if not nmap_data:
logger.info("datos nmap de %s vacios", d)
else:
address = nmap_data.get("address") or {}
if isinstance(address, dict):
address = [address]
ipv4 = next((a for a in address if a.get("@addrtype") == "ipv4"), None)
if ipv4 is None:
ipv4 = address[0] if address else {}
formato_dict[d]["ip"] = ipv4.get("@addr", "")
puertos = (nmap_data.get("ports") or {}).get("port") or []
if isinstance(puertos, dict):
puertos = [puertos]
elif not isinstance(puertos, list):
logger.warning("formato de puertos inesperado en %s: %r", d, puertos)
puertos = []
except (AttributeError, TypeError, KeyError, IndexError) as e:
logger.warning("error procesando IP de %s: %s", d, e)
nmap_data, puertos = {}, []
except Exception as e:
logger.exception("error inesperado procesando nmap de %s: %s", d, e)
nmap_data, puertos = {}, []
for port in puertos:
try:
if not port:
logger.info("datos de puertos de %s vacios", d)
continue
estado = (port.get("state") or {}).get("@state")
if estado != "open":
logger.debug("puerto %s/%s de %s no está abierto (%s)", port.get("@portid"), port.get("@protocol"), d, estado)
continue
service = port.get("service") or {}
formato_dict[d]["data"].append({
"puerto": port.get("@portid"),
"protocolo": port.get("@protocol"),
"servicio": service.get("@name"),
"tecnologia": service.get("@product"),
})
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)
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")
except AttributeError:
logger.warning("item de nuclei inválido en %s", d)
continue
if not url:
logger.debug("item de nuclei en %s sin 'matched-at', se omite", d)
continue
path = urlparse(url).path.lower()
if (any(ext in path for ext in [".php", ".asp", ".jsp"]) or
any(seg in path for seg in ["login", "admin", "xmlrpc"])):
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/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__).with_name(".env"))
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_CONEXIONES_ACTIVAS_START_SERVER = float(get_env("TIMEOUT_CONEXIONES_ACTIVAS_START_SERVER"))
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"
#__________________
# core/.env
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
#__________________
# .gitignore
downloads/*
!downloads/.gitkeep
.env
#__________________
#.dockerignore
.env
.git
__pycache__
downloads/
#__________________
#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
DOWNLOADS_DIR=/downloads
WORKFLOW_CONCURRENCY=5
SOCKET_AUTH_TOKEN=9f1c2a7e4b3d8f0a6c5e2d1b9a8f7e6d5c4b3a2918f7e6d5c4b3a2918f7e6d5
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
restart: unless-stopped
ports:
- "80:80"
- "9000:9000"
volumes:
- ./nginx/nginx.conf:/etc/nginx/nginx.conf:ro
- ./nginx/.htpasswd:/etc/nginx/.htpasswd:ro
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:
#__________________
# docker-compose.dev.yml
#este es para probar local
services:
# ===================== REVERSE PROXY =====================
nginx:
image: nginx:1.27-alpine
restart: unless-stopped
ports:
- "80:80"
volumes:
- ./nginx/nginx.conf:/etc/nginx/nginx.conf:ro
- ./nginx/.htpasswd:/etc/nginx/.htpasswd:ro
networks:
- public
- internal
depends_on:
app:
condition: service_healthy
# ===================== APP (construida localmente, no descargada) =====================
app:
build:
context: .
dockerfile: Dockerfile
restart: unless-stopped
expose:
- "8000"
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 # monta tu código en vivo, para no reconstruir en cada cambio
- downloads_data:/downloads
# ===================== 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/nginx.conf
user nginx;
worker_processes auto;
events {
worker_connections 1024;
}
stream {
server {
listen 9000;
allow 203.0.113.10; # revisar esta 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
#__________________
#postgres-init:/docker-entrypoint-initdb
#!/bin/bash
set -e
psql -v ON_ERROR_STOP=1 --username "$POSTGRES_USER" <<-EOSQL
CREATE DATABASE proyect_framework;
CREATE DATABASE temporal;
CREATE DATABASE temporal_visibility;
EOSQL
To embed this project on your website, copy the following code and paste it into your website's HTML: