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)
- Lee
self.scope['user']→ rechaza anónimos conawait self.close(). - Obtiene la org del usuario mediante
await self.get_user_org_id(). - Si org es
None, rechaza la conexión (log + close). Un usuario sin organización no puede escuchar el realtime. - 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 →JSONDecodeErrorno atrapado.- Payload malformado (no dict, array, etc.) → proceso sin validar tipo.
target_idsdel cliente pasa directo a query sin coerce →ValueError/DataErrorsi 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 integersin try/catch enbatch_check_access.
Procesamiento por tipo
type: "subscribe"
- Coerce
data.get("target_ids")via_coerce_target_ids(). - Llama a
batch_check_access(target_ids)→ query que filtra pororganization=self.org_id→ devuelve solo PKs que pertenecen a la org del usuario. - Para cada PK accesible, suscribe al grupo
target_<id>(p.ej.,target_42). - Responde con lista completa de objetivos suscritos.
Aislamiento multi-tenant:
- ✅
batch_check_accessfiltra 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"
- Coerce
target_ids. - 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
| Aspecto | Garantía |
|---|---|
| Autenticación | Solo usuarios loggeados (self.scope['user'].is_authenticated) |
| Organización | Usuario debe tener org válida; rechazo si falta (self.org_id is None) |
| Aislamiento realtime | batch_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 handling | JSON/tipo errores devuelven msg civilizado, no crash |
| Logging | Rechazos 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_accessfiltra 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_idsjunk sin sanitizar →ValueErroren 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]]
Referenciado desde
- Auditoría #261 ciclo 4 — el Observatory deja de quedarse dormida
- Auditoría Suprema monitoring sa6 — páginas del Observatory + tiempo real (s108)
- Deuda #274 Tanda A — alertas en tiempo real y datos honestos en Wireless (v1.85.9)
- Las alertas del Observatory se conectan a la vía viva del Agente (task #229)
- Las gráficas del Observatory salían vacías a ratos: el suavizado se leía fuera del hilo del ORM (task #297)