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, unique: [ period: 60, keys: [:device_id], states: [:available, :scheduled, :executing, :retryable] ] alias Towerops.Devices alias Towerops.Snmp alias Towerops.Snmp.ArpDiscovery alias Towerops.Snmp.Client alias Towerops.Snmp.MacDiscovery alias Towerops.Snmp.NeighborDiscovery alias Towerops.Workers.PollingOffset require Logger @default_poll_interval 60 @impl Oban.Worker def perform(%Oban.Job{args: %{"device_id" => device_id}}) do case Devices.get_device(device_id) do nil -> Logger.debug("Device #{device_id} no longer exists, skipping poll") :ok device -> if device.snmp_enabled do poll_device(device) end # Schedule next poll poll_interval = get_poll_interval(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 :ok 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 = DateTime.truncate(DateTime.utc_now(), :second) # Run all polling operations in parallel using async tasks 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) ] # Wait for all tasks with timeout, handle failures gracefully tasks |> Task.yield_many(30_000) |> Enum.each(fn {task, result} -> case result do {:ok, _} -> :ok {:exit, reason} -> Logger.error("Polling task failed: #{inspect(reason)}") nil -> Task.shutdown(task, :brutal_kill) 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) 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) 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_neighbors(device.id, cutoff) Enum.each(neighbors, fn neighbor_data -> Snmp.upsert_neighbor(neighbor_data) end) Logger.debug("Polled and saved #{length(neighbors)} neighbors for #{device.name}") _ = Phoenix.PubSub.broadcast( Towerops.PubSub, "device:#{device.id}", {:neighbors_updated, device.id} ) 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_arp_entries(device.id, cutoff) Snmp.upsert_arp_entries(device.id, arp_entries, snmp_device.interfaces) 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} ) 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_mac_addresses(device.id, cutoff) Snmp.upsert_mac_addresses(device.id, mac_entries, snmp_device.interfaces) 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_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 # 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 storage_entries |> Task.async_stream( fn storage -> result = poll_storage_value(storage, client_opts) handle_storage_poll_result(storage, result, timestamp) end, max_concurrency: 2, timeout: 40_000, on_timeout: :kill_task ) |> Stream.run() 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 handle_storage_poll_result(storage, {:ok, values}, timestamp) do Snmp.create_storage_reading(%{ storage_id: storage.id, used_bytes: values.used_bytes, total_bytes: values.total_bytes, usage_percent: values.usage_percent, checked_at: timestamp }) Snmp.update_storage(storage, %{ used_bytes: values.used_bytes, total_bytes: values.total_bytes, last_checked_at: timestamp }) end defp handle_storage_poll_result(storage, {:error, reason}, _timestamp) do Logger.debug("Failed to poll storage #{storage.description}: #{inspect(reason)}") end defp poll_processors(processors, client_opts, timestamp) do processors |> Task.async_stream( fn processor -> result = poll_processor_value(processor, client_opts) handle_processor_poll_result(processor, result, timestamp) end, max_concurrency: 2, timeout: 40_000, on_timeout: :kill_task ) |> Stream.run() end defp poll_processor_value(processor, client_opts) do oid = case processor.processor_type do "hr_processor" -> "1.3.6.1.2.1.25.3.3.1.2.#{processor.processor_index}" "cisco_cpu" -> "1.3.6.1.4.1.9.9.109.1.1.1.1.8.#{processor.processor_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 handle_processor_poll_result(processor, {:ok, load_percent}, timestamp) do Snmp.create_processor_reading(%{ processor_id: processor.id, load_percent: load_percent, status: "ok", checked_at: timestamp }) Snmp.update_processor(processor, %{ load_percent: load_percent, last_checked_at: timestamp }) end defp handle_processor_poll_result(processor, {:error, reason}, timestamp) do Logger.debug("Failed to poll processor #{processor.description}: #{inspect(reason)}") Snmp.create_processor_reading(%{ processor_id: processor.id, load_percent: nil, status: "error", checked_at: timestamp }) end defp poll_sensors(sensors, client_opts, timestamp) do sensors |> Task.async_stream( fn sensor -> result = poll_sensor_value(sensor, client_opts) handle_sensor_poll_result(sensor, result, timestamp) end, max_concurrency: 2, timeout: 40_000, on_timeout: :kill_task ) |> Stream.run() end defp poll_state_sensors(state_sensors, client_opts, timestamp, device_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) end, max_concurrency: 2, timeout: 40_000, on_timeout: :kill_task ) |> Stream.run() 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) 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, timestamp) end end defp handle_state_sensor_poll_result(state_sensor, {:error, reason}, timestamp, _device_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, 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, 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}) 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 defp handle_sensor_poll_result(sensor, {:ok, value}, timestamp) do # For state sensors, map numeric value to state description state_descr = get_state_description(sensor, value) Snmp.create_sensor_reading(%{ sensor_id: sensor.id, value: value, status: "ok", state_descr: state_descr, checked_at: timestamp }) detect_sensor_changes(sensor, value, timestamp) Snmp.update_sensor(sensor, %{ last_value: value, last_checked_at: timestamp }) end defp handle_sensor_poll_result(sensor, {:error, :non_numeric}, timestamp) do Logger.debug("Non-numeric SNMP value for sensor #{sensor.sensor_descr}") Snmp.create_sensor_reading(%{ sensor_id: sensor.id, value: nil, status: "unknown", checked_at: timestamp }) end defp handle_sensor_poll_result(sensor, {:error, reason}, timestamp) do log_sensor_error(sensor, reason) Snmp.create_sensor_reading(%{ sensor_id: sensor.id, value: nil, status: "error", checked_at: timestamp }) end # 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 = decoded_value / sensor.sensor_divisor {: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) do interfaces |> Task.async_stream( fn interface -> {:ok, stat_data} = get_interface_stats(client_opts, interface.if_index) stats = Map.merge(stat_data, %{ interface_id: interface.id, checked_at: timestamp }) Snmp.create_interface_stat(stats) end, max_concurrency: 2, timeout: 40_000, on_timeout: :kill_task ) |> Stream.run() 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, 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, 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) 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 message = if interface.if_oper_status do "Interface #{interface.if_name} changed from #{String.upcase(interface.if_oper_status)} to #{String.upcase(current_attrs.if_oper_status)}" else "Interface #{interface.if_name} is now #{String.upcase(current_attrs.if_oper_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" defp format_speed_change_message(if_name, true, _old_speed, new_speed) do "Interface #{if_name} speed detected: #{format_speed(new_speed)}" end defp 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 defp format_mac_change_message(if_name, true, _old_mac, new_mac) do "Interface #{if_name} MAC address detected: #{new_mac}" end defp 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) do Enum.each(events, fn {:event, event_attrs} -> _ = Phoenix.PubSub.broadcast( Towerops.PubSub, "device:#{event_attrs.device_id}", {:device_event, event_attrs} ) _ = Phoenix.PubSub.broadcast( Towerops.PubSub, "device:events", {:device_event, event_attrs} ) Logger.debug("Broadcast event: #{event_attrs.message}") end) end defp parse_interface_integer(value) when is_integer(value), do: value defp parse_interface_integer(_), do: nil defp parse_if_status(1), do: "up" defp parse_if_status(2), do: "down" defp parse_if_status(3), do: "testing" defp parse_if_status(_), do: "unknown" defp format_mac_address(<<>>) when is_binary(<<>>), do: nil defp format_mac_address(nil), do: nil defp 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 defp format_mac_address(_), do: nil defp format_speed(speed_bps) when is_integer(speed_bps) do cond do speed_bps >= 1_000_000_000 -> "#{Float.round(speed_bps / 1_000_000_000, 1)} Gbps" speed_bps >= 1_000_000 -> "#{Float.round(speed_bps / 1_000_000, 1)} Mbps" speed_bps >= 1_000 -> "#{Float.round(speed_bps / 1_000, 1)} Kbps" true -> "#{speed_bps} bps" end end defp 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 = 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)} end defp fetch_octet_counters(client_opts, hc_oids, std_oids) do hc_results = Enum.map(hc_oids, &try_fetch_hc_counter(client_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) else Logger.debug("HC counters not available, using standard 32-bit counters") Enum.map(std_oids, &fetch_standard_counter(client_opts, &1)) 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 defp decode_snmp_value(value) when is_number(value) do value end defp decode_snmp_value(value) when is_binary(value) do size = byte_size(value) result = case size do 8 -> <> = 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 defp 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) [ ip: device.ip_address, community: snmp_config.community, version: snmp_config.version, port: device.snmp_port || 161 ] end defp 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 offset = PollingOffset.calculate_offset(device_id, interval_seconds) %{device_id: device_id} |> new(schedule_in: interval_seconds + offset) |> Oban.insert() end defp 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 defp detect_sensor_changes(sensor, current_value, timestamp) do if sensor.last_value == nil do :ok else device_id = get_device_id_from_sensor(sensor) thresholds = extract_thresholds(sensor.metadata) events = [] |> maybe_add_threshold_event(sensor, current_value, timestamp, device_id, thresholds) |> maybe_add_change_event(sensor, current_value, timestamp, device_id) broadcast_sensor_events(events) end end defp extract_thresholds(metadata) do %{ warning_high: metadata["warning_high"], critical_high: metadata["critical_high"], warning_low: metadata["warning_low"], critical_low: metadata["critical_low"] } end defp maybe_add_threshold_event(events, sensor, current_value, timestamp, device_id, thresholds) do threshold_event = check_threshold_violation(sensor, current_value, timestamp, device_id, thresholds) if threshold_event, do: [threshold_event | events], else: events end defp check_threshold_violation(sensor, current_value, timestamp, device_id, thresholds) do threshold_checks = [ &check_critical_high/5, &check_critical_low/5, &check_warning_high/5, &check_warning_low/5 ] Enum.find_value(threshold_checks, fn check_fn -> check_fn.(sensor, current_value, timestamp, device_id, thresholds) end) || check_returned_to_normal(sensor, current_value, timestamp, device_id, thresholds) end defp check_critical_high(sensor, current_value, timestamp, device_id, thresholds) do if thresholds.critical_high && current_value >= thresholds.critical_high do build_threshold_event( sensor, current_value, timestamp, device_id, "critical_high", thresholds.critical_high, "critical", "critically high" ) end end defp check_critical_low(sensor, current_value, timestamp, device_id, thresholds) do if thresholds.critical_low && current_value <= thresholds.critical_low do build_threshold_event( sensor, current_value, timestamp, device_id, "critical_low", thresholds.critical_low, "critical", "critically low" ) end end defp check_warning_high(sensor, current_value, timestamp, device_id, thresholds) do if thresholds.warning_high && current_value >= thresholds.warning_high do build_threshold_event( sensor, current_value, timestamp, device_id, "warning_high", thresholds.warning_high, "warning", "high" ) end end defp check_warning_low(sensor, current_value, timestamp, device_id, thresholds) do if thresholds.warning_low && current_value <= thresholds.warning_low do build_threshold_event( sensor, current_value, timestamp, device_id, "warning_low", thresholds.warning_low, "warning", "low" ) end end defp build_threshold_event( sensor, current_value, timestamp, device_id, threshold_type, threshold_value, severity, description ) do event_type = if severity == "critical", do: "sensor_threshold_critical", else: "sensor_threshold_warning" build_sensor_event( device_id, sensor, event_type, severity, "#{sensor.sensor_descr} is #{description}: #{format_sensor_value(current_value, sensor.sensor_unit)} (threshold: #{threshold_value}#{sensor.sensor_unit})", %{ current_value: current_value, previous_value: sensor.last_value, threshold_type: threshold_type, threshold_value: threshold_value }, timestamp ) end defp check_returned_to_normal(sensor, current_value, timestamp, device_id, thresholds) do was_over_threshold = value_was_over_threshold?(sensor.last_value, thresholds) is_now_normal = value_is_normal?(current_value, thresholds) if was_over_threshold && is_now_normal do build_sensor_event( device_id, sensor, "sensor_threshold_normal", "info", "#{sensor.sensor_descr} returned to normal: #{format_sensor_value(current_value, sensor.sensor_unit)}", %{current_value: current_value, previous_value: sensor.last_value}, timestamp ) end end defp value_was_over_threshold?(last_value, thresholds) do (thresholds.warning_high && last_value >= thresholds.warning_high) || (thresholds.critical_high && last_value >= thresholds.critical_high) || (thresholds.warning_low && last_value <= thresholds.warning_low) || (thresholds.critical_low && last_value <= thresholds.critical_low) end defp value_is_normal?(current_value, thresholds) do (thresholds.warning_high == nil || current_value < thresholds.warning_high) && (thresholds.critical_high == nil || current_value < thresholds.critical_high) && (thresholds.warning_low == nil || current_value > thresholds.warning_low) && (thresholds.critical_low == nil || current_value > thresholds.critical_low) end defp maybe_add_change_event(events, sensor, current_value, timestamp, device_id) do if sensor.sensor_unit == "%" do change_event = check_significant_change(sensor, current_value, timestamp, device_id) if change_event, do: [change_event | events], else: events else events end end defp check_significant_change(sensor, current_value, timestamp, device_id) do change_percent = abs(current_value - sensor.last_value) cond do change_percent >= 30 && current_value > sensor.last_value -> build_sensor_event( device_id, sensor, "sensor_value_spike", "warning", "#{sensor.sensor_descr} spiked from #{format_sensor_value(sensor.last_value, sensor.sensor_unit)} to #{format_sensor_value(current_value, sensor.sensor_unit)}", %{current_value: current_value, previous_value: sensor.last_value, change: change_percent}, timestamp ) change_percent >= 30 && current_value < sensor.last_value -> build_sensor_event( device_id, sensor, "sensor_value_drop", "info", "#{sensor.sensor_descr} dropped from #{format_sensor_value(sensor.last_value, sensor.sensor_unit)} to #{format_sensor_value(current_value, sensor.sensor_unit)}", %{current_value: current_value, previous_value: sensor.last_value, change: change_percent}, timestamp ) true -> nil end end defp broadcast_sensor_events(events) do Enum.each(events, fn event -> _ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:#{event.device_id}", {:device_event, event}) _ = Phoenix.PubSub.broadcast(Towerops.PubSub, "device:events", {:device_event, event}) Logger.debug("Sensor event: #{event.message}") end) end defp build_sensor_event(device_id, sensor, event_type, severity, message, metadata, timestamp) do base_metadata = %{ sensor_id: sensor.id, sensor_name: sensor.sensor_descr, sensor_type: sensor.sensor_type, sensor_unit: sensor.sensor_unit } %{ device_id: device_id, event_type: event_type, severity: severity, message: message, metadata: Map.merge(base_metadata, metadata), occurred_at: timestamp } end defp get_device_id_from_sensor(sensor) do case Towerops.Repo.get(Towerops.Snmp.Device, sensor.snmp_device_id) do nil -> nil device -> device.device_id end end defp format_sensor_value(value, unit) when is_number(value) do "#{Float.round(value, 1)}#{unit}" end defp format_sensor_value(_, _), do: "N/A" end