towerops/lib/towerops/workers/device_poller_worker.ex

1929 lines
62 KiB
Elixir

defmodule Towerops.Workers.DevicePollerWorker do
@moduledoc """
Oban worker that performs SNMP polling for a device.
This worker continuously collects:
- Sensor readings (temperature, voltage, CPU, memory, etc.)
- Interface statistics (bandwidth, errors, discards, etc.)
- Neighbor discovery (LLDP/CDP topology information)
- ARP and MAC tables
- Processor and storage statistics
Uses Oban's unique job feature to ensure only one poller per device
runs cluster-wide. Poll interval is configurable per device (default: 60 seconds).
"""
use Oban.Pro.Worker,
queue: :pollers,
max_attempts: 1,
unique: [
period: :infinity,
keys: [:device_id],
states: [:available, :scheduled, :executing, :retryable, :suspended]
],
replace: [
scheduled: [:scheduled_at],
available: [:scheduled_at]
]
import Ecto.Query
alias Towerops.Agents
alias Towerops.Devices
alias Towerops.Gaiia.SubscriberMatching
alias Towerops.JobMonitoring.Events
alias Towerops.Snmp
alias Towerops.Snmp.ArpDiscovery
alias Towerops.Snmp.Client
alias Towerops.Snmp.MacDiscovery
alias Towerops.Snmp.NeighborDiscovery
alias Towerops.Snmp.RadioState
alias Towerops.Snmp.SensorChangeDetector
alias Towerops.Snmp.SensorScale
alias Towerops.Snmp.WirelessClientDiscovery
alias Towerops.Workers.PollingOffset
require Logger
@default_poll_interval 60
@impl Oban.Pro.Worker
@spec process(Oban.Job.t()) :: :ok
def process(%Oban.Job{args: %{"device_id" => device_id}} = job) do
_ = Events.broadcast_job_event(job, :started)
start_time = System.monotonic_time(:second)
result =
case Devices.get_device(device_id) do
nil ->
Logger.debug("Device #{device_id} no longer exists, skipping poll and not rescheduling")
:ok
device ->
# Only reschedule if device still has polling enabled
if should_continue_polling?(device) do
maybe_poll_device(device)
schedule_next_poll_with_error_handling(device_id, device)
else
Logger.debug("Device #{device_id} polling disabled, not rescheduling",
device_id: device_id,
snmp_enabled: device.snmp_enabled
)
end
:ok
end
duration = System.monotonic_time(:second) - start_time
_ = Events.broadcast_job_event(job, :completed, %{duration: duration})
result
end
# Check if device should continue being polled
# Returns false if Phoenix SNMP is disabled, SNMP is disabled on device, or device has any agent
defp should_continue_polling?(device) do
!Client.phoenix_snmp_disabled() && device.snmp_enabled &&
!Agents.device_has_effective_agent?(device.id)
end
# Agent check already done by should_continue_polling?, no need to re-check
defp maybe_poll_device(device) do
poll_device(device)
end
defp schedule_next_poll_with_error_handling(device_id, _stale_device) do
# Re-fetch device to get latest check_interval_seconds
# Configuration could have changed during the polling execution
case Devices.get_device(device_id) do
nil ->
Logger.debug("Device deleted, not rescheduling", device_id: device_id)
:ok
fresh_device ->
poll_interval = get_poll_interval(fresh_device)
case schedule_next_poll(device_id, poll_interval) do
{:ok, _job} ->
:ok
{:error, changeset} ->
Logger.error("Failed to schedule next poll for device #{device_id}: #{inspect(changeset.errors)}")
:ok
end
end
end
@doc """
Starts polling for a device.
"""
def start_polling(device_id) do
device = Devices.get_device!(device_id)
interval = get_poll_interval(device)
offset = PollingOffset.calculate_offset(device_id, interval)
%{device_id: device_id}
|> new(schedule_in: offset)
|> Oban.insert()
end
@doc """
Stops polling for a device by cancelling its jobs.
"""
def stop_polling(device_id) do
Oban.cancel_all_jobs(
from(j in Oban.Job,
where: j.worker == "Towerops.Workers.DevicePollerWorker",
where: fragment("args->>'device_id' = ?", ^device_id),
where: j.state in ["available", "scheduled", "executing", "retryable"]
)
)
end
@doc """
Triggers an immediate poll for a device.
"""
def trigger_poll(device_id) do
# Insert a new job to run immediately
%{device_id: device_id}
|> new()
|> Oban.insert()
end
# Private Functions
defp poll_device(device) do
snmp_device = Snmp.get_device_with_associations(device.id)
if snmp_device && (snmp_device.sensors != [] || snmp_device.interfaces != []) do
poll_interval = get_poll_interval(device)
if should_skip_poll?(device, poll_interval, 5) do
Logger.debug("Skipping poll for #{device.name} - recently polled by another process")
else
perform_device_data_collection(device, snmp_device)
end
end
end
defp perform_device_data_collection(device, snmp_device) do
Devices.update_snmp_poll_time(device)
client_opts = build_client_opts(device)
now = Towerops.Time.now()
# Run all polling operations in parallel using async tasks
task_names = [
"sensors",
"state_sensors",
"interfaces",
"interface_changes",
"neighbors",
"arp",
"mac",
"processors",
"storage",
"mempools",
"wireless_clients",
"transceivers",
"entity_physical"
]
tasks = [
Task.async(fn -> poll_device_sensors(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_state_sensors(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_interfaces(device, snmp_device, client_opts, now) end),
Task.async(fn -> check_device_interface_changes(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_neighbors(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_arp(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_mac(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_processors(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_storage(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_mempools(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_wireless_clients(device, snmp_device, client_opts) end),
Task.async(fn -> poll_device_transceivers(device, snmp_device, client_opts, now) end),
Task.async(fn -> poll_device_entity_physical(device, snmp_device, client_opts, now) end)
]
# Wait for all tasks with timeout, log failures
# Explicitly pair results with names to handle length mismatches
results = Task.yield_many(tasks, 30_000)
# Validate we got results for all tasks
if length(results) != length(task_names) do
Logger.error(
"Task result count mismatch for #{device.name}: expected #{length(task_names)}, got #{length(results)}",
device_id: device.id,
expected: length(task_names),
actual: length(results)
)
end
# Process results with explicit pairing to handle any length mismatch
results
|> Enum.zip(task_names)
|> Enum.each(fn {{task, result}, name} ->
case result do
{:ok, _value} ->
:ok
{:exit, reason} ->
Logger.warning("Polling task '#{name}' crashed for #{device.name}: #{inspect(reason)}",
device_id: device.id,
task: name
)
nil ->
_ = Task.shutdown(task, :brutal_kill)
Logger.warning("Polling task '#{name}' timed out for #{device.name}",
device_id: device.id,
task: name
)
end
end)
# Re-check assignment before processing results
# Device could have been reassigned to an agent during the 30-second polling window
case verify_polling_assignment_unchanged(device.id) do
:ok ->
# Assignment unchanged - process results normally
:ok
:reassigned ->
Logger.info("Device reassigned during polling, discarding results",
device_id: device.id,
device_name: device.name
)
:ok
end
end
# Verify device assignment hasn't changed since polling started
defp verify_polling_assignment_unchanged(device_id) do
case Devices.get_device(device_id) do
nil ->
# Device deleted during polling
:reassigned
device ->
if Agents.device_has_effective_agent?(device.id) do
:reassigned
else
:ok
end
end
end
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}")
# Refresh AP-level radio state (current_channel / current_frequency_mhz)
# from the just-polled frequency sensors so frequency-aware rules read
# a single canonical source.
_ = RadioState.update_from_sensors(snmp_device)
_ =
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()}")
end
defp poll_device_state_sensors(device, snmp_device, client_opts, now) do
poll_state_sensors(snmp_device.state_sensors, client_opts, now, device.id, device.organization_id)
Logger.debug("Polled #{length(snmp_device.state_sensors)} state sensors for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:state_sensors_updated, device.id}
)
rescue
error ->
Logger.error("Error polling state sensors for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_device_interfaces(device, snmp_device, client_opts, now) do
_ = poll_interfaces(snmp_device.interfaces, client_opts, now, snmp_device)
Logger.debug("Polled #{length(snmp_device.interfaces)} interfaces for #{device.name}")
_ =
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()}")
end
defp check_device_interface_changes(device, snmp_device, client_opts, now) do
check_interface_changes(snmp_device.interfaces, device, client_opts, now)
rescue
error ->
Logger.error(
"Error checking interface changes for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}"
)
end
defp poll_device_neighbors(device, snmp_device, client_opts) do
interfaces_with_device = Enum.map(snmp_device.interfaces, &Map.put(&1, :device_id, device.id))
{:ok, neighbors} = NeighborDiscovery.discover_neighbors(client_opts, interfaces_with_device)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_and_upsert_neighbors(device.id, neighbors, cutoff)
Logger.debug("Polled and saved #{length(neighbors)} neighbors for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:neighbors_updated, device.id}
)
# Run topology inference after neighbor data is saved
case Towerops.Topology.process_device(device, device.organization_id) do
{:ok, :changed} ->
Logger.debug("Topology updated for #{device.name}")
{:ok, :unchanged} ->
:ok
end
rescue
error ->
Logger.error("Error polling neighbors for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_device_arp(device, snmp_device, client_opts) do
{:ok, arp_entries} = ArpDiscovery.discover_arp_table(client_opts)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_and_upsert_arp_entries(device.id, arp_entries, snmp_device.interfaces, cutoff)
Logger.debug("Polled and saved #{length(arp_entries)} ARP entries for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:arp_updated, device.id}
)
# Refresh subscriber-to-device links based on new ARP data
SubscriberMatching.refresh_links_for_device(device.id)
rescue
error ->
Logger.error("Error polling ARP table for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_device_mac(device, snmp_device, client_opts) do
{:ok, mac_entries} = MacDiscovery.discover_mac_table(client_opts)
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
Snmp.delete_stale_and_upsert_mac_addresses(device.id, mac_entries, snmp_device.interfaces, cutoff)
Logger.debug("Polled and saved #{length(mac_entries)} MAC FDB entries for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:mac_updated, device.id}
)
rescue
error ->
Logger.error("Error polling MAC FDB table for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_wireless_clients(device, snmp_device, client_opts) do
vendor_profile = WirelessClientDiscovery.detect_vendor_profile(snmp_device)
if vendor_profile do
{:ok, clients} = WirelessClientDiscovery.discover_wireless_clients(client_opts, vendor_profile)
process_wireless_clients(device, clients)
end
rescue
error ->
Logger.error(
"Error polling wireless clients for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}"
)
end
defp process_wireless_clients(_device, []), do: :ok
defp process_wireless_clients(device, clients) do
now = Towerops.Time.now()
cutoff = DateTime.add(now, -5, :minute)
Snmp.delete_stale_and_upsert_wireless_clients(device.id, device.organization_id, clients, cutoff)
# Load upserted wireless clients to create historical readings
wireless_clients = Snmp.list_wireless_clients(device.id)
# Build and insert historical readings
save_wireless_client_readings(device, wireless_clients, now)
Logger.debug("Polled and saved #{length(clients)} wireless clients for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"wireless_clients:device:#{device.id}",
{:wireless_clients_updated, device.id, length(wireless_clients)}
)
# Refresh subscriber-to-device links with new wireless client data
SubscriberMatching.refresh_links_for_device(device.id)
end
defp save_wireless_client_readings(_device, [], _now), do: :ok
defp save_wireless_client_readings(device, wireless_clients, now) do
reading_entries =
Enum.map(wireless_clients, fn client ->
%{
device_id: device.id,
wireless_client_id: client.id,
organization_id: device.organization_id,
mac_address: client.mac_address,
ip_address: client.ip_address,
signal_strength: client.signal_strength,
snr: client.snr,
distance: client.distance,
tx_rate: client.tx_rate,
rx_rate: client.rx_rate,
uptime_seconds: client.uptime_seconds,
metadata: client.metadata || %{},
checked_at: now
}
end)
{count, _} = Snmp.create_wireless_client_readings_batch(reading_entries)
Logger.debug("Inserted #{count} wireless client readings for #{device.name}")
end
defp poll_device_transceivers(device, snmp_device, client_opts, now) do
transceivers = snmp_device.transceivers || []
if transceivers != [] do
poll_transceivers(transceivers, client_opts, now)
Logger.debug("Polled transceivers for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:transceivers_updated, device.id}
)
end
rescue
error ->
Logger.error("Error polling transceivers for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_device_entity_physical(device, snmp_device, client_opts, now) do
entities = snmp_device.entity_physical || []
if entities != [] do
poll_entity_physical_status(entities, client_opts, now)
Logger.debug("Polled entity physical status for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:hardware_inventory_updated, device.id}
)
end
rescue
error ->
Logger.error(
"Error polling entity physical for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}"
)
end
defp poll_device_processors(device, snmp_device, client_opts, now) do
processors = snmp_device.processors || []
if processors != [] do
poll_processors(processors, client_opts, now)
Logger.debug("Polled processors for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:processors_updated, device.id}
)
end
rescue
error ->
Logger.error("Error polling processors for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_device_storage(device, snmp_device, client_opts, now) do
storage_entries = snmp_device.storage || []
if storage_entries != [] do
poll_storage(storage_entries, client_opts, now)
Logger.debug("Polled #{length(storage_entries)} storage entries for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:storage_updated, device.id}
)
end
rescue
error ->
Logger.error("Error polling storage for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
defp poll_device_mempools(device, snmp_device, client_opts, now) do
mempools = snmp_device.mempools || []
if mempools != [] do
poll_mempools(mempools, client_opts, now)
Logger.debug("Polled #{length(mempools)} memory pools for #{device.name}")
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:mempools_updated, device.id}
)
end
rescue
error ->
Logger.error("Error polling memory pools for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}")
end
# All the helper functions from the original PollerWorker GenServer
# (These remain unchanged - I'll include them all below)
defp poll_storage(storage_entries, client_opts, timestamp) do
poll_results =
storage_entries
|> Task.async_stream(
fn storage ->
result = poll_storage_value(storage, client_opts)
{storage, result}
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, pair} -> [pair]
{:exit, _reason} -> []
end)
# Build batch reading entries (only for successful polls)
reading_entries =
Enum.flat_map(poll_results, fn
{storage, {:ok, values}} ->
[
%{
storage_id: storage.id,
used_bytes: values.used_bytes,
total_bytes: values.total_bytes,
usage_percent: values.usage_percent,
checked_at: timestamp
}
]
{_storage, {:error, _reason}} ->
[]
end)
_ = Snmp.create_storage_readings_batch(reading_entries)
# Process metadata updates individually
Enum.each(poll_results, fn
{storage, {:ok, values}} ->
Snmp.update_storage(storage, %{
used_bytes: values.used_bytes,
total_bytes: values.total_bytes,
last_checked_at: timestamp
})
{storage, {:error, reason}} ->
Logger.debug("Failed to poll storage #{storage.description}: #{inspect(reason)}")
end)
end
defp poll_storage_value(storage, client_opts) do
used_oid = "1.3.6.1.2.1.25.2.3.1.6.#{storage.storage_index}"
size_oid = "1.3.6.1.2.1.25.2.3.1.5.#{storage.storage_index}"
alloc_oid = "1.3.6.1.2.1.25.2.3.1.4.#{storage.storage_index}"
with {:ok, used_units} <- Client.get(client_opts, used_oid),
{:ok, size_units} <- Client.get(client_opts, size_oid),
{:ok, alloc_units} <- Client.get(client_opts, alloc_oid),
true <- is_integer(used_units) and is_integer(size_units) and is_integer(alloc_units),
true <- size_units > 0 do
used_bytes = used_units * alloc_units
total_bytes = size_units * alloc_units
usage_percent = used_bytes / total_bytes * 100
{:ok, %{used_bytes: used_bytes, total_bytes: total_bytes, usage_percent: usage_percent}}
else
{:error, reason} -> {:error, reason}
false -> {:error, :invalid_values}
end
end
defp poll_processors(processors, client_opts, timestamp) do
poll_results =
processors
|> Task.async_stream(
fn processor ->
result = poll_processor_value(processor, client_opts)
{processor, result}
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, pair} -> [pair]
{:exit, _reason} -> []
end)
# Build batch reading entries
reading_entries =
Enum.map(poll_results, fn {processor, result} ->
build_processor_reading_entry(processor, result, timestamp)
end)
_ = Snmp.create_processor_readings_batch(reading_entries)
# Process metadata updates individually
Enum.each(poll_results, fn {processor, result} ->
handle_processor_poll_metadata(processor, result, timestamp)
end)
end
defp poll_processor_value(processor, client_opts) do
# Extract numeric index from processor_index (e.g., "hr_196608" -> "196608")
numeric_index = processor.processor_index |> String.split("_") |> List.last()
oid =
case processor.processor_type do
"hr_processor" ->
"1.3.6.1.2.1.25.3.3.1.2.#{numeric_index}"
"cisco_cpu" ->
"1.3.6.1.4.1.9.9.109.1.1.1.1.8.#{numeric_index}"
"ucd_cpu" ->
poll_ucd_cpu(processor, client_opts)
_ ->
{:error, :unknown_processor_type}
end
case oid do
{:ok, value} -> {:ok, value}
{:error, _} = error -> error
oid_string when is_binary(oid_string) -> poll_snmp_load(client_opts, oid_string)
end
end
defp poll_snmp_load(client_opts, oid) do
case Client.get(client_opts, oid) do
{:ok, value} when is_integer(value) ->
{:ok, value / 1.0}
{:ok, _non_integer} ->
{:error, :non_integer}
{:error, reason} ->
{:error, reason}
end
end
defp poll_ucd_cpu(_processor, client_opts) do
user_oid = "1.3.6.1.4.1.2021.11.9.0"
system_oid = "1.3.6.1.4.1.2021.11.10.0"
with {:ok, user} <- Client.get(client_opts, user_oid),
{:ok, system} <- Client.get(client_opts, system_oid),
true <- is_integer(user) and is_integer(system) do
{:ok, (user + system) / 1.0}
else
_ -> {:error, :failed_to_calculate}
end
end
defp build_processor_reading_entry(processor, {:ok, load_percent}, timestamp) do
%{
processor_id: processor.id,
load_percent: load_percent,
status: "ok",
checked_at: timestamp
}
end
defp build_processor_reading_entry(processor, {:error, reason}, timestamp) do
Logger.debug("Failed to poll processor #{processor.description}: #{inspect(reason)}")
%{
processor_id: processor.id,
load_percent: nil,
status: "error",
checked_at: timestamp
}
end
defp handle_processor_poll_metadata(processor, {:ok, load_percent}, timestamp) do
Snmp.update_processor(processor, %{
load_percent: load_percent,
last_checked_at: timestamp
})
end
defp handle_processor_poll_metadata(_processor, {:error, _reason}, _timestamp), do: :ok
defp poll_mempools(mempools, client_opts, timestamp) do
poll_results =
mempools
|> Task.async_stream(
fn mempool ->
result = poll_mempool_value(mempool, client_opts)
{mempool, result}
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, pair} -> [pair]
{:exit, _reason} -> []
end)
# Build batch reading entries (only for successful polls)
reading_entries =
Enum.flat_map(poll_results, fn
{mempool, {:ok, values}} ->
[
%{
mempool_id: mempool.id,
used_bytes: values.used_bytes,
total_bytes: values.total_bytes,
free_bytes: values.free_bytes,
usage_percent: values.usage_percent,
checked_at: timestamp
}
]
{_mempool, {:error, _reason}} ->
[]
end)
_ = Snmp.create_mempool_readings_batch(reading_entries)
# Process metadata updates individually
Enum.each(poll_results, fn
{mempool, {:ok, values}} ->
Snmp.update_mempool(mempool, %{
used_bytes: values.used_bytes,
total_bytes: values.total_bytes,
free_bytes: values.free_bytes,
usage_percent: values.usage_percent,
last_checked_at: timestamp
})
{mempool, {:error, reason}} ->
Logger.debug("Failed to poll mempool #{mempool.description}: #{inspect(reason)}")
end)
end
defp poll_mempool_value(mempool, client_opts) do
used_oid = "1.3.6.1.2.1.25.2.3.1.6.#{mempool.mempool_index}"
size_oid = "1.3.6.1.2.1.25.2.3.1.5.#{mempool.mempool_index}"
alloc_oid = "1.3.6.1.2.1.25.2.3.1.4.#{mempool.mempool_index}"
with {:ok, used_units} <- Client.get(client_opts, used_oid),
{:ok, size_units} <- Client.get(client_opts, size_oid),
{:ok, alloc_units} <- Client.get(client_opts, alloc_oid),
true <- is_integer(used_units) and is_integer(size_units) and is_integer(alloc_units),
true <- size_units > 0 do
used_bytes = used_units * alloc_units
total_bytes = size_units * alloc_units
free_bytes = total_bytes - used_bytes
usage_percent = used_bytes / total_bytes * 100
{:ok,
%{
used_bytes: used_bytes,
total_bytes: total_bytes,
free_bytes: free_bytes,
usage_percent: usage_percent
}}
else
{:error, reason} -> {:error, reason}
false -> {:error, :invalid_values}
end
end
defp poll_sensors(sensors, client_opts, timestamp) do
# Poll all sensors concurrently and collect results
poll_results =
sensors
|> Task.async_stream(
fn sensor ->
result = poll_sensor_value(sensor, client_opts)
{sensor, result}
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, pair} ->
[pair]
{:exit, reason} ->
Logger.warning("Sensor poll task failed: #{inspect(reason)}")
[]
end)
# Build batch reading entries
reading_entries =
Enum.map(poll_results, fn {sensor, result} ->
build_sensor_reading_entry(sensor, result, timestamp)
end)
# Batch insert all sensor readings
{_count, _} = Snmp.create_sensor_readings_batch(reading_entries)
# Process metadata updates and change detection individually
Enum.each(poll_results, fn {sensor, result} ->
handle_sensor_poll_metadata(sensor, result, timestamp)
end)
end
defp poll_state_sensors(state_sensors, client_opts, timestamp, device_id, org_id) do
state_sensors
|> Task.async_stream(
fn state_sensor ->
result = poll_state_sensor_value(state_sensor, client_opts)
handle_state_sensor_poll_result(state_sensor, result, timestamp, device_id, org_id)
state_sensor.sensor_descr
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Enum.each(fn
{:ok, _} -> :ok
{:exit, reason} -> Logger.warning("State sensor poll task failed: #{inspect(reason)}")
end)
end
defp poll_state_sensor_value(state_sensor, client_opts) do
case Client.get(client_opts, state_sensor.sensor_oid) do
{:ok, raw_value} when is_integer(raw_value) ->
{:ok, raw_value}
{:ok, _non_integer} ->
{:error, :non_integer}
{:error, reason} ->
{:error, reason}
end
end
defp handle_state_sensor_poll_result(state_sensor, {:ok, new_state_value}, timestamp, device_id, org_id) do
old_state_value = state_sensor.state_value
old_state_descr = state_sensor.state_descr
old_status = state_sensor.status
new_status = entity_state_to_status(new_state_value)
new_descr = entity_state_to_descr(new_state_value)
Snmp.update_state_sensor(state_sensor, %{
state_value: new_state_value,
state_descr: new_descr,
status: new_status,
last_checked_at: timestamp
})
# Track any state value change, not just status changes
# This catches changes like "6X (64QAM)" -> "8X (256QAM)" that have the same status
if old_state_value != nil && old_state_value != new_state_value do
state_changes = %{
old_state_value: old_state_value,
new_state_value: new_state_value,
old_state_descr: old_state_descr,
new_state_descr: new_descr,
old_status: old_status,
new_status: new_status
}
broadcast_state_sensor_change(state_sensor, state_changes, device_id, org_id, timestamp)
end
end
defp handle_state_sensor_poll_result(state_sensor, {:error, reason}, timestamp, _device_id, _org_id) do
Logger.debug("Failed to poll state sensor #{state_sensor.sensor_descr}: #{inspect(reason)}")
Snmp.update_state_sensor(state_sensor, %{
status: "unknown",
last_checked_at: timestamp
})
end
defp broadcast_state_sensor_change(state_sensor, state_changes, device_id, org_id, timestamp) do
%{
old_state_value: old_state_value,
new_state_value: new_state_value,
old_state_descr: old_state_descr,
new_state_descr: new_state_descr,
old_status: old_status,
new_status: new_status
} = state_changes
severity = determine_state_change_severity(old_status, new_status)
event_type = determine_state_change_event_type(new_status)
# Format message like: "RX Modulation Rate changed from 6X (64QAM)(6) to 8X (256QAM)(8)"
message =
"#{state_sensor.sensor_descr} changed from #{old_state_descr}(#{old_state_value}) to #{new_state_descr}(#{new_state_value})"
event = %{
device_id: device_id,
org_id: org_id,
event_type: event_type,
severity: severity,
message: message,
metadata: %{
state_sensor_id: state_sensor.id,
sensor_name: state_sensor.sensor_descr,
entity_type: state_sensor.entity_type,
old_state_value: old_state_value,
new_state_value: new_state_value,
old_state_descr: old_state_descr,
new_state_descr: new_state_descr,
old_status: old_status,
new_status: new_status
},
occurred_at: timestamp
}
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:#{device_id}", {:device_event, event})
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:events", {:device_event, event})
_ =
if org_id do
_ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device_events:org:#{org_id}", {:device_event, event})
end
Logger.info("State sensor change: #{event.message}")
end
defp determine_state_change_severity(_old_status, "critical"), do: "critical"
defp determine_state_change_severity(_old_status, "warning"), do: "warning"
defp determine_state_change_severity("critical", "ok"), do: "info"
defp determine_state_change_severity("warning", "ok"), do: "info"
defp determine_state_change_severity(_old_status, _new_status), do: "info"
defp determine_state_change_event_type("critical"), do: "state_sensor_critical"
defp determine_state_change_event_type("warning"), do: "state_sensor_warning"
defp determine_state_change_event_type("ok"), do: "state_sensor_normal"
defp determine_state_change_event_type(_), do: "state_sensor_unknown"
defp entity_state_to_status(1), do: "unknown"
defp entity_state_to_status(2), do: "ok"
defp entity_state_to_status(3), do: "warning"
defp entity_state_to_status(4), do: "unknown"
defp entity_state_to_status(_), do: "unknown"
defp entity_state_to_descr(1), do: "unknown"
defp entity_state_to_descr(2), do: "enabled"
defp entity_state_to_descr(3), do: "disabled"
defp entity_state_to_descr(4), do: "testing"
defp entity_state_to_descr(_), do: "unknown"
defp poll_sensor_value(sensor, client_opts) do
if sensor.metadata["calculation"] == "percentage" do
poll_percentage_sensor(sensor, client_opts)
else
poll_simple_sensor(sensor, client_opts)
end
end
# Build a reading entry map for batch insert, or nil if not applicable
defp build_sensor_reading_entry(sensor, {:ok, value}, timestamp) do
state_descr = get_state_description(sensor, value)
%{
sensor_id: sensor.id,
value: value,
status: "ok",
state_descr: state_descr,
checked_at: timestamp
}
end
defp build_sensor_reading_entry(sensor, {:error, :non_numeric}, timestamp) do
Logger.debug("Non-numeric SNMP value for sensor #{sensor.sensor_descr}")
%{
sensor_id: sensor.id,
value: nil,
status: "unknown",
checked_at: timestamp
}
end
defp build_sensor_reading_entry(sensor, {:error, reason}, timestamp) do
log_sensor_error(sensor, reason)
%{
sensor_id: sensor.id,
value: nil,
status: "error",
checked_at: timestamp
}
end
# Process metadata updates and change detection (must remain per-sensor)
defp handle_sensor_poll_metadata(sensor, {:ok, value}, timestamp) do
SensorChangeDetector.detect_and_broadcast(sensor, value, timestamp)
Snmp.update_sensor(sensor, %{
last_value: value,
last_checked_at: timestamp
})
end
defp handle_sensor_poll_metadata(_sensor, {:error, _reason}, _timestamp), do: :ok
# Get state description for state sensors from metadata
defp get_state_description(sensor, value) do
if sensor.sensor_type == "state" && is_map(sensor.metadata) do
states = Map.get(sensor.metadata, "states", %{})
# Convert value to string for lookup (metadata keys are strings)
Map.get(states, to_string(trunc(value)), nil)
end
end
defp log_sensor_error(sensor, :no_such_object) do
Logger.debug("Sensor #{sensor.sensor_descr} OID not found on device")
end
defp log_sensor_error(sensor, :no_such_instance) do
Logger.debug("Sensor #{sensor.sensor_descr} instance not found on device")
end
defp log_sensor_error(sensor, :end_of_mib_view) do
Logger.debug("Sensor #{sensor.sensor_descr} end of MIB view")
end
defp log_sensor_error(sensor, reason) do
Logger.warning("Failed to poll sensor #{sensor.sensor_descr}: #{inspect(reason)}")
end
defp poll_simple_sensor(sensor, client_opts) do
case Client.get(client_opts, sensor.sensor_oid) do
{:ok, raw_value} ->
decoded_value = decode_snmp_value(raw_value)
if is_number(decoded_value) do
divided = decoded_value / sensor.sensor_divisor
normalized = SensorScale.normalize(sensor.sensor_type, divided)
{:ok, Float.round(normalized / 1.0, 2)}
else
{:error, :non_numeric}
end
{:error, reason} ->
{:error, reason}
end
end
defp poll_percentage_sensor(sensor, client_opts) do
size_oid = sensor.metadata["size_oid"]
used_oid = sensor.sensor_oid
with {:ok, used_raw} <- Client.get(client_opts, used_oid),
{:ok, size_raw} <- Client.get(client_opts, size_oid),
used when is_number(used) <- decode_snmp_value(used_raw),
size when is_number(size) <- decode_snmp_value(size_raw),
true <- size > 0 do
percentage = used / size * 100
capped_percentage =
if percentage > 100 do
Logger.warning(
"Percentage sensor #{sensor.sensor_descr} (#{sensor.sensor_type}) " <>
"exceeded 100%: #{Float.round(percentage, 2)}% " <>
"(used=#{used} from #{used_oid}, size=#{size} from #{size_oid}). " <>
"Capping at 100%. This may indicate a device firmware bug."
)
100.0
else
percentage
end
{:ok, capped_percentage}
else
{:error, reason} -> {:error, reason}
false -> {:error, :division_by_zero}
_ -> {:error, :non_numeric}
end
end
defp poll_interfaces(interfaces, client_opts, timestamp, snmp_device) do
# Detect if this is an AirFiber device that needs proprietary counter OIDs
af_overrides = detect_airfiber_counter_overrides(snmp_device, client_opts)
entries =
interfaces
|> Task.async_stream(
fn interface ->
stat_data =
case get_airfiber_stats(af_overrides, interface, client_opts) do
{:ok, data} ->
data
_ ->
{:ok, data} = get_interface_stats(client_opts, interface.if_index)
data
end
Map.merge(stat_data, %{
interface_id: interface.id,
checked_at: timestamp
})
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, entry} -> [entry]
{:exit, _reason} -> []
end)
Snmp.create_interface_stats_batch(entries)
end
# Detect AirFiber devices and return the appropriate proprietary OIDs
# for traffic counters. LibreNMS uses UBNT-AirFIBER-MIB (rxOctetsOK/txOctetsOK)
# and UBNT-AFLTU-MIB (afLTUethRxBytes/afLTUethTxBytes) for these devices
# because their IF-MIB counters are unreliable.
defp detect_airfiber_counter_overrides(nil, _client_opts), do: nil
defp detect_airfiber_counter_overrides(snmp_device, client_opts) do
sys_oid = snmp_device.sys_object_id || ""
# Ubiquiti uses two enterprise OIDs:
# - 1.3.6.1.4.1.41112 (new UBNT enterprise OID)
# - 1.3.6.1.4.1.10002 (old UBNT enterprise OID, used by AirFiber AF11/AF24)
is_ubnt =
String.starts_with?(sys_oid, "1.3.6.1.4.1.41112") or
String.starts_with?(sys_oid, "1.3.6.1.4.1.10002")
if is_ubnt do
# Probe for the AirFiber statistics table — if it exists, this device
# needs proprietary counters. Use v1 since these devices may not support v2c.
v1_opts = Keyword.put(client_opts, :version, "1")
# Check for LTU first (more specific), fall back to regular AirFiber
with {:ltu, {:error, _}} <- {:ltu, Client.get(v1_opts, "1.3.6.1.4.1.41112.1.10.1.2.2.0")},
{:ok, val} when val != nil <- Client.get(v1_opts, "1.3.6.1.4.1.41112.1.3.3.1.1.1") do
:airfiber
else
{:ltu, {:ok, val}} when val != nil -> :airfiber_ltu
_ -> nil
end
end
end
# Fetch traffic stats using AirFiber proprietary MIBs (matching LibreNMS behavior).
# LibreNMS applies these counters to the eth0 interface specifically, but TowerOps
# should apply them to whichever interface is being tracked (eth0, br0, or air0)
# since the proprietary counters represent the radio link's actual throughput.
# Only skip lo and sit0 (loopback/tunnel with no real traffic).
defp get_airfiber_stats(nil, _interface, _client_opts), do: :not_airfiber
defp get_airfiber_stats(_af_type, %{if_descr: descr}, _client_opts) when descr in ["lo", "sit0"], do: :skip_loopback
defp get_airfiber_stats(:airfiber, _interface, client_opts) do
# UBNT-AirFIBER-MIB: airFiberStatistics table (1.3.6.1.4.1.41112.1.3.3.1)
# Field order from MIB SEQUENCE: 1=index, 2=txFramesOK, 3=rxFramesOK,
# 4=rxFrameCrcErr, 5=rxAlignErr, 6=txOctetsOK, 7=rxOctetsOK,
# 8=txPauseFrames, 9=rxPauseFrames, 10=rxErroredFrames, 11=txErroredFrames
#
# These devices may only support SNMPv1 — use v1 for all queries.
v1_opts = Keyword.put(client_opts, :version, "1")
oids = [
# rxOctetsOK
{"1.3.6.1.4.1.41112.1.3.3.1.7.1", :if_in_octets},
# txOctetsOK
{"1.3.6.1.4.1.41112.1.3.3.1.6.1", :if_out_octets},
# rxErroredFrames
{"1.3.6.1.4.1.41112.1.3.3.1.10.1", :if_in_errors},
# txErroredFrames
{"1.3.6.1.4.1.41112.1.3.3.1.11.1", :if_out_errors}
]
fetch_proprietary_stats(v1_opts, oids, true)
end
defp get_airfiber_stats(:airfiber_ltu, _interface, client_opts) do
# UBNT-AFLTU-MIB: afLTUeth table (1.3.6.1.4.1.41112.1.10.1.6)
oids = [
# afLTUethRxBytes
{"1.3.6.1.4.1.41112.1.10.1.6.1.6.0", :if_in_octets},
# afLTUethTxBytes
{"1.3.6.1.4.1.41112.1.10.1.6.1.4.0", :if_out_octets}
]
fetch_proprietary_stats(client_opts, oids, true)
end
defp fetch_proprietary_stats(client_opts, oids, is_hc) do
results =
Enum.map(oids, fn {oid, key} ->
case Client.get(client_opts, oid) do
{:ok, value} -> {key, decode_snmp_value(value)}
_ -> {key, nil}
end
end)
{:ok, Map.new(results ++ [is_hc: is_hc])}
end
defp check_interface_changes(interfaces, device, client_opts, timestamp) do
interfaces
|> Task.async_stream(
fn interface ->
case get_interface_attributes(client_opts, interface.if_index) do
{:ok, current_attrs} ->
detect_and_log_changes(interface, current_attrs, device.id, device.organization_id, timestamp)
{:error, reason} ->
Logger.debug("Failed to get interface attributes for #{interface.if_name}: #{inspect(reason)}")
end
end,
max_concurrency: 2,
timeout: 40_000,
on_timeout: :kill_task
)
|> Stream.run()
end
defp get_interface_attributes(client_opts, if_index) do
oids = [
"1.3.6.1.2.1.2.2.1.5.#{if_index}",
"1.3.6.1.2.1.2.2.1.6.#{if_index}",
"1.3.6.1.2.1.2.2.1.7.#{if_index}",
"1.3.6.1.2.1.2.2.1.8.#{if_index}"
]
case Client.get_multiple(client_opts, oids) do
{:ok, [speed, phys_addr, admin_status, oper_status]} ->
{:ok,
%{
if_speed: parse_interface_integer(speed),
if_phys_address: format_mac_address(phys_addr),
if_admin_status: parse_if_status(admin_status),
if_oper_status: parse_if_status(oper_status)
}}
{:error, reason} = error ->
Logger.debug("Failed to get interface attributes: #{inspect(reason)}")
error
end
end
defp detect_and_log_changes(interface, current_attrs, device_id, org_id, timestamp) do
events =
[]
|> maybe_add_oper_status_event(interface, current_attrs, device_id, timestamp)
|> maybe_add_admin_status_event(interface, current_attrs, device_id, timestamp)
|> maybe_add_speed_change_event(interface, current_attrs, device_id, timestamp)
|> maybe_add_mac_change_event(interface, current_attrs, device_id, timestamp)
if events != [] do
broadcast_interface_events(events, org_id)
Snmp.update_interface(interface, current_attrs)
end
end
defp maybe_add_oper_status_event(events, interface, current_attrs, device_id, timestamp) do
if interface.if_oper_status == current_attrs.if_oper_status do
events
else
event = build_oper_status_event(interface, current_attrs, device_id, timestamp)
[{:event, event} | events]
end
end
defp build_oper_status_event(interface, current_attrs, device_id, timestamp) do
# Handle nil status values from SNMP failures. current_attrs.if_oper_status
# is always set (parse_if_status/1 covers all inputs) — only the stored
# interface field can be nil before the first successful poll.
old_status = interface.if_oper_status || "unknown"
new_status = current_attrs.if_oper_status
message =
if interface.if_oper_status do
"Interface #{interface.if_name} changed from #{String.upcase(old_status)} to #{String.upcase(new_status)}"
else
"Interface #{interface.if_name} is now #{String.upcase(new_status)}"
end
%{
device_id: device_id,
event_type: if(current_attrs.if_oper_status == "up", do: "interface_up", else: "interface_down"),
severity: if(current_attrs.if_oper_status == "up", do: "info", else: "warning"),
message: message,
metadata: %{
interface_id: interface.id,
interface_name: interface.if_name,
old_status: interface.if_oper_status,
new_status: current_attrs.if_oper_status
},
occurred_at: timestamp
}
end
defp maybe_add_admin_status_event(events, interface, current_attrs, device_id, timestamp) do
if interface.if_admin_status == current_attrs.if_admin_status do
events
else
event = build_admin_status_event(interface, current_attrs, device_id, timestamp)
[{:event, event} | events]
end
end
defp build_admin_status_event(interface, current_attrs, device_id, timestamp) do
%{
device_id: device_id,
event_type: "interface_admin_status_change",
severity: "info",
message:
"Interface #{interface.if_name} admin status changed from #{interface.if_admin_status} to #{current_attrs.if_admin_status}",
metadata: %{
interface_id: interface.id,
interface_name: interface.if_name,
old_status: interface.if_admin_status,
new_status: current_attrs.if_admin_status
},
occurred_at: timestamp
}
end
defp maybe_add_speed_change_event(events, interface, current_attrs, device_id, timestamp) do
if interface.if_speed != current_attrs.if_speed && current_attrs.if_speed != nil do
event = build_speed_change_event(interface, current_attrs, device_id, timestamp)
[{:event, event} | events]
else
events
end
end
defp build_speed_change_event(interface, current_attrs, device_id, timestamp) do
is_initial = interface.if_speed == nil
severity = determine_speed_change_severity(is_initial, interface.if_speed, current_attrs.if_speed)
message = format_speed_change_message(interface.if_name, is_initial, interface.if_speed, current_attrs.if_speed)
%{
device_id: device_id,
event_type: "interface_speed_change",
severity: severity,
message: message,
metadata: %{
interface_id: interface.id,
interface_name: interface.if_name,
old_speed: interface.if_speed,
new_speed: current_attrs.if_speed
},
occurred_at: timestamp
}
end
defp determine_speed_change_severity(true, _old_speed, _new_speed), do: "info"
defp determine_speed_change_severity(false, old_speed, new_speed) when new_speed < old_speed, do: "warning"
defp determine_speed_change_severity(false, _old_speed, _new_speed), do: "info"
@doc false
def format_speed_change_message(if_name, true, _old_speed, new_speed) do
"Interface #{if_name} speed detected: #{format_speed(new_speed)}"
end
def format_speed_change_message(if_name, false, old_speed, new_speed) do
"Interface #{if_name} speed changed from #{format_speed(old_speed)} to #{format_speed(new_speed)}"
end
defp maybe_add_mac_change_event(events, interface, current_attrs, device_id, timestamp) do
if interface.if_phys_address != current_attrs.if_phys_address && current_attrs.if_phys_address != nil do
event = build_mac_change_event(interface, current_attrs, device_id, timestamp)
[{:event, event} | events]
else
events
end
end
defp build_mac_change_event(interface, current_attrs, device_id, timestamp) do
is_initial = interface.if_phys_address == nil
message =
format_mac_change_message(interface.if_name, is_initial, interface.if_phys_address, current_attrs.if_phys_address)
%{
device_id: device_id,
event_type: "interface_mac_change",
severity: if(is_initial, do: "info", else: "warning"),
message: message,
metadata: %{
interface_id: interface.id,
interface_name: interface.if_name,
old_mac: interface.if_phys_address,
new_mac: current_attrs.if_phys_address
},
occurred_at: timestamp
}
end
@doc false
def format_mac_change_message(if_name, true, _old_mac, new_mac) do
"Interface #{if_name} MAC address detected: #{new_mac}"
end
def format_mac_change_message(if_name, false, old_mac, new_mac) do
"Interface #{if_name} MAC address changed from #{old_mac} to #{new_mac}"
end
defp broadcast_interface_events(events, org_id) do
Enum.each(events, fn {:event, event_attrs} ->
event_with_org = Map.put(event_attrs, :org_id, org_id)
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{event_attrs.device_id}",
{:device_event, event_with_org}
)
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:events",
{:device_event, event_with_org}
)
_ =
if org_id do
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device_events:org:#{org_id}",
{:device_event, event_with_org}
)
end
Logger.debug("Broadcast event: #{event_attrs.message}")
end)
end
@doc false
def parse_interface_integer(value) when is_integer(value), do: value
def parse_interface_integer(_), do: nil
@doc false
def parse_if_status(1), do: "up"
def parse_if_status(2), do: "down"
def parse_if_status(3), do: "testing"
def parse_if_status(_), do: "unknown"
@doc false
def format_mac_address(<<>>) when is_binary(<<>>), do: nil
def format_mac_address(nil), do: nil
def format_mac_address(mac) when is_binary(mac) do
mac
|> :binary.bin_to_list()
|> Enum.map_join(":", &String.pad_leading(Integer.to_string(&1, 16), 2, "0"))
|> String.downcase()
end
def format_mac_address(_), do: nil
@doc false
def format_speed(speed_bps) when is_integer(speed_bps) and speed_bps >= 1_000_000_000,
do: "#{Float.round(speed_bps / 1_000_000_000, 1)} Gbps"
def format_speed(speed_bps) when is_integer(speed_bps) and speed_bps >= 1_000_000,
do: "#{Float.round(speed_bps / 1_000_000, 1)} Mbps"
def format_speed(speed_bps) when is_integer(speed_bps) and speed_bps >= 1_000,
do: "#{Float.round(speed_bps / 1_000, 1)} Kbps"
def format_speed(speed_bps) when is_integer(speed_bps), do: "#{speed_bps} bps"
def format_speed(_), do: "Unknown"
defp get_interface_stats(client_opts, if_index) do
hc_oids = [
if_in_octets: "1.3.6.1.2.1.31.1.1.1.6.#{if_index}",
if_out_octets: "1.3.6.1.2.1.31.1.1.1.10.#{if_index}"
]
std_oids = [
if_in_octets: "1.3.6.1.2.1.2.2.1.10.#{if_index}",
if_out_octets: "1.3.6.1.2.1.2.2.1.16.#{if_index}"
]
error_oids = [
if_in_errors: "1.3.6.1.2.1.2.2.1.14.#{if_index}",
if_out_errors: "1.3.6.1.2.1.2.2.1.20.#{if_index}",
if_in_discards: "1.3.6.1.2.1.2.2.1.13.#{if_index}",
if_out_discards: "1.3.6.1.2.1.2.2.1.19.#{if_index}"
]
{octet_results, is_hc} = fetch_octet_counters(client_opts, hc_oids, std_oids)
error_results =
Enum.map(error_oids, fn {key, oid} ->
case Client.get(client_opts, oid) do
{:ok, value} -> {key, decode_snmp_value(value)}
_ -> {key, nil}
end
end)
{:ok, Map.new(octet_results ++ error_results ++ [is_hc: is_hc])}
end
defp fetch_octet_counters(client_opts, hc_oids, std_oids) do
# Counter64 (HC) can't be encoded in SNMPv1 responses, so upgrade to v2c
# for HC queries. Most devices support v2c with the same community string.
hc_opts = upgrade_to_v2c_for_hc(client_opts)
hc_results = Enum.map(hc_oids, &try_fetch_hc_counter(hc_opts, &1))
hc_available = Enum.all?(hc_results, fn {_key, value} -> value != :not_available end)
if hc_available do
{Enum.map(hc_results, fn {key, value} -> {key, value} end), true}
else
Logger.debug("HC counters not available, using standard 32-bit counters")
{Enum.map(std_oids, &fetch_standard_counter(client_opts, &1)), false}
end
end
defp upgrade_to_v2c_for_hc(client_opts) do
case Keyword.get(client_opts, :version) do
v when v in ["1", :v1] -> Keyword.put(client_opts, :version, "2c")
_ -> client_opts
end
end
defp try_fetch_hc_counter(client_opts, {key, oid}) do
case Client.get(client_opts, oid) do
{:ok, value} ->
decoded = decode_snmp_value(value)
if decoded == nil, do: {key, :not_available}, else: {key, decoded}
_ ->
{key, :not_available}
end
end
defp fetch_standard_counter(client_opts, {key, oid}) do
case Client.get(client_opts, oid) do
{:ok, value} -> {key, decode_snmp_value(value)}
_ -> {key, nil}
end
end
@doc false
def decode_snmp_value(value) when is_number(value) do
value
end
def decode_snmp_value(value) when is_binary(value) do
size = byte_size(value)
result =
case size do
4 ->
<<counter::unsigned-big-integer-size(32)>> = value
counter
8 ->
<<counter::unsigned-big-integer-size(64)>> = value
counter
16 ->
<<_prefix::binary-size(8), counter::unsigned-big-integer-size(64)>> = value
counter
s when s > 8 ->
offset = s - 8
<<_prefix::binary-size(^offset), counter::unsigned-big-integer-size(64)>> = value
counter
_ ->
Logger.warning("Unknown SNMP binary value format, size: #{size}")
nil
end
result
rescue
error ->
Logger.error("Failed to decode SNMP binary value (size: #{byte_size(value)}): #{inspect(error)}")
nil
end
def decode_snmp_value(other) do
Logger.warning("Unexpected SNMP value type: #{inspect(other)}")
nil
end
defp build_client_opts(device) do
snmp_config = Devices.get_snmp_config(device)
base_opts = [
ip: device.ip_address,
version: snmp_config.version,
port: device.snmp_port || 161
]
# Add version-specific credentials
if snmp_config.version == "3" do
v3_config = Devices.get_snmpv3_config(device)
base_opts ++
[
security_name: v3_config.username,
security_level: v3_config.security_level,
auth_protocol: v3_config.auth_protocol,
auth_password: v3_config.auth_password,
priv_protocol: v3_config.priv_protocol,
priv_password: v3_config.priv_password
]
else
base_opts ++ [community: snmp_config.community]
end
end
@doc """
Polls DOM (Digital Optical Monitoring) metrics for transceivers.
Only polls transceivers with supports_dom=true. For each transceiver, attempts to
read optical power, temperature, voltage, and bias current from vendor-specific MIBs.
"""
def poll_transceivers(transceivers, client_opts, timestamp) do
# Filter to only DOM-capable transceivers
dom_transceivers = Enum.filter(transceivers, & &1.supports_dom)
if dom_transceivers == [] do
:ok
else
poll_transceiver_dom_metrics(dom_transceivers, client_opts, timestamp)
end
end
defp poll_transceiver_dom_metrics(transceivers, client_opts, timestamp) do
# Truncate timestamp to seconds for Ecto :utc_datetime compatibility
truncated_timestamp = DateTime.truncate(timestamp, :second)
# Poll each transceiver's DOM metrics
poll_results =
transceivers
|> Task.async_stream(
fn transceiver ->
result = poll_transceiver_dom(transceiver, client_opts)
{transceiver, result}
end,
max_concurrency: 2,
timeout: 10_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, pair} -> [pair]
{:exit, _reason} -> []
end)
# Build reading entries
reading_entries =
Enum.map(poll_results, fn {transceiver, metrics} ->
build_transceiver_reading_entry(transceiver, metrics, truncated_timestamp)
end)
# Insert readings (even if some metrics are nil)
_ =
case reading_entries do
[] ->
:ok
entries ->
Snmp.create_transceiver_readings_batch(entries)
end
:ok
end
# Poll DOM metrics for a single transceiver
# Returns map with rx_power_dbm, tx_power_dbm, bias_current_ma, temperature_celsius, voltage_v
defp poll_transceiver_dom(transceiver, client_opts) do
# For now, use simulated vendor-specific OIDs (CISCO-ENTITY-SENSOR-MIB style)
# In production, this would dispatch to vendor-specific profiles
base_oid = "1.3.6.1.4.1.9.9.999.1.1.1.1"
port_index = transceiver.port_index
rx_power = poll_transceiver_metric(client_opts, "#{base_oid}.1.#{port_index}")
tx_power = poll_transceiver_metric(client_opts, "#{base_oid}.2.#{port_index}")
bias_current = poll_transceiver_metric(client_opts, "#{base_oid}.3.#{port_index}")
temperature = poll_transceiver_metric(client_opts, "#{base_oid}.4.#{port_index}")
voltage = poll_transceiver_metric(client_opts, "#{base_oid}.5.#{port_index}")
%{
rx_power_dbm: rx_power,
tx_power_dbm: tx_power,
bias_current_ma: bias_current,
temperature_celsius: temperature,
voltage_v: voltage
}
end
defp poll_transceiver_metric(client_opts, oid) do
case Client.get(client_opts, oid) do
{:ok, value} when is_binary(value) ->
case Float.parse(value) do
{float_val, _} -> float_val
:error -> nil
end
{:ok, value} when is_integer(value) ->
value / 1.0
{:ok, value} when is_float(value) ->
value
{:error, _} ->
nil
end
end
defp build_transceiver_reading_entry(transceiver, metrics, timestamp) do
%{
transceiver_id: transceiver.id,
rx_power_dbm: metrics.rx_power_dbm,
tx_power_dbm: metrics.tx_power_dbm,
bias_current_ma: metrics.bias_current_ma,
temperature_celsius: metrics.temperature_celsius,
voltage_v: metrics.voltage_v,
measured_at: timestamp
}
end
@doc """
Polls operational and administrative status for entity physical components.
Monitors status changes for power supplies, fans, modules, and other physical
entities. Records readings when status changes are detected.
"""
def poll_entity_physical_status(entities, client_opts, timestamp) do
# Truncate timestamp to seconds for Ecto :utc_datetime compatibility
truncated_timestamp = DateTime.truncate(timestamp, :second)
# Poll each entity's operational and admin status
poll_results =
entities
|> Task.async_stream(
fn entity ->
result = poll_entity_status(entity, client_opts)
{entity, result}
end,
max_concurrency: 2,
timeout: 10_000,
on_timeout: :kill_task
)
|> Enum.flat_map(fn
{:ok, pair} -> [pair]
{:exit, _reason} -> []
end)
# Build reading entries
reading_entries =
Enum.map(poll_results, fn {entity, status} ->
build_entity_physical_reading_entry(entity, status, truncated_timestamp)
end)
# Insert readings (even if some statuses are nil)
_ =
case reading_entries do
[] ->
:ok
entries ->
Snmp.create_entity_physical_readings_batch(entries)
end
:ok
end
# Poll operational and admin status for a single entity
# Returns map with operational_status and admin_status
defp poll_entity_status(entity, client_opts) do
# ENTITY-MIB (RFC 4133) OIDs for status
# entPhysicalOperStatus: 1.3.6.1.2.1.47.1.1.1.1.5
# entPhysicalAdminStatus: 1.3.6.1.2.1.47.1.1.1.1.6 (not always supported)
oper_status_oid = "1.3.6.1.2.1.47.1.1.1.1.5.#{entity.entity_index}"
admin_status_oid = "1.3.6.1.2.1.47.1.1.1.1.6.#{entity.entity_index}"
oper_status = poll_entity_status_value(client_opts, oper_status_oid)
admin_status = poll_entity_status_value(client_opts, admin_status_oid)
%{
operational_status: oper_status,
admin_status: admin_status
}
end
defp poll_entity_status_value(client_opts, oid) do
case Client.get(client_opts, oid) do
{:ok, value} when is_integer(value) ->
map_status_code(value)
{:ok, value} when is_binary(value) ->
case Integer.parse(value) do
{int_value, _} -> map_status_code(int_value)
:error -> nil
end
_ ->
nil
end
end
# Map RFC 2578 OperStatus enum values
defp map_status_code(1), do: "up"
defp map_status_code(2), do: "down"
defp map_status_code(3), do: "testing"
defp map_status_code(4), do: "unknown"
defp map_status_code(5), do: "dormant"
defp map_status_code(6), do: "notPresent"
defp map_status_code(7), do: "lowerLayerDown"
# Admin status values (1-3 overlap with oper status, plus shuttingDown)
# locked(1), unlocked(2), shuttingDown(3), unknown(4)
defp map_status_code(_), do: "unknown"
defp build_entity_physical_reading_entry(entity, status, timestamp) do
%{
entity_physical_id: entity.id,
operational_status: status.operational_status,
admin_status: status.admin_status,
measured_at: timestamp
}
end
@doc false
def get_poll_interval(device) do
min_interval = Application.get_env(:towerops, :snmp_min_poll_interval, 300)
max(device.check_interval_seconds || @default_poll_interval, min_interval)
end
defp schedule_next_poll(device_id, interval_seconds) do
# Use same offset calculation as start_polling to maintain consistent intervals
# Offset is deterministic hash-based value from 0 to interval_seconds
# This spreads polls evenly across the interval window
offset = PollingOffset.calculate_offset(device_id, interval_seconds)
%{device_id: device_id}
|> new(schedule_in: offset)
|> Oban.insert()
end
@doc false
def should_skip_poll?(device, poll_interval_seconds, grace_period_seconds) do
case device.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
seconds_since_poll < min_interval
end
end
end