⚡ Automatizaciones · Temporal Workflows

Orquestación de Workflows: NutriLoop Onboarding

Diseño de arquitectura Temporal para automatizar el onboarding B2B de extremo a extremo — saga, estado persistente y aprobación humana.

Cliente: NutriLoop SaaS
Stack: Temporal + Python 3.11
Patrones: Saga, Entity, Fan-Out, Async-Callback
SLA: <3 min onboarding end-to-end
Impacto estimado
4h → 3min
Tiempo de onboarding
De proceso manual a automatizado
99.7%
Tasa de éxito esperada
Con retries automáticos Temporal
5
Sistemas integrados
Stripe · HubSpot · Core · Slack · Email
0
Intervenciones manuales
Excepto aprobación >1.000 asientos
Arquitectura del workflow
NutriLoopOnboardingWorkflow SAGA PATTERN
🎯
TRIGGER
CRM webhook
Deal cerrado en HubSpot → start_workflow()
🔀
¿>1.000 asientos?
Compliance gate
workflow.wait_condition() hasta señal
💳
Stripe
create_subscription()
Retry 3×, backoff exp.
🏗️
Provisionar
create_workspace()
Idempotente (upsert por org_id)
📊
HubSpot
sync_crm()
No retryable: DuplicateError
💬
Slack + Email
Fan-out async
Actividades paralelas
COMPLETED
~2.8 min
Estado final persistido
Patrón Saga — pasos y compensaciones

Cada paso registra su compensación antes de ejecutarse. Si falla el paso N, se ejecutan las compensaciones N-1 … 1 en orden inverso (LIFO).

Paso
Acción (Activity)
Compensación (Rollback)
1
Reservar plan
stripe.create_subscription()
stripe.cancel_subscription()
2
Cobrar primera cuota
stripe.charge_setup_fee()
stripe.refund_charge()
3
Provisionar workspace
core.create_workspace()
core.delete_workspace()
4
Asignar roles y asientos
core.assign_seats()
core.revoke_seats()
5
Sync CRM
hubspot.create_company()
hubspot.archive_company()
6
Notificar cliente
slack.invite_to_channel() + email.send_welcome()
— (fire-and-forget, no reversible)
Código de referencia — Workflow
🐍 Python · Temporal SDK
nutriloop/workflows/onboarding.py
from temporalio import workflow
from temporalio.common import RetryPolicy
from dataclasses import dataclass
from datetime import timedelta

# ── Tipos de entrada/salida ──────────────────────────────────────────
@dataclass
class OnboardingInput:
    org_id: str
    plan: str
    seats: int
    admin_email: str
    hubspot_deal_id: str

@workflow.defn
class NutriLoopOnboardingWorkflow:
    """
    Orquesta el onboarding B2B completo.
    Implementa Saga con compensaciones LIFO para rollback atómico.
    """

    @workflow.run
    async def run(self, inp: OnboardingInput) -> dict:
        compensations: list = []

        # ── Compliance gate: aprobación humana si > 1.000 asientos ─────────
        if inp.seats > 1000:
            # Bloquea hasta recibir la señal 'compliance_approved'
            await workflow.wait_condition(
                lambda: self._compliance_ok,
                timeout=timedelta(hours=48)          # escalado automático si no hay respuesta
            )

        try:
            # ── Paso 1: Stripe ──────────────────────────────────────────────────
            sub = await workflow.execute_activity(
                create_stripe_subscription,
                args=[inp.org_id, inp.plan, inp.seats],
                retry_policy=RetryPolicy(
                    maximum_attempts=3,
                    backoff_coefficient=2.0,
                    non_retryable_error_types=["CardDeclinedError"]
                ),
                start_to_close_timeout=timedelta(seconds=30)
            )
            compensations.append((cancel_stripe_subscription, sub["subscription_id"]))

            # ── Paso 2: Provisionar workspace ───────────────────────────────────
            ws = await workflow.execute_activity(
                create_workspace,
                args=[inp.org_id, inp.seats, sub["subscription_id"]],
                retry_policy=RetryPolicy(maximum_attempts=5),
                start_to_close_timeout=timedelta(minutes=2)
            )
            compensations.append((delete_workspace, ws["workspace_id"]))

            # ── Paso 3: Fan-out — HubSpot + Slack + Email en paralelo ───────────
            slack_h, email_h, crm_h = await asyncio.gather(
                workflow.execute_activity(invite_to_slack, args=[inp.admin_email, ws["channel"]]),
                workflow.execute_activity(send_welcome_email, args=[inp.admin_email]),
                workflow.execute_activity(sync_hubspot, args=[inp.hubspot_deal_id, ws]),
            )

            return {"status": "success", "workspace_id": ws["workspace_id"]}

        except Exception as err:
            # ── Rollback LIFO ───────────────────────────────────────────────────
            for fn, arg in reversed(compensations):
                await workflow.execute_activity(fn, args=[arg])
            raise

    @workflow.signal
    def compliance_approved(self, approved_by: str):
        self._compliance_ok = True

    @workflow.query
    def current_state(self) -> str:
        return self._state
Código de referencia — Actividades
🐍 Python · Activities
nutriloop/activities/stripe_activities.py
from temporalio import activity
import stripe

@activity.defn
async def create_stripe_subscription(org_id: str, plan: str, seats: int) -> dict:
    """
    Actividad idempotente: usa org_id como idempotency_key.
    Si Stripe ya procesó la request, devuelve el objeto existente.
    """
    activity.heartbeat("Creando suscripción Stripe...")         # heartbeat para detección de stall

    try:
        sub = stripe.Subscription.create(
            customer=org_id,
            items=[{"price": price_id_for(plan, seats)}],
            idempotency_key=f"onboard-{org_id}"           # garantía de idempotencia
        )
        return {"subscription_id": sub.id, "status": sub.status}

    except stripe.error.CardError as e:
        # Error NO retryable — Temporal no volverá a intentarlo
        raise CardDeclinedError(str(e)) from e

    except stripe.error.APIConnectionError as e:
        # Error transitorio — Temporal reintentará con backoff exponencial
        raise                                               # re-raise para retry automático

@activity.defn
async def cancel_stripe_subscription(subscription_id: str) -> None:
    """Compensación idempotente: cancelar está bien llamar N veces."""
    try:
        stripe.Subscription.delete(subscription_id)
    except stripe.error.InvalidRequestError:
        pass                                                # ya cancelada → ok
Reglas de determinismo — Workflows
✅ Permitido en Workflow
Usar workflow.now() para obtener la hora actual
Usar workflow.random() para valores aleatorios
Funciones puras, cálculos locales sin I/O
Llamar actividades con execute_activity()
Señales y queries para comunicación externa
Esperar condiciones con wait_condition()
Lanzar child workflows con execute_child_workflow()
❌ Prohibido en Workflow
datetime.now() — no determinista entre replays
random.random() — produce valores distintos en cada replay
Llamadas directas a APIs, DB o red
Threading, locks o primitivas de sincronización
Variables globales o estado estático mutable
time.sleep() — usar asyncio.sleep() de Temporal
Librerías no deterministas (uuid4() sin seed)
Patrones aplicados
Patrón Dónde se aplica Beneficio Complejidad
Saga + Compensaciones Pasos 1–5 de onboarding Transacciones distribuidas sin 2-phase commit MEDIA
Async Callback (Human-in-the-loop) Compliance gate >1.000 asientos Bloqueo sin consumir CPU; timeout automático a 48h BAJA
Fan-Out / Fan-In Slack + Email + HubSpot en paralelo Reducción de latencia de 3 pasos secuenciales a 1 BAJA
Entity Workflow WorkspaceLifecycleWorkflow (post-onboarding) Un workflow por organización; recibe señales de upgrade/downgrade MEDIA
Retry con backoff exponencial Todas las actividades de red Recuperación automática de errores transitorios BÁSICA
Idempotency keys Stripe, Core provisioning Seguridad ante retries: cobros y workspaces únicos por org_id BAJA
Heartbeat en actividades largas create_workspace (hasta 2 min) Detección de stall; progreso visible en Temporal UI BAJA
Child Workflows (escala) BulkOnboardingWorkflow (>100 empresas) 1 parent × N child; cada child independiente y bounded ALTA
📊 Métricas de monitorización
workflow_execution_duration
Duración total: target <3 min. Alerta >10 min.
activity_failure_rate
Por actividad (Stripe, Core, HubSpot). Alerta >2%.
activity_retry_count
Reintentos por tipo. Spike indica degradación de servicio externo.
pending_workflow_count
Flujos esperando human-approval. SLA 48h para compliance.
saga_compensation_runs
Rollbacks ejecutados. Indica fallos irrecuperables.
🔧 Configuración de retry policies
Python
# Actividades críticas (Stripe, Core)
STRICT_RETRY = RetryPolicy(
    initial_interval=timedelta(seconds=1),
    backoff_coefficient=2.0,
    maximum_interval=timedelta(seconds=30),
    maximum_attempts=3,
    non_retryable_error_types=[
        "CardDeclinedError",
        "ValidationError"
    ]
)

# Actividades de notificación (Slack, Email)
SOFT_RETRY = RetryPolicy(
    maximum_attempts=5,
    backoff_coefficient=1.5
)

Los errores de negocio (tarjeta rechazada, email inválido) se marcan como non_retryable para evitar loops infinitos innecesarios. Los errores de red/timeout se reintentan automáticamente.


Framework de decisión: ¿Workflow o Activity?
🔀 WORKFLOW — Orquestación
Lógica de negocio, decisiones, coordinación de pasos.
✓ ¿Hay más de 2 pasos secuenciales o paralelos?
✓ ¿Debe ejecutarse en minutos, horas o días?
✓ ¿Necesita sobrevivir a reinicios de servidor?
✓ ¿Tiene lógica de rollback (saga)?
✓ ¿Espera aprobación humana o eventos externos?
⚡ ACTIVITY — Ejecución
Toda interacción con sistemas externos, I/O, red, DB.
✓ ¿Toca una API externa (Stripe, HubSpot, Slack)?
✓ ¿Lee o escribe en base de datos?
✓ ¿Puede tener efectos secundarios (enviar email)?
✓ ¿Tarda segundos o minutos de forma variable?
✓ ¿Necesita ser idempotente y reiniciable?