CreaRack-SL

Consumer monitoring · realtime WebSocket — canal de tiempo real del Observatory

Descripción

Consumer de Django Channels que gestiona el canal WebSocket de tiempo real del Observatory: suscripción a metricas/alertas de MonitoringTarget, broadcast de actualizaciones en vivo, aislamiento multi-tenant por organización.

Archivo: monitoring/consumers.py (285 LOC)
Refactorizado en: commit f263a03 (sesión s108, auditoría sa6)

Conexión (connect)

  1. Lee self.scope['user'] → rechaza anónimos con await self.close().
  2. Obtiene la org del usuario mediante await self.get_user_org_id().
  3. Si org es None, rechaza la conexión (log + close). Un usuario sin organización no puede escuchar el realtime.
  4. Acepta la conexión; se mantienen en memoria: self.org_id, self.subscribed_targets (set de PKs).

Seguridad de conexión:

  • ✅ Solo usuarios autenticados.
  • ✅ Solo usuarios con org válida.
  • ✅ Log de rechazos para auditoría.

Recepción de mensajes (receive)

Procesa mensajes JSON del cliente:

{
  "type": "subscribe" | "unsubscribe" | "ping",
  "target_ids": [1, 2, 3]  // para subscribe/unsubscribe (client-controlled)
}

Validación robusta

Antes (auditoría s6 identificó):

  • json.loads(text_data) sin try/except → JSONDecodeError no atrapado.
  • Payload malformado (no dict, array, etc.) → proceso sin validar tipo.
  • target_ids del cliente pasa directo a query sin coerce → ValueError/DataError si incluye strings no-numéricos.

Después (refactorizado):

try:
    data = json.loads(text_data)
except json.JSONDecodeError:
    await self.send({"type": "error", "message": "Invalid JSON"})
    return

if not isinstance(data, dict):
    await self.send({"type": "error", "message": "Expected a JSON object"})
    return

Devuelve error civilizado, cierra stream normalmente (no crash).

_coerce_target_ids(raw) → List[int]

Nuevo helper: sanitiza target_ids del cliente a enteros limpios.

@staticmethod
def _coerce_target_ids(raw):
    """Coerce to list of ints. Drops non-numeric; avoids DB ValueError."""
    if raw is None:
        return []
    if not isinstance(raw, (list, tuple)):
        raw = [raw]
    ids = []
    for item in raw:
        try:
            ids.append(int(item))
        except (TypeError, ValueError):
            continue  # Drop junk, don't crash
    return ids

Propósito:

  • Client malicioso o buggy no puede inyectar strings en la query (batch_check_access).
  • Convierte target_ids=[1, "abc", 2.5] → [1, 2] (silenciosamente).
  • Evita DataError: invalid input syntax for integer sin try/catch en batch_check_access.

Procesamiento por tipo

type: "subscribe"

  1. Coerce data.get("target_ids") via _coerce_target_ids().
  2. Llama a batch_check_access(target_ids) → query que filtra por organization=self.org_id → devuelve solo PKs que pertenecen a la org del usuario.
  3. Para cada PK accesible, suscribe al grupo target_<id> (p.ej., target_42).
  4. Responde con lista completa de objetivos suscritos.

Aislamiento multi-tenant:

  • ✅ batch_check_access filtra por org → cliente NO puede suscribirse a un target de otra org.
  • ✅ Grupo derivado de PK (no de input) → grupo global imposible.
  • ✅ Verificado en auditoría s6 — no había fuga.

type: "unsubscribe"

  1. Coerce target_ids.
  2. Para cada PK que ya está suscrito (self.subscribed_targets), desuscribe del grupo.

type: "ping"

Responde {"type": "pong"} — heartbeat para detectar conexiones muertas.

type: <unknown>

Devuelve error: {"type": "error", "message": f"Unknown message type: {msg_type!r}"}.

Antes: un tipo desconocido caía en silencio (sin respuesta). Ahora es explícito.

Broadcast de eventos

Funciones de clase que envían eventos a canales de suscripción:

broadcast_metric_update(org_id, target_id, metric_data)

Actualización de métrica (CPU, memoria, etc.) → grupo target_<id>.

broadcast_alert(org_id, target_id, alert_data)

Alerta crítica → grupo target_<id>.

broadcast_job_progress(org_id, job_id, progress)

Progreso de operación batch → grupo job_<id>.

broadcast_agent_status(org_id, agent_id, status)

Cambio de estado del agente → grupo agent_<id>.

broadcast_insight(org_id, insight_data)

Nueva perspectiva de IA → grupo org_<org_id> (globla por org).

Patrón de nombres:

  • target_<id> — métrica/alerta del objetivo específico.
  • job_<id> — progreso de operación específica.
  • agent_<id> — estado del agente específico.
  • org_<id> — eventos globales de la org.

Todos los broadcasts reciben org_id explícitamente (verifican que el channel_name pertenece a esa org antes de enviar, o confían en que el listener valida su suscripción). Aislamiento verificado en auditoría s6.

Desconexión (disconnect)

Limpia self.subscribed_targets: desuscribe de todos los grupos.

async def disconnect(self):
    for tid in list(self.subscribed_targets):
        group_name = f"target_{tid}"
        await self.channel_layer.group_discard(group_name, self.channel_name)

Modelo de seguridad

AspectoGarantía
AutenticaciónSolo usuarios loggeados (self.scope['user'].is_authenticated)
OrganizaciónUsuario debe tener org válida; rechazo si falta (self.org_id is None)
Aislamiento realtimebatch_check_access(org=self.org_id) filtra targets antes de group_add
Sanitización input_coerce_target_ids() convierte junk a ints; nada malformado llega a DB
Error handlingJSON/tipo errores devuelven msg civilizado, no crash
LoggingRechazos de conexión (anónimo, sin org) se registran

Auditoría s6 (sesión s108)

Hallazgos sobre realtime WebSocket:

  • ❌ (no real) Fuga cross-tenant en group_add — verificado en auditoría: batch_check_access filtra por org correctamente.
  • ❌ (no real) Conexión anónima capaz de recibir broadcasts — verificado: connect() rechaza anónimos.
  • ✅ (arreglado) Payload malformado → crash del consumer — ahora maneja gracefully.
  • ✅ (arreglado) target_ids junk sin sanitizar → ValueError en query — ahora coercionado.
  • ✅ (arreglado) Tipo de mensaje desconocido cae en silencio — ahora devuelve error.

Decisión de arquitectura: se mantuvo el aislamiento por org ya existente (fue correcto). Se reforzó la robustez de validación + logging.

Carve-outs (backlog Etapa 3)

  • Mover broadcast a trabajo Huey (async task) en lugar de síncrono en request.
  • Naming defensivo en profundidad: grupo WebSocket target_<org_id>_<id> (además de RLS + permisos) — redundancia de defensa.

Véase también

  • [[feature—monitoring—auditoria-suprema-sa6-observatory-realtime]]
  • [[concept—saas—multi-tenancy]]
  • [[concept—security—authentication]]
  • [[entity—monitoring—model—monitoring-target]]