#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):
if isinstance(resultado, Exception):
logger.error("lote de nuclei falló: %s", 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)
if not dict_final:
raise RuntimeError("nuclei falló para todos los dominios")
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
To embed this project on your website, copy the following code and paste it into your website's HTML: