diff --git a/lib/towerops/equipment.ex b/lib/towerops/equipment.ex index bdce2cba..3d182b1f 100644 --- a/lib/towerops/equipment.ex +++ b/lib/towerops/equipment.ex @@ -123,4 +123,16 @@ defmodule Towerops.Equipment do |> Ecto.Changeset.change(changes) |> Repo.update() end + + @doc """ + Updates the last SNMP poll timestamp for equipment. + Used for distributed coordination to prevent duplicate polling across pods. + """ + def update_snmp_poll_time(%EquipmentSchema{} = equipment) do + now = DateTime.truncate(DateTime.utc_now(), :second) + + equipment + |> Ecto.Changeset.change(%{last_snmp_poll_at: now}) + |> Repo.update() + end end diff --git a/lib/towerops/equipment/equipment.ex b/lib/towerops/equipment/equipment.ex index 31b4cd57..cbdd2047 100644 --- a/lib/towerops/equipment/equipment.ex +++ b/lib/towerops/equipment/equipment.ex @@ -22,6 +22,7 @@ defmodule Towerops.Equipment.Equipment do field :snmp_community, :string field :snmp_port, :integer, default: 161 field :last_discovery_at, :utc_datetime + field :last_snmp_poll_at, :utc_datetime belongs_to :site, Towerops.Sites.Site diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/snmp/poller_worker.ex index 430c21ff..ed200279 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/snmp/poller_worker.ex @@ -82,33 +82,46 @@ defmodule Towerops.Snmp.PollerWorker do device = Snmp.get_device_with_associations(equipment_id) if device do - client_opts = build_client_opts(equipment) - now = DateTime.truncate(DateTime.utc_now(), :second) + # Check if recently polled by another pod (distributed coordination) + poll_interval = get_poll_interval(equipment) + grace_period = 5 - # Poll all sensors - try do - poll_sensors(device.sensors, client_opts, now) - Logger.debug("Polled #{length(device.sensors)} sensors for #{equipment.name}") - rescue - error -> - Logger.error( - "Error polling sensors for #{equipment.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}" - ) + if !should_skip_poll?(equipment, poll_interval, grace_period) do + # Update poll timestamp before polling (optimistic locking) + Equipment.update_snmp_poll_time(equipment) + + client_opts = build_client_opts(equipment) + now = DateTime.truncate(DateTime.utc_now(), :second) + + # Poll all sensors + try do + poll_sensors(device.sensors, client_opts, now) + Logger.debug("Polled #{length(device.sensors)} sensors for #{equipment.name}") + rescue + error -> + Logger.error( + "Error polling sensors for #{equipment.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}" + ) + end + + # Poll all interfaces + try do + poll_interfaces(device.interfaces, client_opts, now) + Logger.debug("Polled #{length(device.interfaces)} interfaces for #{equipment.name}") + rescue + error -> + Logger.error( + "Error polling interfaces for #{equipment.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}" + ) + end + else + Logger.debug( + "Skipping poll for #{equipment.name} - recently polled by another process" + ) end - # Poll all interfaces - try do - poll_interfaces(device.interfaces, client_opts, now) - Logger.debug("Polled #{length(device.interfaces)} interfaces for #{equipment.name}") - rescue - error -> - Logger.error( - "Error polling interfaces for #{equipment.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}" - ) - end - - # Schedule next poll - schedule_next_poll(get_poll_interval(equipment)) + # Schedule next poll regardless of whether we polled or skipped + schedule_next_poll(poll_interval) end end end @@ -266,6 +279,21 @@ defmodule Towerops.Snmp.PollerWorker do Process.send_after(self(), :poll_data, interval_seconds * 1000) end + defp should_skip_poll?(equipment, poll_interval_seconds, grace_period_seconds) do + case equipment.last_snmp_poll_at do + nil -> + false + + last_poll -> + now = DateTime.utc_now() + seconds_since_poll = DateTime.diff(now, last_poll, :second) + min_interval = poll_interval_seconds - grace_period_seconds + + # Skip if polled within the minimum interval + seconds_since_poll < min_interval + end + end + defp via_tuple(equipment_id) do {:via, Registry, {PollerRegistry, equipment_id}} end diff --git a/priv/repo/migrations/20260104192555_add_last_snmp_poll_to_equipment.exs b/priv/repo/migrations/20260104192555_add_last_snmp_poll_to_equipment.exs new file mode 100644 index 00000000..4055cb5e --- /dev/null +++ b/priv/repo/migrations/20260104192555_add_last_snmp_poll_to_equipment.exs @@ -0,0 +1,11 @@ +defmodule Towerops.Repo.Migrations.AddLastSnmpPollToEquipment do + use Ecto.Migration + + def change do + alter table(:equipment) do + add :last_snmp_poll_at, :utc_datetime + end + + create index(:equipment, [:last_snmp_poll_at]) + end +end