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