CultivaFlow · Redis Patterns

Implementación completa: Caché · Sesiones · Rate Limiting · Pub/Sub · Bloqueo Distribuido · Leaderboards

Redis 7.2 Python 3.12 FastAPI · 3 workers
Latencia caché
<10ms
▼ vs 800ms PostgreSQL
Usuarios activos/día
8.000
3 workers FastAPI
RAM Redis
2 GB
⚠ maxmemory configurado
TTL sesiones
8 h
Expiración deslizante

⚡ Cache-Aside

Catálogos de contactos (50K entries). Lectura desde Redis, fallback a PG, escritura lazy.

Perf
97%

🔑 Session Store

SSO propio. Sliding expiration 8h. Logout global via user_sessions set.

Perf
99%

🛑 Sliding Window RL

Rate limiting por API key. 3 planes: Free 100/min · Pro 1K/min · Agency 10K/min.

Perf
94%

🔒 Distributed Lock

Previene doble envío de campañas. Redlock pattern con Lua script atómico.

Perf
100%

Setup & Conexión

cultivaflow/redis_client.py
cultivaflow/redis_client.py Python
"""redis_client.py — CultivaFlow: connection pool para 3 workers FastAPI."""
import os
import redis
from redis.backoff import ExponentialBackoff
from redis.retry  import Retry

# ─── Connection pool compartido entre workers ─────────────────────────────
REDIS_URL = os.getenv("REDIS_URL", "redis://:secret@redis-cloud:6379/0")

pool = redis.ConnectionPool.from_url(
    REDIS_URL,
    max_connections=30,       # 10 por worker × 3 workers
    decode_responses=True,
    socket_timeout=0.5,       # 500ms — falla rápido, fallback a PG
    socket_connect_timeout=1.0,
)

retry = Retry(ExponentialBackoff(), retries=3)
r = redis.Redis(connection_pool=pool, retry=retry)

# Health check al arrancar
def ping_redis() -> bool:
    try:
        return r.ping()
    except redis.RedisError:
        return False

Caché de Catálogos (Cache-Aside)

800ms → <10ms
cultivaflow/cache/contacts.py Python
"""cache/contacts.py — Cache-aside para catálogos de contactos (hasta 50K)."""
import json, time
from .redis_client import r
from .db import db

CONTACTS_TTL = 3600     # 1h — catálogos cambian poco
PRODUCT_TTL  = 7200     # 2h — datos de campañas

def get_agency_contacts(agency_id: str, page: int = 1) -> dict:
    """Fetch contactos paginados. Cache-aside con fallback a PG.
    Key: cultivaflow:contacts:{agency_id}:page:{page}"""

    cache_key = f"cultivaflow:contacts:{agency_id}:page:{page}"
    cached    = r.get(cache_key)

    if cached:
        return json.loads(cached)  # ⚡ Cache HIT — <1ms

    # Cache MISS — consulta PG (~800ms)
    contacts = db.query(
        "SELECT id,email,name,tags,score FROM contacts"
        " WHERE agency_id=%s ORDER BY score DESC LIMIT 500 OFFSET %s",
        agency_id, (page-1)*500
    )
    payload = {"data": contacts, "page": page, "cached_at": time.time()}
    r.setex(cache_key, CONTACTS_TTL, json.dumps(payload))
    return payload

def invalidate_agency_contacts(agency_id: str):
    """Invalidar todas las páginas de una agencia (SCAN, no KEYS *)."""
    cursor, keys = 0, []
    while True:
        cursor, batch = r.scan(
            cursor, match=f"cultivaflow:contacts:{agency_id}:*", count=100
        )
        keys.extend(batch)
        if cursor == 0: break
    if keys: r.delete(*keys)
🔑

Sesiones con Expiración Deslizante

SSO · 8h · logout global
cultivaflow/auth/sessions.py Python
"""auth/sessions.py — SSO con sliding expiration y logout-all."""
import secrets, json, time
from .redis_client import r

SESSION_TTL = 28800  # 8 horas en segundos

def create_session(user_id: str, agency_id: str, metadata: dict = None) -> str:
    """Crea sesión y registra token en el set del usuario para logout global.

    Returns:
        token: str  →  enviar como cookie HttpOnly
    """
    token    = secrets.token_urlsafe(32)
    session  = {
        "user_id":   user_id,
        "agency_id": agency_id,
        "created_at": time.time(),
        **(metadata or {}),
    }
    pipe = r.pipeline()
    pipe.setex(f"cultivaflow:session:{token}", SESSION_TTL, json.dumps(session))
    pipe.sadd(f"cultivaflow:user_sessions:{user_id}", token)
    pipe.expire(f"cultivaflow:user_sessions:{user_id}", SESSION_TTL)
    pipe.execute()
    return token

def get_session(token: str) -> dict | None:
    """Valida sesión y extiende TTL (sliding expiration)."""
    data = r.get(f"cultivaflow:session:{token}")
    if not data:
        return None
    r.expire(f"cultivaflow:session:{token}", SESSION_TTL)  # ← Sliding
    return json.loads(data)

def destroy_all_sessions(user_id: str):
    """Logout en todos los dispositivos (cambio de contraseña / breach)."""
    tokens = r.smembers(f"cultivaflow:user_sessions:{user_id}")
    pipe   = r.pipeline()
    for t in tokens:
        pipe.delete(f"cultivaflow:session:{t}")
    pipe.delete(f"cultivaflow:user_sessions:{user_id}")
    pipe.execute()
🛑

Rate Limiting · Ventana Deslizante

3 planes · Sorted Sets
Plan Límite Ventana Key Redis Coste TTL
Free 100 req 60 s cultivaflow:rl:free:{api_key} 60 s
Pro 1.000 req 60 s cultivaflow:rl:pro:{api_key} 60 s
Agency 10.000 req 60 s cultivaflow:rl:agency:{api_key} 60 s
cultivaflow/middleware/rate_limiter.py Python
"""middleware/rate_limiter.py — Sliding window con Sorted Sets."""
import time
from .redis_client import r

PLANS = {
    "free":   (100,   60),   # (limite, ventana_s)
    "pro":    (1000,  60),
    "agency": (10000, 60),
}

def check_rate_limit(api_key: str, plan: str) -> dict:
    """Comprueba y registra una petición. Devuelve estado y headers.

    Returns:
        {"allowed": bool, "remaining": int, "reset_in": float}
    """
    limit, window = PLANS.get(plan, PLANS["free"])
    now  = time.time()
    wstart = now - window
    rk   = f"cultivaflow:rl:{plan}:{api_key}"

    pipe = r.pipeline()
    pipe.zremrangebyscore(rk, 0, wstart)      # purga expirados
    pipe.zadd(rk, {f"{now}": now})              # añade petición actual
    pipe.zcard(rk)                               # cuenta en ventana
    pipe.expire(rk, window)                      # auto-cleanup
    _, _, count, _ = pipe.execute()

    return {
        "allowed":   count <= limit,
        "remaining": max(0, limit - count),
        "reset_in":  window,
        "limit":     limit,
    }
📡

Pub/Sub · Notificaciones en Tiempo Real

SSE → dashboard
cultivaflow/realtime/pubsub.py Python
"""realtime/pubsub.py — Pub/Sub para eventos de campaña → SSE frontend."""
import json, threading
from .redis_client import r

# Canales CultivaFlow
# cultivaflow:events:campaign:{agency_id}   →  campaign.sent, campaign.failed
# cultivaflow:events:lead:{agency_id}        →  lead.converted, lead.scored
# cultivaflow:events:automation:{agency_id}  →  automation.triggered

def emit(agency_id: str, entity: str, event_type: str, payload: dict):
    """Publica un evento en el canal de la agencia.

    Args:
        entity: "campaign" | "lead" | "automation"
        event_type: "sent" | "failed" | "converted" | "triggered"
    """
    channel = f"cultivaflow:events:{entity}:{agency_id}"
    message = {"type": event_type, "payload": payload, "ts": __import__("time").time()}
    r.publish(channel, json.dumps(message))

def subscribe_agency(agency_id: str, on_event):
    """Suscribe a todos los eventos de una agencia (pattern: cultivaflow:events:*:{agency_id}).

    Args:
        on_event: callback(channel, data_dict)
    """
    ps      = r.pubsub()
    pattern = f"cultivaflow:events:*:{agency_id}"
    ps.psubscribe(pattern)

    def _listener():
        for msg in ps.listen():
            if msg["type"] == "pmessage":
                on_event(msg["channel"], json.loads(msg["data"]))

    threading.Thread(target=_listener, daemon=True).start()
    return ps
🔒

Bloqueo Distribuido · Redlock Pattern

Anti-doble-envío campañas
cultivaflow/workers/lock.py Python
"""workers/lock.py — Distributed lock para evitar doble envío de campañas."""
import secrets
from contextlib import contextmanager
from .redis_client import r

# Lua script: check + delete atómico — solo el dueño puede liberar
_RELEASE_LUA = """
if redis.call('GET', KEYS[1]) == ARGV[1] then
    return redis.call('DEL', KEYS[1])
end
return 0
"""

def acquire(resource: str, ttl: int = 30) -> str | None:
    """Adquiere lock. Devuelve token o None si ya está bloqueado.

    Args:
        resource: ej. "campaign:send:camp_abc123"
        ttl:      segundos antes de expirar automáticamente (evita deadlock)
    """
    token = secrets.token_urlsafe(16)
    key   = f"cultivaflow:lock:{resource}"
    return token if r.set(key, token, nx=True, ex=ttl) else None

def release(resource: str, token: str) -> bool:
    """Libera lock sólo si somos el dueño (Lua atómico)."""
    return r.eval(_RELEASE_LUA, 1, f"cultivaflow:lock:{resource}", token) == 1

@contextmanager
def campaign_lock(campaign_id: str):
    """Context manager para envío de campañas — garantiza ejecución única.

    Usage:
        with campaign_lock("camp_abc123") as acquired:
            if not acquired: raise AlreadySendingError()
            send_campaign(...)
    """
    resource = f"campaign:send:{campaign_id}"
    token    = acquire(resource, ttl=300)  # 5min max por envío
    try:
        yield token is not None
    finally:
        if token: release(resource, token)
🏆

Leaderboard de Campañas en Tiempo Real

Sorted Sets · ZREVRANGE
cultivaflow/analytics/leaderboard.py Python
"""analytics/leaderboard.py — Top 10 campañas por open rate en tiempo real."""
from .redis_client import r

def record_open(agency_id: str, campaign_id: str, open_rate: float):
    """Actualiza el score de una campaña (open rate × 100 para evitar floats)."""
    key   = f"cultivaflow:leaderboard:{agency_id}:open_rate"
    score = round(open_rate * 10000)  # 42.75% → 42750 (entero exacto)
    r.zadd(key, {campaign_id: score})
    r.expire(key, 86400)              # reseteamos daily leaderboard

def top_campaigns(agency_id: str, n: int = 10) -> list:
    """Devuelve top-N campañas con scores (open rate real)."""
    key     = f"cultivaflow:leaderboard:{agency_id}:open_rate"
    entries = r.zrevrange(key, 0, n-1, withscores=True)
    return [
        {"campaign_id": cid, "open_rate": score / 10000}
        for cid, score in entries
    ]

def campaign_rank(agency_id: str, campaign_id: str) -> int | None:
    """Posición de una campaña en el leaderboard (0-indexed desde el top)."""
    key  = f"cultivaflow:leaderboard:{agency_id}:open_rate"
    return r.zrevrank(key, campaign_id)
🗂

Key Naming Schema

cultivaflow:{entidad}:{id}:{campo}

Todas las claves siguen: cultivaflow:{entidad}:{id}[:{subcampo}]

cultivaflow:contacts:{agency_id}:page:{n} STRING TTL: 3.600s (1h)
cultivaflow:session:{token} STRING TTL: 28.800s (8h, deslizante)
cultivaflow:user_sessions:{user_id} SET TTL: 28.800s
cultivaflow:rl:{plan}:{api_key} ZSET TTL: 60s (auto)
cultivaflow:lock:{resource} STRING TTL: 30–300s (configurable)
cultivaflow:leaderboard:{agency_id}:open_rate ZSET TTL: 86.400s (24h)

Guidelines de Producción

🚀

Configuración de Producción

redis.conf · Redis Cloud
redis.conf (extracto recomendado) Config
# Memoria
maxmemory         2gb
maxmemory-policy  allkeys-lru   # evicta LRU cuando llega al límite

# Persistencia (no perder sesiones en restart)
appendonly        yes
appendfsync       everysec       # equilibrio rendimiento/durabilidad

# Seguridad
requirepass       your-strong-password
bind              127.0.0.1      # no exponer al exterior
protected-mode    yes

# Timeouts y keepalive
timeout           300
tcp-keepalive     60

# Monitoring
slowlog-log-slower-than 10000   # log operaciones >10ms
latency-monitor-threshold 50    # alertas >50ms