Current section
Files
Jump to
Current section
Files
lib/dashboard/pubsub_bridge.ex
defmodule ExESDBDashboard.PubSubBridge do
@moduledoc """
Bridge that subscribes to ex-esdb-gater PubSub events and translates them
to dashboard-specific events on ExESDBDashboard.PubSub.
This encapsulates the integration between the gater system and the dashboard,
allowing the dashboard package to be completely self-contained while still
receiving real-time updates from the cluster.
"""
use GenServer
require Logger
@gater_pubsub_instances [
:ex_esdb_events, # Core event data (application-specific)
:ex_esdb_system, # General system events and configuration
:ex_esdb_health, # Health monitoring and cluster status
:ex_esdb_metrics, # Performance metrics and analytics
:ex_esdb_alerts, # Critical alerts and notifications
:ex_esdb_lifecycle, # Process lifecycle events
:ex_esdb_security, # Security events and authentication
:ex_esdb_audit, # Audit trail and compliance tracking
:ex_esdb_diagnostics, # Diagnostic information and debugging
:ex_esdb_logging # Log aggregation and centralized logging
]
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@impl true
def init(_opts) do
Logger.info("[Dashboard.PubSubBridge] Starting PubSub bridge for dashboard events")
# Subscribe to relevant gater PubSub instances
subscribe_to_gater_events()
{:ok, %{subscriptions: [], last_activity: DateTime.utc_now()}}
end
# Handle health events from gater
@impl true
def handle_info({:health_update, health_data}, state) do
# Translate to dashboard format
dashboard_health = translate_health_data(health_data)
# Broadcast to dashboard PubSub
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"cluster_health_updates",
{:cluster_health_update, dashboard_health}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed health update: #{inspect(dashboard_health)}")
{:noreply, update_activity(state)}
end
# Handle system events from gater
def handle_info({:system_event, event_data}, state) do
# Broadcast to dashboard PubSub
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"system_events",
{:system_event, event_data}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed system event: #{inspect(event_data)}")
{:noreply, update_activity(state)}
end
# Handle lifecycle events from gater
def handle_info({:lifecycle_event, event_data}, state) do
# Broadcast to dashboard PubSub
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"cluster_lifecycle_events",
{:cluster_lifecycle, event_data}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed lifecycle event: #{inspect(event_data)}")
{:noreply, update_activity(state)}
end
# Handle metrics events from gater
def handle_info({:metrics_update, metrics_data}, state) do
# Could translate to performance data or system events
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"system_events",
{:system_event, %{type: "metrics_update", data: metrics_data}}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed metrics update")
{:noreply, update_activity(state)}
end
# Handle alerts from gater
def handle_info({:alert, alert_data}, state) do
# Broadcast as system event with enhanced information
enhanced_alert = %{
type: "alert",
message: alert_data.message,
severity: alert_data.severity,
timestamp: Map.get(alert_data, :timestamp, DateTime.utc_now()),
node: Map.get(alert_data, :node, :unknown),
category: Map.get(alert_data, :category, "system")
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"system_events",
{:system_event, enhanced_alert}
)
# Also broadcast to dedicated alerts channel
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_alerts",
{:operational_alert, enhanced_alert}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed alert: #{inspect(enhanced_alert)}")
{:noreply, update_activity(state)}
end
# Handle connection status changes
def handle_info({:connection_status_change, status}, state) do
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"connection_status",
{:connection_status, status}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed connection status: #{inspect(status)}")
{:noreply, update_activity(state)}
end
# Handle security events from gater
def handle_info({:security_event, security_data}, state) do
enhanced_security = %{
type: "security",
event_type: Map.get(security_data, :event_type, "unknown"),
actor: Map.get(security_data, :actor, "system"),
resource: Map.get(security_data, :resource, "unknown"),
timestamp: Map.get(security_data, :timestamp, DateTime.utc_now()),
node: Map.get(security_data, :node, :unknown),
details: Map.get(security_data, :details, %{})
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_security",
{:security_event, enhanced_security}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed security event: #{inspect(enhanced_security)}")
{:noreply, update_activity(state)}
end
# Handle audit events from gater
def handle_info({:audit_event, audit_data}, state) do
enhanced_audit = %{
type: "audit",
actor: Map.get(audit_data, :actor, "system"),
action: Map.get(audit_data, :action, "unknown"),
resource: Map.get(audit_data, :resource, "unknown"),
timestamp: Map.get(audit_data, :timestamp, DateTime.utc_now()),
node: Map.get(audit_data, :node, :unknown),
result: Map.get(audit_data, :result, "success"),
details: Map.get(audit_data, :details, %{})
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_audit",
{:audit_event, enhanced_audit}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed audit event: #{inspect(enhanced_audit)}")
{:noreply, update_activity(state)}
end
# Handle diagnostic events from gater
def handle_info({:diagnostic_event, diagnostic_data}, state) do
enhanced_diagnostic = %{
type: "diagnostic",
component: Map.get(diagnostic_data, :component, "system"),
diagnostic_type: Map.get(diagnostic_data, :diagnostic_type, "trace"),
data: Map.get(diagnostic_data, :data, %{}),
timestamp: Map.get(diagnostic_data, :timestamp, DateTime.utc_now()),
node: Map.get(diagnostic_data, :node, :unknown)
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_diagnostics",
{:diagnostic_event, enhanced_diagnostic}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed diagnostic event: #{inspect(enhanced_diagnostic)}")
{:noreply, update_activity(state)}
end
# Handle logging events from gater
def handle_info({:log_event, log_data}, state) do
enhanced_log = %{
type: "log",
level: Map.get(log_data, :level, "info"),
message: Map.get(log_data, :message, ""),
component: Map.get(log_data, :component, "system"),
timestamp: Map.get(log_data, :timestamp, DateTime.utc_now()),
node: Map.get(log_data, :node, :unknown),
metadata: Map.get(log_data, :metadata, %{})
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_logging",
{:log_event, enhanced_log}
)
# Also add to system events for general dashboard updates
if log_data.level in ["error", "warn"] do
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"system_events",
{:system_event, %{type: "log_#{log_data.level}", message: log_data.message}}
)
end
Logger.debug("[Dashboard.PubSubBridge] Relayed log event: #{inspect(enhanced_log)}")
{:noreply, update_activity(state)}
end
# Handle enhanced metrics with operational patterns
def handle_info({:operational_metrics, metrics_data}, state) do
enhanced_metrics = %{
type: "operational_metrics",
component: Map.get(metrics_data, :component, "system"),
metrics: Map.get(metrics_data, :metrics, %{}),
timestamp: Map.get(metrics_data, :timestamp, DateTime.utc_now()),
node: Map.get(metrics_data, :node, :unknown),
metric_type: Map.get(metrics_data, :metric_type, "performance")
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_metrics",
{:operational_metrics, enhanced_metrics}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed operational metrics: #{inspect(enhanced_metrics)}")
{:noreply, update_activity(state)}
end
# Handle configuration changes
def handle_info({:config_change, config_data}, state) do
enhanced_config = %{
type: "config_change",
component: Map.get(config_data, :component, "system"),
changes: Map.get(config_data, :changes, %{}),
timestamp: Map.get(config_data, :timestamp, DateTime.utc_now()),
node: Map.get(config_data, :node, :unknown),
changed_by: Map.get(config_data, :changed_by, "system")
}
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"operational_config",
{:config_change, enhanced_config}
)
# Also broadcast to system events for general awareness
Phoenix.PubSub.broadcast(
ExESDBDashboard.PubSub,
"system_events",
{:system_event, %{type: "config_change", component: config_data.component}}
)
Logger.debug("[Dashboard.PubSubBridge] Relayed config change: #{inspect(enhanced_config)}")
{:noreply, update_activity(state)}
end
# Generic event handler for unknown events
def handle_info(event, state) do
Logger.debug("[Dashboard.PubSubBridge] Received unknown event: #{inspect(event)}")
{:noreply, state}
end
# Private functions
defp subscribe_to_gater_events do
# Subscribe to relevant topics on each gater PubSub instance
subscriptions = [
# Health monitoring
{:ex_esdb_health, "cluster.health"},
{:ex_esdb_health, "store.health"},
{:ex_esdb_health, "component_health"},
{:ex_esdb_health, "node_health"},
# System events and configuration
{:ex_esdb_system, "cluster.lifecycle"},
{:ex_esdb_system, "store.lifecycle"},
{:ex_esdb_system, "config"},
{:ex_esdb_system, "lifecycle"},
# Alerts and notifications
{:ex_esdb_alerts, "cluster.alerts"},
{:ex_esdb_alerts, "system.alerts"},
{:ex_esdb_alerts, "critical_alerts"},
# Performance metrics
{:ex_esdb_metrics, "performance"},
{:ex_esdb_metrics, "system.metrics"},
{:ex_esdb_metrics, "persistence"},
{:ex_esdb_metrics, "emitter"},
# Process lifecycle
{:ex_esdb_lifecycle, "process.lifecycle"},
{:ex_esdb_lifecycle, "cluster_membership"},
{:ex_esdb_lifecycle, "node_lifecycle"},
# Security events
{:ex_esdb_security, "authentication"},
{:ex_esdb_security, "authorization"},
{:ex_esdb_security, "security_events"},
# Audit trail
{:ex_esdb_audit, "data_change"},
{:ex_esdb_audit, "access_log"},
{:ex_esdb_audit, "audit_trail"},
# Diagnostic information
{:ex_esdb_diagnostics, "debug_trace"},
{:ex_esdb_diagnostics, "diagnostics"},
{:ex_esdb_diagnostics, "system_info"},
# Log aggregation
{:ex_esdb_logging, "system_logs"},
{:ex_esdb_logging, "error_logs"},
{:ex_esdb_logging, "log_aggregation"}
]
successful_subs =
Enum.reduce(subscriptions, [], fn {pubsub_instance, topic}, acc ->
try do
Phoenix.PubSub.subscribe(pubsub_instance, topic)
Logger.debug("[Dashboard.PubSubBridge] Subscribed to #{topic} on #{pubsub_instance}")
[{pubsub_instance, topic} | acc]
rescue
error ->
Logger.debug("[Dashboard.PubSubBridge] Failed to subscribe to #{topic} on #{pubsub_instance}: #{inspect(error)}")
acc
end
end)
Logger.info("[Dashboard.PubSubBridge] Successfully subscribed to #{length(successful_subs)} gater PubSub topics")
successful_subs
end
defp translate_health_data(health_data) when is_map(health_data) do
# Convert gater health format to dashboard format
status = case Map.get(health_data, :status) do
:healthy -> :healthy
:ok -> :healthy
:degraded -> :degraded
:warning -> :degraded
:unhealthy -> :unhealthy
:error -> :unhealthy
:critical -> :unhealthy
_ -> :unknown
end
%{
status: status,
message: Map.get(health_data, :message, "System status update"),
timestamp: Map.get(health_data, :timestamp, DateTime.utc_now()),
source: Map.get(health_data, :source, "gater")
}
end
defp translate_health_data(health_atom) when is_atom(health_atom) do
translate_health_data(%{status: health_atom})
end
defp translate_health_data(health_data) do
Logger.warning("[Dashboard.PubSubBridge] Unknown health data format: #{inspect(health_data)}")
%{status: :unknown, message: "Unknown health status", timestamp: DateTime.utc_now()}
end
defp update_activity(state) do
Map.put(state, :last_activity, DateTime.utc_now())
end
end