Promoted pure presentation and utility helpers from `defp` to `def @doc false` across ~20 LiveViews, Oban workers, and sync modules so they're reachable from unit tests. Refactored several `cond` blocks into idiomatic function heads with guards. Added ~250 new test cases in new files under test/towerops and test/towerops_web, including DB-backed tests for CnMaestro.Sync and AlertNotificationWorker, and removed dead LiveView tab components and CapacityLive (no callers anywhere in lib/test). Configured mix.exs test_coverage.ignore_modules to exclude vendored third-party code (SnmpKit, protobuf-generated Towerops.Agent.*, Absinthe GraphQL types, Phoenix HTML modules, Inspect protocol impls) from coverage calculations — these are not our project code. Coverage: 66.93% → 70.09%. Full suite: 10,127 tests, 0 failures.
1921 lines
61 KiB
Elixir
1921 lines
61 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.Worker,
|
|
queue: :pollers,
|
|
max_attempts: 1,
|
|
unique: [
|
|
period: :infinity,
|
|
keys: [:device_id],
|
|
states: [:available, :scheduled, :retryable]
|
|
],
|
|
replace: [
|
|
scheduled: [:scheduled_at],
|
|
available: [:scheduled_at]
|
|
]
|
|
|
|
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.SensorChangeDetector
|
|
alias Towerops.Snmp.WirelessClientDiscovery
|
|
alias Towerops.Workers.PollingOffset
|
|
|
|
require Logger
|
|
|
|
@default_poll_interval 60
|
|
|
|
@impl Oban.Worker
|
|
@spec perform(Oban.Job.t()) :: :ok
|
|
def perform(%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
|
|
import Ecto.Query
|
|
|
|
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}")
|
|
|
|
_ =
|
|
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
|
|
value = Float.round(decoded_value / sensor.sensor_divisor, 1)
|
|
{:ok, value}
|
|
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
|