diff --git a/lib/towerops/application.ex b/lib/towerops/application.ex index f9be1d3e..da69b987 100644 --- a/lib/towerops/application.ex +++ b/lib/towerops/application.ex @@ -26,7 +26,7 @@ defmodule Towerops.Application do ToweropsWeb.Telemetry, Towerops.Repo, {DNSCluster, query: Application.get_env(:towerops, :dns_cluster_query) || :ignore}, - {Phoenix.PubSub, name: Towerops.PubSub}, + {Phoenix.PubSub, name: Towerops.PubSub, adapter: Phoenix.PubSub.Redis, redis_config: redis_pubsub_config()}, # Start event logger (subscribes to PubSub) Towerops.Devices.EventLogger, # Start monitoring supervisor @@ -78,6 +78,16 @@ defmodule Towerops.Application do ] end + defp redis_pubsub_config do + redis_config = Application.get_env(:towerops, :redis, []) + + [ + host: Keyword.get(redis_config, :host, "localhost"), + port: Keyword.get(redis_config, :port, 6379), + node_name: node() + ] + end + @doc """ Returns the deployment timestamp. diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/snmp/poller_worker.ex index 127b1481..9bb93c65 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/snmp/poller_worker.ex @@ -137,6 +137,13 @@ defmodule Towerops.Snmp.PollerWorker do defp poll_device_sensors(device, snmp_device, client_opts, now) do poll_sensors(snmp_device.sensors, client_opts, now) Logger.debug("Polled #{length(snmp_device.sensors)} sensors for #{device.name}") + + # Broadcast sensor update event to device-specific topic + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{device.id}", + {:sensors_updated, device.id} + ) rescue error -> Logger.error("Error polling sensors for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}") @@ -145,6 +152,13 @@ defmodule Towerops.Snmp.PollerWorker do defp poll_device_interfaces(device, snmp_device, client_opts, now) do poll_interfaces(snmp_device.interfaces, client_opts, now) Logger.debug("Polled #{length(snmp_device.interfaces)} interfaces for #{device.name}") + + # Broadcast interface stats update event to device-specific topic + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{device.id}", + {:interfaces_updated, device.id} + ) rescue error -> Logger.error("Error polling interfaces for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}") @@ -175,6 +189,13 @@ defmodule Towerops.Snmp.PollerWorker do end) Logger.debug("Polled and saved #{length(neighbors)} neighbors for #{device.name}") + + # Broadcast neighbor update event to device-specific topic + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{device.id}", + {:neighbors_updated, device.id} + ) rescue error -> Logger.error("Error polling neighbors for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}") @@ -518,6 +539,15 @@ defmodule Towerops.Snmp.PollerWorker do defp broadcast_interface_events(events) do Enum.each(events, fn {:event, event_attrs} -> + # Broadcast to device-specific topic for real-time updates + _ = + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{event_attrs.device_id}", + {:device_event, event_attrs} + ) + + # Also broadcast to generic events topic for global event logger _ = Phoenix.PubSub.broadcast( Towerops.PubSub, @@ -923,7 +953,12 @@ defmodule Towerops.Snmp.PollerWorker do defp broadcast_sensor_events(events) do Enum.each(events, fn event -> + # Broadcast to device-specific topic for real-time updates + _ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:#{event.device_id}", {:device_event, event}) + + # Also broadcast to generic events topic for global event logger _ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:events", {:device_event, event}) + Logger.debug("Sensor event: #{event.message}") end) end diff --git a/lib/towerops/workers/monitor_worker.ex b/lib/towerops/workers/monitor_worker.ex index 0d78b669..34df0627 100644 --- a/lib/towerops/workers/monitor_worker.ex +++ b/lib/towerops/workers/monitor_worker.ex @@ -41,6 +41,13 @@ defmodule Towerops.Workers.MonitorWorker do checked_at: DateTime.utc_now() }) + # Broadcast monitoring check update to device-specific topic + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{device_id}", + {:monitoring_check_updated, device_id} + ) + :ok {:error, reason} -> @@ -53,6 +60,13 @@ defmodule Towerops.Workers.MonitorWorker do checked_at: DateTime.utc_now() }) + # Broadcast monitoring check update to device-specific topic + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{device_id}", + {:monitoring_check_updated, device_id} + ) + :ok end else diff --git a/lib/towerops_web/live/device_live/show.ex b/lib/towerops_web/live/device_live/show.ex index faaa294f..5cc09570 100644 --- a/lib/towerops_web/live/device_live/show.ex +++ b/lib/towerops_web/live/device_live/show.ex @@ -72,6 +72,33 @@ defmodule ToweropsWeb.DeviceLive.Show do |> put_flash(:info, "Discovery completed")} end + @impl true + def handle_info({:sensors_updated, _device_id}, socket) do + {:noreply, load_equipment_data(socket, socket.assigns.device.id)} + end + + @impl true + def handle_info({:interfaces_updated, _device_id}, socket) do + {:noreply, load_equipment_data(socket, socket.assigns.device.id)} + end + + @impl true + def handle_info({:neighbors_updated, _device_id}, socket) do + {:noreply, load_equipment_data(socket, socket.assigns.device.id)} + end + + @impl true + def handle_info({:monitoring_check_updated, _device_id}, socket) do + {:noreply, load_equipment_data(socket, socket.assigns.device.id)} + end + + @impl true + def handle_info({:device_event, _event_attrs}, socket) do + # Device event logged (interface changes, sensor thresholds, etc.) + # Reload to update events list + {:noreply, load_equipment_data(socket, socket.assigns.device.id)} + end + # Private functions defp load_equipment_data(socket, device_id) do diff --git a/mix.exs b/mix.exs index ad785ee4..b6286cb3 100644 --- a/mix.exs +++ b/mix.exs @@ -70,6 +70,7 @@ defmodule Towerops.MixProject do {:dns_cluster, "~> 0.2.0"}, {:libcluster, "~> 3.4"}, {:bandit, "~> 1.5"}, + {:phoenix_pubsub_redis, "~> 3.0"}, {:ecto_psql_extras, "~> 0.6"}, {:exq, "~> 0.19"}, {:mox, "~> 1.0", only: :test}, diff --git a/mix.lock b/mix.lock index abc8aa05..e680f6d7 100644 --- a/mix.lock +++ b/mix.lock @@ -45,9 +45,11 @@ "phoenix_live_reload": {:hex, :phoenix_live_reload, "1.6.2", "b18b0773a1ba77f28c52decbb0f10fd1ac4d3ae5b8632399bbf6986e3b665f62", [:mix], [{:file_system, "~> 0.2.10 or ~> 1.0", [hex: :file_system, repo: "hexpm", optional: false]}, {:phoenix, "~> 1.4", [hex: :phoenix, repo: "hexpm", optional: false]}], "hexpm", "d1f89c18114c50d394721365ffb428cce24f1c13de0467ffa773e2ff4a30d5b9"}, "phoenix_live_view": {:hex, :phoenix_live_view, "1.1.20", "4f20850ee700b309b21906a0e510af1b916b454b4f810fb8581ada016eb42dfc", [:mix], [{:igniter, ">= 0.6.16 and < 1.0.0-0", [hex: :igniter, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:lazy_html, "~> 0.1.0", [hex: :lazy_html, repo: "hexpm", optional: true]}, {:phoenix, "~> 1.6.15 or ~> 1.7.0 or ~> 1.8.0-rc", [hex: :phoenix, repo: "hexpm", optional: false]}, {:phoenix_html, "~> 3.3 or ~> 4.0", [hex: :phoenix_html, repo: "hexpm", optional: false]}, {:phoenix_template, "~> 1.0", [hex: :phoenix_template, repo: "hexpm", optional: false]}, {:phoenix_view, "~> 2.0", [hex: :phoenix_view, repo: "hexpm", optional: true]}, {:plug, "~> 1.15", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.2 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "c16abd605a21f778165cb0079946351ef20ef84eb1ef467a862fb9a173b1d27d"}, "phoenix_pubsub": {:hex, :phoenix_pubsub, "2.2.0", "ff3a5616e1bed6804de7773b92cbccfc0b0f473faf1f63d7daf1206c7aeaaa6f", [:mix], [], "hexpm", "adc313a5bf7136039f63cfd9668fde73bba0765e0614cba80c06ac9460ff3e96"}, + "phoenix_pubsub_redis": {:hex, :phoenix_pubsub_redis, "3.0.1", "d4d856b1e57a21358e448543e1d091e07e83403dde4383b8be04ed9d2c201cbc", [:mix], [{:phoenix_pubsub, "~> 2.0", [hex: :phoenix_pubsub, repo: "hexpm", optional: false]}, {:poolboy, "~> 1.5.1 or ~> 1.6", [hex: :poolboy, repo: "hexpm", optional: false]}, {:redix, "~> 0.10.0 or ~> 1.0", [hex: :redix, repo: "hexpm", optional: false]}], "hexpm", "0b36a17ff6e9a56159f8df8933d62b5c1f0695eae995a02e0c86c035ace6a309"}, "phoenix_template": {:hex, :phoenix_template, "1.0.4", "e2092c132f3b5e5b2d49c96695342eb36d0ed514c5b252a77048d5969330d639", [:mix], [{:phoenix_html, "~> 2.14.2 or ~> 3.0 or ~> 4.0", [hex: :phoenix_html, repo: "hexpm", optional: true]}], "hexpm", "2c0c81f0e5c6753faf5cca2f229c9709919aba34fab866d3bc05060c9c444206"}, "plug": {:hex, :plug, "1.19.1", "09bac17ae7a001a68ae393658aa23c7e38782be5c5c00c80be82901262c394c0", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:plug_crypto, "~> 1.1.1 or ~> 1.2 or ~> 2.0", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.3 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "560a0017a8f6d5d30146916862aaf9300b7280063651dd7e532b8be168511e62"}, "plug_crypto": {:hex, :plug_crypto, "2.1.1", "19bda8184399cb24afa10be734f84a16ea0a2bc65054e23a62bb10f06bc89491", [:mix], [], "hexpm", "6470bce6ffe41c8bd497612ffde1a7e4af67f36a15eea5f921af71cf3e11247c"}, + "poolboy": {:hex, :poolboy, "1.5.2", "392b007a1693a64540cead79830443abf5762f5d30cf50bc95cb2c1aaafa006b", [:rebar3], [], "hexpm", "dad79704ce5440f3d5a3681c8590b9dc25d1a561e8f5a9c995281012860901e3"}, "postgrex": {:hex, :postgrex, "0.22.0", "fb027b58b6eab1f6de5396a2abcdaaeb168f9ed4eccbb594e6ac393b02078cbd", [:mix], [{:db_connection, "~> 2.9", [hex: :db_connection, repo: "hexpm", optional: false]}, {:decimal, "~> 1.5 or ~> 2.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:table, "~> 0.1.0", [hex: :table, repo: "hexpm", optional: true]}], "hexpm", "a68c4261e299597909e03e6f8ff5a13876f5caadaddd0d23af0d0a61afcc5d84"}, "protobuf": {:hex, :protobuf, "0.16.0", "d1878725105d49162977cf3408ccc3eac4f3532e26e5a9e250f2c624175d10f6", [:mix], [{:jason, "~> 1.2", [hex: :jason, repo: "hexpm", optional: true]}], "hexpm", "f0d0d3edd8768130f24cc2cfc41320637d32c80110e80d13f160fa699102c828"}, "ranch": {:hex, :ranch, "2.2.0", "25528f82bc8d7c6152c57666ca99ec716510fe0925cb188172f41ce93117b1b0", [:make, :rebar3], [], "hexpm", "fa0b99a1780c80218a4197a59ea8d3bdae32fbff7e88527d7d8a4787eff4f8e7"},