Skip to content

Repository files navigation

⚙️ flujo

Tests

Un motor de workflows de automatización construido desde cero. Defines el pipeline en YAML — pasos HTTP, transformaciones, condiciones, IA — y el motor lo ejecuta: resuelve el orden por dependencias (DAG), reintenta lo que falla con backoff exponencial, toma ramas condicionales y deja un diario auditable de cada corrida.

¿Por qué construir esto? En Manto automatizo procesos para pymes con n8n todos los días. Este proyecto es la respuesta a la pregunta "¿y qué hay dentro de una herramienta así?": un motor mínimo pero completo, con las mismas piezas que usa cualquier orquestador serio — grafo de dependencias, política de reintentos, plantillas, plugins y observabilidad.

Un workflow en YAML

nombre: dolar-alerta
descripcion: Consulta el dolar oficial (BCRP) y alerta si supera el umbral.

entradas:
  umbral: 3.50

pasos:
  - id: consultar_bcrp
    tipo: http
    reintentos: {max_intentos: 3, backoff_segundos: 1}
    con:
      url: https://estadisticas.bcrp.gob.pe/estadisticas/series/api/PD04640PD/json/...

  - id: extraer
    tipo: transformar
    depende_de: [consultar_bcrp]
    con:
      asignar:
        valor: "{{ pasos.consultar_bcrp.salida.json.periods.-1.values.0 }}"

  - id: supera_umbral
    tipo: condicion
    depende_de: [extraer]
    con:
      si: "{{ pasos.extraer.salida.valor }} > {{ entradas.umbral }}"

  - id: alertar
    tipo: log
    depende_de: [supera_umbral]
    solo_si: supera_umbral
    con:
      mensaje: "ALERTA: dolar a S/ {{ pasos.extraer.salida.valor }}"

Y el motor lo corre — contra la API real del BCRP:

$ python -m flujo.cli correr workflows/dolar_alerta.yml --entrada umbral=3.0

== dolar-alerta ==
  [log] ALERTA: dolar a S/ 3.408 (17.Jul.26), sobre el umbral de S/ 3.0

  OK  consultar_bcrp       exitoso
  OK  extraer              exitoso
  OK  supera_umbral        exitoso
  OK  alertar              exitoso

  Run: EXITOSO  (523 ms)

Con el umbral por defecto (3.50) el dólar no lo supera, la condición da falso y alertar se omite — la rama simplemente no se toma.

Qué hay dentro del motor

flowchart LR
    Y[YAML] --> V[Validación<br/>Pydantic]
    V --> D[Orden topológico<br/>+ detección de ciclos]
    D --> E[Ejecución por paso]
    E --> P[Plantillas<br/>resueltas]
    P --> R[Ejecutor del paso<br/>+ reintentos/backoff]
    R --> C[Contexto<br/>compartido]
    C --> E
    E --> J[Diario JSONL<br/>del run]
Loading
Pieza Módulo Qué resuelve
DSL validado modelos.py Ids duplicados, dependencias rotas o referencias inexistentes fallan en la carga, no a mitad de la corrida.
DAG motor.py Orden topológico (Kahn) + detección de ciclos antes de ejecutar nada.
Plantillas plantillas.py {{ pasos.x.salida.monto }} conserva el tipo (número/lista/dict) cuando es la plantilla completa; interpola como texto cuando está incrustada. Soporta índices de lista, incluidos negativos.
Condiciones seguras condiciones.py Comparaciones sin eval() — el YAML puede venir de un usuario y no debe poder ejecutar código.
Reintentos motor.py Backoff exponencial por paso (base × 2^(intento−1)).
Semántica de fallos motor.py en_fallo: abortar detiene el run; continuar lo registra y sigue (run parcial). Un paso con dependencias no exitosas se omite: nunca corre sobre datos a medias.
Plugins registro.py Un tipo de paso nuevo es una función dict -> dict y una línea de registro.
Observabilidad diario.py Diario JSONL por run: qué corrió, cuándo, con cuántos intentos, qué falló.

Extenderlo es una línea

from flujo.registro import registrar

def notificar_slack(params: dict) -> dict:
    ...  # tu integración
    return {"enviado": True}

registrar("slack", notificar_slack)

Desde ese momento cualquier YAML puede usar tipo: slack. Los tests usan exactamente este mecanismo para probar los reintentos con un paso que falla dos veces y funciona a la tercera.

Tipos de paso incluidos

Tipo Hace Salida
http Llamada HTTP con timeout {status, json|texto}
transformar Mapea valores del contexto el dict asignado
condicion Evalúa una expresión {resultado: bool}
ia Pregunta a un LLM (Groq) {respuesta}
log Imprime un mensaje {mensaje}

El paso ia requiere GROQ_API_KEY (gratis en console.groq.com); sin la clave falla con un mensaje claro y el workflow decide qué hacer con eso (reintentos / en_fallo) — ver workflows/resumen_con_ia.yml.

Uso

# CLI
python -m flujo.cli correr workflows/triaje_leads.yml            # demo local, sin red
python -m flujo.cli correr workflows/triaje_leads.yml --entrada puntaje=45
python -m flujo.cli validar workflows/*.yml                      # valida sin ejecutar

# API (documentación interactiva en /docs)
uvicorn flujo.api.main:app --reload

# Docker
docker build -t flujo . && docker run -p 8000:8000 flujo

Correrlo tú mismo

python -m venv .venv
.venv\Scripts\activate          # Windows  (Linux/Mac: source .venv/bin/activate)
pip install -r requirements.txt

pytest tests/ -v                # 32 tests
python -m flujo.cli correr workflows/triaje_leads.yml

Qué demuestra este proyecto

  • Código de infraestructura, no solo de aplicación: grafos (Kahn + ciclos), un mini-lenguaje de plantillas, un evaluador seguro, arquitectura de plugins.
  • Semántica de errores pensada: abortar vs continuar vs omitir, reintentos con backoff, fallos que se propagan de forma predecible.
  • Seguridad por diseño: condiciones sin eval(), validación completa antes de ejecutar.
  • 32 tests que cubren el motor de punta a punta, incluida la extensibilidad.
  • Empaque completo: CLI, API REST (OpenAPI), Docker, CI.

Stack

Python · Pydantic · PyYAML · FastAPI · Groq (LLM) · pytest · Docker · GitHub Actions

Licencia

MIT — úsalo como quieras.

About

Motor de workflows de automatizacion definidos en YAML: DAG, reintentos con backoff, plantillas, condiciones sin eval y pasos de IA. Un mini-n8n desde cero.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages