feat: add Horde for cluster-aware distributed supervision

Prevents duplicate polling in multi-node production environment:
- Replace DynamicSupervisor with Horde.DynamicSupervisor (distributed)
- Replace Registry with Horde.Registry (cluster-aware)
- Each device monitor runs on exactly ONE node across the cluster
- Each SNMP poller runs on exactly ONE node across the cluster
- Automatic failover when nodes go down
- Automatic rebalancing when nodes join/leave

Architecture:
- Horde.Registry for process naming (distributed lookup)
- Horde.DynamicSupervisor for DeviceMonitor processes
- Horde.DynamicSupervisor for PollerWorker processes
- Task.Supervisor remains local per-node for parallel operations
- members: :auto enables automatic cluster discovery via libcluster

Production impact:
- 2 replicas = devices split ~50/50 across both pods
- No duplicate SNMP polls to same device
- If pod dies, surviving pod takes over all devices
- When new pod starts, devices rebalance automatically

Dependencies added:
- horde ~> 0.9.0 (distributed supervisor/registry)
- delta_crdt (Horde dependency for CRDT-based state)
- libring (Horde dependency for consistent hashing)
This commit is contained in:
Graham McIntire 2026-01-18 10:14:50 -06:00
parent fed1b3ecfc
commit 3cc69c9250
No known key found for this signature in database
5 changed files with 31 additions and 20 deletions

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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"},

View file

@ -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"},