diff --git a/lib/towerops/monitoring/device_monitor.ex b/lib/towerops/monitoring/device_monitor.ex index dfadb1fc..0753d46c 100644 --- a/lib/towerops/monitoring/device_monitor.ex +++ b/lib/towerops/monitoring/device_monitor.ex @@ -27,10 +27,11 @@ defmodule Towerops.Monitoring.DeviceMonitor do @doc """ Triggers an immediate check for the device. + Cluster-aware: finds the monitor regardless of which node it's on. """ def trigger_check(device_id) do - # Check if monitor process is running - case Registry.lookup(Towerops.Monitoring.Registry, device_id) do + # Check if monitor process is running (cluster-aware lookup) + case Horde.Registry.lookup(Towerops.Monitoring.Registry, device_id) do [{_pid, _}] -> # Process exists, send cast GenServer.cast(via_tuple(device_id), :check_now) @@ -235,6 +236,6 @@ defmodule Towerops.Monitoring.DeviceMonitor do end defp via_tuple(device_id) do - {:via, Registry, {Towerops.Monitoring.Registry, device_id}} + {:via, Horde.Registry, {Towerops.Monitoring.Registry, device_id}} end end diff --git a/lib/towerops/monitoring/supervisor.ex b/lib/towerops/monitoring/supervisor.ex index 87e31c37..46a6e625 100644 --- a/lib/towerops/monitoring/supervisor.ex +++ b/lib/towerops/monitoring/supervisor.ex @@ -7,9 +7,10 @@ defmodule Towerops.Monitoring.Supervisor do alias Towerops.Monitoring.DeviceMonitor alias Towerops.Snmp.NeighborCleanupWorker alias Towerops.Snmp.PollerRegistry - alias Towerops.Snmp.PollerSupervisor alias Towerops.Snmp.PollerWorker + # Note: DynamicSupervisor names are process names for Horde, not module aliases + require Logger def start_link(init_arg) do @@ -19,16 +20,16 @@ defmodule Towerops.Monitoring.Supervisor do @impl true def init(_init_arg) do base_children = [ - # Registry for naming monitor processes - {Registry, keys: :unique, name: Towerops.Monitoring.Registry}, - # Registry for naming SNMP poller processes - {Registry, keys: :unique, name: PollerRegistry}, - # Task.Supervisor for parallel polling operations + # Horde.Registry for distributed process naming (cluster-aware) + {Horde.Registry, keys: :unique, name: Towerops.Monitoring.Registry, members: :auto}, + # Horde.Registry for SNMP poller processes (cluster-aware) + {Horde.Registry, keys: :unique, name: PollerRegistry, members: :auto}, + # Task.Supervisor for parallel polling operations (local to each node) {Task.Supervisor, name: Towerops.Snmp.PollerTaskSupervisor}, - # DynamicSupervisor for monitor workers - {DynamicSupervisor, name: Towerops.Monitoring.DynamicSupervisor, strategy: :one_for_one}, - # DynamicSupervisor for SNMP poller workers - {DynamicSupervisor, name: PollerSupervisor, strategy: :one_for_one} + # Horde.DynamicSupervisor for monitor workers (distributed across cluster) + {Horde.DynamicSupervisor, name: DynamicSupervisor, strategy: :one_for_one, members: :auto}, + # Horde.DynamicSupervisor for SNMP poller workers (distributed across cluster) + {Horde.DynamicSupervisor, name: PollerSupervisor, strategy: :one_for_one, members: :auto} ] # Only start cleanup worker in non-test environments @@ -70,20 +71,22 @@ defmodule Towerops.Monitoring.Supervisor do @doc """ Starts monitoring for a specific device. + Cluster-aware: only one monitor per device across all nodes. """ def start_monitor(device_id) do spec = {DeviceMonitor, device_id} - DynamicSupervisor.start_child(Towerops.Monitoring.DynamicSupervisor, spec) + Horde.DynamicSupervisor.start_child(DynamicSupervisor, spec) end @doc """ Stops monitoring for a specific device. + Cluster-aware: stops the monitor regardless of which node it's running on. """ def stop_monitor(device_id) do - case Registry.lookup(Towerops.Monitoring.Registry, device_id) do + case Horde.Registry.lookup(Towerops.Monitoring.Registry, device_id) do [{pid, _}] -> - DynamicSupervisor.terminate_child(Towerops.Monitoring.DynamicSupervisor, pid) + Horde.DynamicSupervisor.terminate_child(DynamicSupervisor, pid) [] -> :ok @@ -101,19 +104,21 @@ defmodule Towerops.Monitoring.Supervisor do @doc """ Starts an SNMP poller for a specific device. + Cluster-aware: only one poller per device across all nodes. """ def start_snmp_poller(device_id) do spec = {PollerWorker, device_id} - DynamicSupervisor.start_child(PollerSupervisor, spec) + Horde.DynamicSupervisor.start_child(PollerSupervisor, spec) end @doc """ Stops an SNMP poller for a specific device. + Cluster-aware: stops the poller regardless of which node it's running on. """ def stop_snmp_poller(device_id) do - case Registry.lookup(PollerRegistry, device_id) do + case Horde.Registry.lookup(PollerRegistry, device_id) do [{pid, _}] -> - DynamicSupervisor.terminate_child(PollerSupervisor, pid) + Horde.DynamicSupervisor.terminate_child(PollerSupervisor, pid) [] -> :ok diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/snmp/poller_worker.ex index 9c5cd023..fcff0108 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/snmp/poller_worker.ex @@ -1054,6 +1054,6 @@ defmodule Towerops.Snmp.PollerWorker do end defp via_tuple(device_id) do - {:via, Registry, {PollerRegistry, device_id}} + {:via, Horde.Registry, {PollerRegistry, device_id}} end end diff --git a/mix.exs b/mix.exs index b6286cb3..1275f48c 100644 --- a/mix.exs +++ b/mix.exs @@ -69,6 +69,7 @@ defmodule Towerops.MixProject do {:jason, "~> 1.2"}, {:dns_cluster, "~> 0.2.0"}, {:libcluster, "~> 3.4"}, + {:horde, "~> 0.9.0"}, {:bandit, "~> 1.5"}, {:phoenix_pubsub_redis, "~> 3.0"}, {:ecto_psql_extras, "~> 0.6"}, diff --git a/mix.lock b/mix.lock index e680f6d7..08be4601 100644 --- a/mix.lock +++ b/mix.lock @@ -8,6 +8,7 @@ "credo": {:hex, :credo, "1.7.15", "283da72eeb2fd3ccf7248f4941a0527efb97afa224bcdef30b4b580bc8258e1c", [:mix], [{:bunt, "~> 0.2.1 or ~> 1.0", [hex: :bunt, repo: "hexpm", optional: false]}, {:file_system, "~> 0.2 or ~> 1.0", [hex: :file_system, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}], "hexpm", "291e8645ea3fea7481829f1e1eb0881b8395db212821338e577a90bf225c5607"}, "db_connection": {:hex, :db_connection, "2.9.0", "a6a97c5c958a2d7091a58a9be40caf41ab496b0701d21e1d1abff3fa27a7f371", [:mix], [{:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "17d502eacaf61829db98facf6f20808ed33da6ccf495354a41e64fe42f9c509c"}, "decimal": {:hex, :decimal, "2.3.0", "3ad6255aa77b4a3c4f818171b12d237500e63525c2fd056699967a3e7ea20f62", [:mix], [], "hexpm", "a4d66355cb29cb47c3cf30e71329e58361cfcb37c34235ef3bf1d7bf3773aeac"}, + "delta_crdt": {:hex, :delta_crdt, "0.6.5", "c7bb8c2c7e60f59e46557ab4e0224f67ba22f04c02826e273738f3dcc4767adc", [:mix], [{:merkle_map, "~> 0.2.0", [hex: :merkle_map, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "c6ae23a525d30f96494186dd11bf19ed9ae21d9fe2c1f1b217d492a7cc7294ae"}, "dialyxir": {:hex, :dialyxir, "1.4.7", "dda948fcee52962e4b6c5b4b16b2d8fa7d50d8645bbae8b8685c3f9ecb7f5f4d", [:mix], [{:erlex, ">= 0.2.8", [hex: :erlex, repo: "hexpm", optional: false]}], "hexpm", "b34527202e6eb8cee198efec110996c25c5898f43a4094df157f8d28f27d9efe"}, "dns_cluster": {:hex, :dns_cluster, "0.2.0", "aa8eb46e3bd0326bd67b84790c561733b25c5ba2fe3c7e36f28e88f384ebcb33", [:mix], [], "hexpm", "ba6f1893411c69c01b9e8e8f772062535a4cf70f3f35bcc964a324078d8c8240"}, "ecto": {:hex, :ecto, "3.13.5", "9d4a69700183f33bf97208294768e561f5c7f1ecf417e0fa1006e4a91713a834", [:mix], [{:decimal, "~> 2.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "df9efebf70cf94142739ba357499661ef5dbb559ef902b68ea1f3c1fabce36de"}, @@ -25,11 +26,14 @@ "gen_smtp": {:hex, :gen_smtp, "1.3.0", "62c3d91f0dcf6ce9db71bcb6881d7ad0d1d834c7f38c13fa8e952f4104a8442e", [:rebar3], [{:ranch, ">= 1.8.0", [hex: :ranch, repo: "hexpm", optional: false]}], "hexpm", "0b73fbf069864ecbce02fe653b16d3f35fd889d0fdd4e14527675565c39d84e6"}, "gettext": {:hex, :gettext, "1.0.2", "5457e1fd3f4abe47b0e13ff85086aabae760497a3497909b8473e0acee57673b", [:mix], [{:expo, "~> 0.5.1 or ~> 1.0", [hex: :expo, repo: "hexpm", optional: false]}], "hexpm", "eab805501886802071ad290714515c8c4a17196ea76e5afc9d06ca85fb1bfeb3"}, "heroicons": {:git, "https://github.com/tailwindlabs/heroicons.git", "0435d4ca364a608cc75e2f8683d374e55abbae26", [tag: "v2.2.0", sparse: "optimized", depth: 1]}, + "horde": {:hex, :horde, "0.9.1", "547507dfe1228471d9a31cece7653be81df2cf8f210816f672ce507dc9606128", [:mix], [{:delta_crdt, "~> 0.6.2", [hex: :delta_crdt, repo: "hexpm", optional: false]}, {:libring, "~> 1.7", [hex: :libring, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.0 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:telemetry_poller, "~> 0.5.0 or ~> 1.0", [hex: :telemetry_poller, repo: "hexpm", optional: false]}], "hexpm", "14e1be4d6a84b066101641628a64f849d6df4f0311e93b7a9c39709e9c0bf4ba"}, "hpax": {:hex, :hpax, "1.0.3", "ed67ef51ad4df91e75cc6a1494f851850c0bd98ebc0be6e81b026e765ee535aa", [:mix], [], "hexpm", "8eab6e1cfa8d5918c2ce4ba43588e894af35dbd8e91e6e55c817bca5847df34a"}, "idna": {:hex, :idna, "6.1.1", "8a63070e9f7d0c62eb9d9fcb360a7de382448200fbbd1b106cc96d3d8099df8d", [:rebar3], [{:unicode_util_compat, "~> 0.7.0", [hex: :unicode_util_compat, repo: "hexpm", optional: false]}], "hexpm", "92376eb7894412ed19ac475e4a86f7b413c1b9fbb5bd16dccd57934157944cea"}, "jason": {:hex, :jason, "1.4.4", "b9226785a9aa77b6857ca22832cffa5d5011a667207eb2a0ad56adb5db443b8a", [:mix], [{:decimal, "~> 1.0 or ~> 2.0", [hex: :decimal, repo: "hexpm", optional: true]}], "hexpm", "c5eb0cab91f094599f94d55bc63409236a8ec69a21a67814529e8d5f6cc90b3b"}, "lazy_html": {:hex, :lazy_html, "0.1.8", "677a8642e644eef8de98f3040e2520d42d0f0f8bd6c5cd49db36504e34dffe91", [:make, :mix], [{:cc_precompiler, "~> 0.1", [hex: :cc_precompiler, repo: "hexpm", optional: false]}, {:elixir_make, "~> 0.9.0", [hex: :elixir_make, repo: "hexpm", optional: false]}, {:fine, "~> 0.1.0", [hex: :fine, repo: "hexpm", optional: false]}], "hexpm", "0d8167d930b704feb94b41414ca7f5779dff9bca7fcf619fcef18de138f08736"}, "libcluster": {:hex, :libcluster, "3.5.0", "5ee4cfde4bdf32b2fef271e33ce3241e89509f4344f6c6a8d4069937484866ba", [:mix], [{:jason, "~> 1.1", [hex: :jason, repo: "hexpm", optional: false]}, {:telemetry, "~> 1.3", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "ebf6561fcedd765a4cd43b4b8c04b1c87f4177b5fb3cbdfe40a780499d72f743"}, + "libring": {:hex, :libring, "1.7.0", "4f245d2f1476cd7ed8f03740f6431acba815401e40299208c7f5c640e1883bda", [:mix], [], "hexpm", "070e3593cb572e04f2c8470dd0c119bc1817a7a0a7f88229f43cf0345268ec42"}, + "merkle_map": {:hex, :merkle_map, "0.2.2", "f36ff730cca1f2658e317a3c73406f50bbf5ac8aff54cf837d7ca2069a6e251c", [:mix], [], "hexpm", "383107f0503f230ac9175e0631647c424efd027e89ea65ab5ea12eeb54257aaf"}, "mime": {:hex, :mime, "2.0.7", "b8d739037be7cd402aee1ba0306edfdef982687ee7e9859bee6198c1e7e2f128", [:mix], [], "hexpm", "6171188e399ee16023ffc5b76ce445eb6d9672e2e241d2df6050f3c771e80ccd"}, "mint": {:hex, :mint, "1.7.1", "113fdb2b2f3b59e47c7955971854641c61f378549d73e829e1768de90fc1abf1", [:mix], [{:castore, "~> 0.1.0 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:hpax, "~> 0.1.1 or ~> 0.2.0 or ~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}], "hexpm", "fceba0a4d0f24301ddee3024ae116df1c3f4bb7a563a731f45fdfeb9d39a231b"}, "mix_audit": {:hex, :mix_audit, "2.1.5", "c0f77cee6b4ef9d97e37772359a187a166c7a1e0e08b50edf5bf6959dfe5a016", [:make, :mix], [{:jason, "~> 1.4", [hex: :jason, repo: "hexpm", optional: false]}, {:yaml_elixir, "~> 2.11", [hex: :yaml_elixir, repo: "hexpm", optional: false]}], "hexpm", "87f9298e21da32f697af535475860dc1d3617a010e0b418d2ec6142bc8b42d69"},