Reliability improvements: agent channel, polling, discovery, batch ops

- Add Task.Supervisor for supervised background tasks in agent channel
- Discovery worker retries transient failures (max 3 attempts)
- Log silent polling timeouts with named task identification
- Track errors in batch ARP/MAC upserts instead of silently ignoring
- Discovery gracefully degrades on interface/sensor timeout (partial results)
- Synchronous job dispatch on agent join to prevent race condition
- Surface SNMP result, monitoring check, and backup processing errors
- Atomic delete-stale + upsert for neighbors/ARP/MAC (transaction)
- Discovery transaction timeout (60s) with duration logging
- Heartbeat timeout detection (5min) closes stale agent channels
- Live poll timeout (60s) notifies waiting LiveView on expiry
- Disconnect agent channel when token is disabled or deleted
- Debounce rapid assignment changes to prevent redundant job dispatches
- Sensor index deduplication via YAML index template substitution
This commit is contained in:
Graham McIntire 2026-02-09 16:28:12 -06:00
parent 3b8e1ea639
commit 4cfab3b1ee
No known key found for this signature in database
9 changed files with 696 additions and 114 deletions

View file

@ -212,9 +212,22 @@ defmodule Towerops.Agents do
"""
@spec update_agent_token(AgentToken.t(), map()) :: {:ok, AgentToken.t()} | {:error, Ecto.Changeset.t()}
def update_agent_token(%AgentToken{} = agent_token, attrs) do
agent_token
|> AgentToken.update_changeset(attrs)
|> Repo.update()
result =
agent_token
|> AgentToken.update_changeset(attrs)
|> Repo.update()
# If token was disabled, disconnect any active agent channel
with {:ok, updated} <- result,
true <- agent_token.enabled == true and updated.enabled == false do
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"agent:#{updated.id}:lifecycle",
:token_disabled
)
end
result
end
@doc """
@ -293,6 +306,13 @@ defmodule Towerops.Agents do
"""
@spec delete_agent_token(Ecto.UUID.t()) :: {:ok, AgentToken.t()} | {:error, any()}
def delete_agent_token(id) do
# Disconnect any active agent channel before deleting
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"agent:#{id}:lifecycle",
:token_disabled
)
Repo.transaction(fn ->
agent_token = Repo.get!(AgentToken, id)

View file

@ -493,6 +493,7 @@ defmodule Towerops.Profiles.YamlProfiles do
sensor_unit: Map.get(sensor_def, "unit"),
sensor_divisor: Map.get(sensor_def, "divisor", 1),
precision: Map.get(sensor_def, "precision", 1),
index_template: Map.get(sensor_def, "index"),
table: true
}
end
@ -593,6 +594,7 @@ defmodule Towerops.Profiles.YamlProfiles do
sensor_unit: Map.get(sensor_def, "unit") || default_unit,
sensor_divisor: Map.get(sensor_def, "divisor", 1),
precision: Map.get(sensor_def, "precision", 1),
index_template: Map.get(sensor_def, "index"),
table: true
}
end
@ -620,6 +622,7 @@ defmodule Towerops.Profiles.YamlProfiles do
sensor_descr: extract_descr(sensor_def, "state"),
sensor_unit: "",
states: state_map,
index_template: Map.get(sensor_def, "index"),
table: true
}
end

View file

@ -27,6 +27,8 @@ defmodule Towerops.Snmp do
alias Towerops.Snmp.StorageReading
alias Towerops.Snmp.Vlan
require Logger
@doc """
Tests SNMP connectivity to a device.
@ -527,6 +529,20 @@ defmodule Towerops.Snmp do
|> Repo.delete_all()
end
@doc """
Atomically deletes stale neighbors and upserts new ones in a single transaction.
Prevents race conditions where concurrent delete/upsert could remove fresh data.
"""
def delete_stale_and_upsert_neighbors(device_id, neighbors, cutoff) do
Repo.transaction(fn ->
delete_stale_neighbors(device_id, cutoff)
Enum.each(neighbors, fn neighbor_data ->
upsert_neighbor(neighbor_data)
end)
end)
end
@doc """
Lists all discovered devices across an organization that haven't been added yet.
@ -1520,14 +1536,26 @@ defmodule Towerops.Snmp do
|> Repo.delete_all()
end
@doc """
Atomically deletes stale ARP entries and upserts new ones in a single transaction.
Prevents race conditions where concurrent delete/upsert could remove fresh data.
"""
def delete_stale_and_upsert_arp_entries(device_id, arp_entries, interfaces, cutoff) do
Repo.transaction(fn ->
delete_stale_arp_entries(device_id, cutoff)
upsert_arp_entries(device_id, arp_entries, interfaces)
end)
end
@doc """
Bulk upserts ARP entries for a device.
Returns the number of entries processed.
Returns `{success_count, error_count}` tuple.
"""
def upsert_arp_entries(device_id, arp_entries, interfaces) do
interface_map = Map.new(interfaces, fn i -> {i.if_index, i.id} end)
Enum.each(arp_entries, fn entry ->
arp_entries
|> Enum.reduce({0, 0}, fn entry, {s, e} ->
interface_id = Map.get(interface_map, entry.if_index)
attrs =
@ -1535,10 +1563,19 @@ defmodule Towerops.Snmp do
|> Map.put(:device_id, device_id)
|> Map.put(:interface_id, interface_id)
upsert_arp_entry(attrs)
case upsert_arp_entry(attrs) do
{:ok, _} -> {s + 1, e}
{:error, _} -> {s, e + 1}
end
end)
|> tap(fn {_s, errors} ->
if errors > 0 do
Logger.warning("#{errors} ARP upsert failures for device #{device_id}",
device_id: device_id,
error_count: errors
)
end
end)
length(arp_entries)
end
# MAC Address queries
@ -1581,9 +1618,20 @@ defmodule Towerops.Snmp do
|> Repo.delete_all()
end
@doc """
Atomically deletes stale MAC addresses and upserts new ones in a single transaction.
Prevents race conditions where concurrent delete/upsert could remove fresh data.
"""
def delete_stale_and_upsert_mac_addresses(device_id, mac_entries, interfaces, cutoff) do
Repo.transaction(fn ->
delete_stale_mac_addresses(device_id, cutoff)
upsert_mac_addresses(device_id, mac_entries, interfaces)
end)
end
@doc """
Bulk upserts MAC address entries for a device.
Returns the number of entries processed.
Returns `{success_count, error_count}` tuple.
"""
def upsert_mac_addresses(device_id, mac_entries, interfaces) do
# Build a map from bridge port index to interface ID
@ -1591,7 +1639,8 @@ defmodule Towerops.Snmp do
# For now, we'll use a simple mapping if they're available
interface_map = Map.new(interfaces, fn i -> {i.if_index, i.id} end)
Enum.each(mac_entries, fn entry ->
mac_entries
|> Enum.reduce({0, 0}, fn entry, {s, e} ->
interface_id = Map.get(interface_map, entry.port_index)
attrs =
@ -1599,10 +1648,19 @@ defmodule Towerops.Snmp do
|> Map.put(:device_id, device_id)
|> Map.put(:interface_id, interface_id)
upsert_mac_address(attrs)
case upsert_mac_address(attrs) do
{:ok, _} -> {s + 1, e}
{:error, _} -> {s, e + 1}
end
end)
|> tap(fn {_s, errors} ->
if errors > 0 do
Logger.warning("#{errors} MAC upsert failures for device #{device_id}",
device_id: device_id,
error_count: errors
)
end
end)
length(mac_entries)
end
# Network Topology Functions

View file

@ -195,6 +195,7 @@ defmodule Towerops.Snmp.Discovery do
Client.test_connection(client_opts)
end
# Critical path: connection, system info, device info must succeed
with {:ok, _} <- connection_result,
Logger.info("Discovering system info...", device_id: device.id),
{:ok, system_info} <- discover_system(client_opts),
@ -203,51 +204,85 @@ defmodule Towerops.Snmp.Discovery do
Logger.info("Selecting device profile...", device_id: device.id),
profile = select_profile(system_info, client_opts),
Logger.info("Profile selected, building device info...", device_id: device.id),
{:ok, device_info} <- build_device_info(client_opts, system_info, profile),
Logger.info("Discovering interfaces...", device_id: device.id),
{:ok, interfaces} <- discover_interfaces_with_timeout(client_opts, profile, timeouts),
Logger.info("Discovering sensors...", device_id: device.id),
{:ok, sensors} <- discover_sensors_with_timeout(client_opts, profile, timeouts),
Logger.info("Discovering VLANs...", device_id: device.id),
{:ok, vlans} <- discover_vlans_with_timeout(client_opts, profile, timeouts),
Logger.info("Discovering IP addresses...", device_id: device.id),
{:ok, ip_addresses} <- discover_ip_addresses_with_timeout(client_opts, timeouts),
_ = log_ip_discovery_results(device, ip_addresses),
Logger.info("Discovering processors...", device_id: device.id),
{:ok, processors} <- discover_processors_with_timeout(client_opts, timeouts),
Logger.info("Discovering storage...", device_id: device.id),
{:ok, storage} <- discover_storage_with_timeout(client_opts, timeouts),
Logger.info("Saving discovery results (including IP addresses and processors)...", device_id: device.id),
{:ok, discovered_device} <-
save_discovery_results(device, device_info, interfaces, sensors, vlans, ip_addresses, processors),
Logger.info("Syncing storage...", device_id: device.id),
:ok <- sync_storage(discovered_device, storage),
Logger.info("Discovering neighbors...", device_id: device.id),
{:ok, neighbors} <- discover_neighbors_with_timeout(client_opts, discovered_device.interfaces, timeouts),
Logger.info("Saving neighbors...", device_id: device.id),
:ok <- save_neighbors(discovered_device.device_id, neighbors),
Logger.info("Discovering ARP entries...", device_id: device.id),
{:ok, arp_entries} <- discover_arp_with_timeout(client_opts, timeouts),
:ok <- save_arp_entries(discovered_device.device_id, arp_entries, discovered_device.interfaces) do
_ = update_device_discovery_time(device)
Logger.debug("SNMP discovery completed successfully for: #{device.name}")
{:ok, device_info} <- build_device_info(client_opts, system_info, profile) do
# Interface and sensor discovery: fall back to empty on timeout with warning
# CRITICAL: Empty list on timeout means sync will NOT delete existing data
# because save_discovery_results skips sync when list is empty and device already exists
Logger.info("Discovering interfaces...", device_id: device.id)
# Determine if this was first discovery or rediscovery
# Check if device existed before (had last_discovery_at set)
is_rediscovery = device.last_discovery_at != nil
interfaces =
case discover_interfaces_with_timeout(client_opts, profile, timeouts) do
{:ok, ifaces} ->
ifaces
# Log discovery event
_ = log_discovery_event(device, discovered_device, is_rediscovery)
{:error, :interface_discovery_timeout} ->
Logger.warning("Interface discovery timed out, using empty list", device_id: device.id)
[]
end
# Broadcast discovery completion for real-time updates
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:discovery_completed, device.id}
)
Logger.info("Discovering sensors...", device_id: device.id)
{:ok, discovered_device}
sensors =
case discover_sensors_with_timeout(client_opts, profile, timeouts) do
{:ok, s} ->
s
{:error, :sensor_discovery_timeout} ->
Logger.warning("Sensor discovery timed out, using empty list", device_id: device.id)
[]
end
# Secondary discovery: VLANs, IPs, processors, storage (already fall back to [] on timeout)
Logger.info("Discovering VLANs...", device_id: device.id)
{:ok, vlans} = discover_vlans_with_timeout(client_opts, profile, timeouts)
Logger.info("Discovering IP addresses...", device_id: device.id)
{:ok, ip_addresses} = discover_ip_addresses_with_timeout(client_opts, timeouts)
_ = log_ip_discovery_results(device, ip_addresses)
Logger.info("Discovering processors...", device_id: device.id)
{:ok, processors} = discover_processors_with_timeout(client_opts, timeouts)
Logger.info("Discovering storage...", device_id: device.id)
{:ok, storage} = discover_storage_with_timeout(client_opts, timeouts)
Logger.info("Saving discovery results (including IP addresses and processors)...", device_id: device.id)
case save_discovery_results(device, device_info, interfaces, sensors, vlans, ip_addresses, processors) do
{:ok, discovered_device} ->
# Post-save discovery: neighbors, ARP (best-effort)
Logger.info("Syncing storage...", device_id: device.id)
_ = sync_storage(discovered_device, storage)
Logger.info("Discovering neighbors...", device_id: device.id)
{:ok, neighbors} = discover_neighbors_with_timeout(client_opts, discovered_device.interfaces, timeouts)
Logger.info("Saving neighbors...", device_id: device.id)
_ = save_neighbors(discovered_device.device_id, neighbors)
Logger.info("Discovering ARP entries...", device_id: device.id)
{:ok, arp_entries} = discover_arp_with_timeout(client_opts, timeouts)
_ = save_arp_entries(discovered_device.device_id, arp_entries, discovered_device.interfaces)
_ = update_device_discovery_time(device)
Logger.debug("SNMP discovery completed successfully for: #{device.name}")
is_rediscovery = device.last_discovery_at != nil
_ = log_discovery_event(device, discovered_device, is_rediscovery)
_ =
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"device:#{device.id}",
{:discovery_completed, device.id}
)
{:ok, discovered_device}
{:error, reason} ->
Logger.error("Failed to save discovery results for #{device.name}: #{inspect(reason)}")
{:error, reason}
end
else
{:error, reason} = error ->
Logger.error("SNMP discovery failed for #{device.name}: #{inspect(reason)}")
@ -598,28 +633,46 @@ defmodule Towerops.Snmp.Discovery do
# This can take several seconds for large discovery data (thousands of OIDs)
sanitized_device_info = prepare_device_info_for_save(device_info, device.id)
Repo.transaction(fn ->
# Upsert Device
snmp_device = upsert_device(device, sanitized_device_info)
start = System.monotonic_time(:millisecond)
# Sync interfaces, sensors, and VLANs (preserving historical data)
synced_interfaces = sync_interfaces(snmp_device, interfaces)
_ = sync_sensors(snmp_device, sensors)
_ = sync_vlans(snmp_device, vlans)
result =
Repo.transaction(
fn ->
# Upsert Device
snmp_device = upsert_device(device, sanitized_device_info)
# Sync IP addresses and processors inside transaction for atomicity
device_with_interfaces = %{
snmp_device
| interfaces: Enum.map(synced_interfaces, &Map.put(&1, :device_id, snmp_device.device_id))
}
# Sync interfaces, sensors, and VLANs (preserving historical data)
synced_interfaces = sync_interfaces(snmp_device, interfaces)
_ = sync_sensors(snmp_device, sensors)
_ = sync_vlans(snmp_device, vlans)
_ = sync_ip_addresses(device_with_interfaces, ip_addresses)
_ = sync_processors(device_with_interfaces, processors)
# Sync IP addresses and processors inside transaction for atomicity
device_with_interfaces = %{
snmp_device
| interfaces: Enum.map(synced_interfaces, &Map.put(&1, :device_id, snmp_device.device_id))
}
# Return device with interfaces loaded and device_id added
# Use snmp_device.device_id which references the Equipment table
device_with_interfaces
end)
_ = sync_ip_addresses(device_with_interfaces, ip_addresses)
_ = sync_processors(device_with_interfaces, processors)
# Return device with interfaces loaded and device_id added
# Use snmp_device.device_id which references the Equipment table
device_with_interfaces
end,
timeout: 60_000
)
duration = System.monotonic_time(:millisecond) - start
if duration > 5_000 do
Logger.warning("Discovery transaction took #{duration}ms",
device_id: device.id,
interface_count: length(interfaces),
sensor_count: length(sensors)
)
end
result
end
# Prepares device_info for database save by sanitizing raw_discovery_data
@ -1080,9 +1133,9 @@ defmodule Towerops.Snmp.Discovery do
Towerops.Snmp.delete_stale_arp_entries(device_id, cutoff)
# Upsert each discovered ARP entry
count = Towerops.Snmp.upsert_arp_entries(device_id, arp_entries, interfaces)
{success, _errors} = Towerops.Snmp.upsert_arp_entries(device_id, arp_entries, interfaces)
Logger.debug("Saved #{count} ARP entries for device #{device_id}")
Logger.debug("Saved #{success} ARP entries for device #{device_id}")
:ok
end

View file

@ -316,9 +316,20 @@ defmodule Towerops.Snmp.Profiles.Dynamic do
# Build a sensor from table walk result
defp build_table_sensor(sensor_def, oid, value, idx, descr_map, oid_index) when is_number(value) do
# Use index_template from YAML if available, otherwise generate default index
sensor_index =
case sensor_def[:index_template] do
nil ->
"#{sensor_def[:sensor_type]}_#{idx}"
template ->
# Replace {{ $index }} with the actual OID index
String.replace(template, "{{ $index }}", oid_index)
end
%{
sensor_type: sensor_def[:sensor_type],
sensor_index: "#{sensor_def[:sensor_type]}_#{idx}",
sensor_index: sensor_index,
sensor_oid: oid,
sensor_descr: build_sensor_descr(sensor_def, idx, descr_map, oid_index),
sensor_unit: sensor_def[:sensor_unit] || "",
@ -350,13 +361,15 @@ defmodule Towerops.Snmp.Profiles.Dynamic do
results
|> Enum.with_index(1)
|> Enum.map(fn {{oid, value}, idx} ->
build_state_sensor(sensor_def, oid, value, idx)
# Extract the index from the OID for template substitution
oid_index = oid |> String.split(".") |> List.last()
build_state_sensor(sensor_def, oid, value, idx, oid_index)
end)
|> Enum.reject(&is_nil/1)
end
# Build a state sensor with value-to-description mapping
defp build_state_sensor(sensor_def, oid, value, idx) when is_number(value) do
defp build_state_sensor(sensor_def, oid, value, idx, oid_index) when is_number(value) do
states = sensor_def[:states] || %{}
# Convert value to integer for state lookup (states are discrete)
int_value = trunc(value)
@ -366,9 +379,20 @@ defmodule Towerops.Snmp.Profiles.Dynamic do
# Convert states map keys to strings for JSON storage
states_for_json = Map.new(states, fn {k, v} -> {to_string(k), v} end)
# Use index_template from YAML if available, otherwise generate default index
sensor_index =
case sensor_def[:index_template] do
nil ->
"#{sensor_def[:oid_name] || "state"}_#{idx}"
template ->
# Replace {{ $index }} with the actual OID index
String.replace(template, "{{ $index }}", oid_index)
end
%{
sensor_type: "state",
sensor_index: "#{sensor_def[:oid_name] || "state"}_#{idx}",
sensor_index: sensor_index,
sensor_oid: oid,
sensor_descr: build_sensor_descr_default(sensor_def, idx),
sensor_unit: "",
@ -379,7 +403,7 @@ defmodule Towerops.Snmp.Profiles.Dynamic do
}
end
defp build_state_sensor(_sensor_def, _oid, _value, _idx), do: nil
defp build_state_sensor(_sensor_def, _oid, _value, _idx, _oid_index), do: nil
# Build sensor description, using the descr_map if available for actual device names
defp build_sensor_descr(sensor_def, idx, descr_map, oid_index) do

View file

@ -313,11 +313,7 @@ defmodule Towerops.Workers.DevicePollerWorker do
{: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)
Snmp.delete_stale_and_upsert_neighbors(device.id, neighbors, cutoff)
Logger.debug("Polled and saved #{length(neighbors)} neighbors for #{device.name}")
@ -336,9 +332,7 @@ defmodule Towerops.Workers.DevicePollerWorker 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)
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}")
@ -357,9 +351,7 @@ defmodule Towerops.Workers.DevicePollerWorker 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)
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}")

View file

@ -44,6 +44,11 @@ defmodule ToweropsWeb.AgentChannel do
@typep base64_string :: String.t()
@typep protobuf_binary :: binary()
# Heartbeat timeout: close channel if no heartbeat in 5 minutes
@heartbeat_timeout_seconds 300
# Check heartbeat every 2 minutes
@heartbeat_check_interval_ms 120_000
@impl true
@spec join(String.t(), map(), Phoenix.Socket.t()) ::
{:ok, Phoenix.Socket.t()} | {:error, map()}
@ -85,6 +90,9 @@ defmodule ToweropsWeb.AgentChannel do
# Subscribe to live poll requests for this agent
_ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:live_poll")
# Subscribe to token lifecycle events (disable/delete triggers disconnect)
_ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:lifecycle")
# Update last_seen_at and IP on join
remote_ip = get_remote_ip(socket)
_ = Agents.update_agent_token_heartbeat(agent_token.id, remote_ip, %{})
@ -97,8 +105,13 @@ defmodule ToweropsWeb.AgentChannel do
{:agent_connected, agent_token.id, agent_token.organization_id}
)
# Send initial job list
send(self(), :send_jobs)
# Track heartbeat time and schedule periodic check
socket = assign(socket, :last_heartbeat_at, DateTime.utc_now())
Process.send_after(self(), :check_heartbeat, @heartbeat_check_interval_ms)
# Send initial job list synchronously during join to avoid race condition
# where agent disconnects before the async :send_jobs message is processed
build_and_push_jobs(socket)
{:ok, socket}
@ -129,24 +142,50 @@ defmodule ToweropsWeb.AgentChannel do
@impl true
@spec handle_info(:send_jobs, Phoenix.Socket.t()) :: {:noreply, Phoenix.Socket.t()}
def handle_info(:send_jobs, socket) do
jobs = build_jobs_for_agent(socket.assigns.agent_token_id)
job_list = %AgentJobList{jobs: jobs}
binary = AgentJobList.encode(job_list)
Logger.info("Sending #{length(jobs)} jobs to agent: #{inspect(Enum.map(jobs, & &1.job_id))}")
push(socket, "jobs", %{binary: Base.encode64(binary)})
build_and_push_jobs(socket)
{:noreply, socket}
end
# Handle PubSub broadcast when device assignments change
def handle_info({:assignments_changed, _event}, socket) do
Logger.info("Agent assignments changed, sending updated job list",
# Periodic heartbeat check — close channel if agent hasn't sent a heartbeat recently
def handle_info(:check_heartbeat, socket) do
seconds_since = DateTime.diff(DateTime.utc_now(), socket.assigns.last_heartbeat_at, :second)
if seconds_since > @heartbeat_timeout_seconds do
Logger.warning("Agent heartbeat timeout after #{seconds_since}s",
agent_token_id: socket.assigns.agent_token_id
)
{:stop, :heartbeat_timeout, socket}
else
Process.send_after(self(), :check_heartbeat, @heartbeat_check_interval_ms)
{:noreply, socket}
end
end
# Handle token disabled/deleted — disconnect the agent
def handle_info(:token_disabled, socket) do
Logger.warning("Agent token disabled, disconnecting",
agent_token_id: socket.assigns.agent_token_id
)
send(self(), :send_jobs)
{:noreply, socket}
{:stop, :token_disabled, socket}
end
# Handle PubSub broadcast when device assignments change.
# Debounces rapid changes — cancels any pending refresh and schedules a new one in 500ms.
def handle_info({:assignments_changed, _event}, socket) do
Logger.info("Agent assignments changed, scheduling job refresh",
agent_token_id: socket.assigns.agent_token_id
)
# Cancel any pending job refresh timer
if timer = socket.assigns[:job_refresh_timer] do
Process.cancel_timer(timer)
end
# Schedule new refresh in 500ms
timer = Process.send_after(self(), :send_jobs, 500)
{:noreply, assign(socket, :job_refresh_timer, timer)}
end
# Handle PubSub broadcast when discovery is requested for a device
@ -251,10 +290,27 @@ defmodule ToweropsWeb.AgentChannel do
)
push(socket, "jobs", %{binary: Base.encode64(binary)})
# Schedule timeout — if no result in 60s, notify the waiting LiveView
Process.send_after(self(), {:live_poll_timeout, reply_topic}, 60_000)
{:noreply, socket}
end
end
# Handle live poll timeout — notify the LiveView that the poll didn't return in time
def handle_info({:live_poll_timeout, reply_topic}, socket) do
Logger.warning("Live poll timed out", reply_topic: reply_topic)
Phoenix.PubSub.broadcast(
Towerops.PubSub,
reply_topic,
{:live_poll_error, :timeout}
)
{:noreply, socket}
end
@impl true
@spec handle_in(String.t(), map(), socket()) :: {:noreply, socket()}
def handle_in("result", %{"binary" => binary_b64}, socket) when is_binary(binary_b64) do
@ -272,7 +328,7 @@ defmodule ToweropsWeb.AgentChannel do
if String.starts_with?(result.job_id, "live_poll:") do
handle_live_poll_result(result, socket)
else
_ = process_snmp_result(socket.assigns.organization_id, result, socket)
process_and_log_snmp_result(socket, result)
end
{:noreply, socket}
@ -315,6 +371,9 @@ defmodule ToweropsWeb.AgentChannel do
metadata
)
# Update heartbeat tracking for timeout detection
socket = assign(socket, :last_heartbeat_at, DateTime.utc_now())
# Broadcast heartbeat for real-time UI updates (especially important for stale agents coming back online)
_ =
Phoenix.PubSub.broadcast(
@ -435,7 +494,16 @@ defmodule ToweropsWeb.AgentChannel do
)
# Store ping result in database
_ = store_monitoring_check(check, socket)
case store_monitoring_check(check, socket) do
:ok ->
:ok
{:error, reason} ->
Logger.error("Failed to store monitoring check",
device_id: check.device_id,
reason: inspect(reason)
)
end
{:noreply, socket}
else
@ -474,7 +542,17 @@ defmodule ToweropsWeb.AgentChannel do
# Check if this is a backup job by the job_id prefix
if String.starts_with?(result.job_id, "backup:") do
_ = process_backup_result(result)
case process_backup_result(result) do
:ok ->
:ok
{:error, reason} ->
Logger.error("Failed to process backup result",
device_id: result.device_id,
job_id: result.job_id,
reason: inspect(reason)
)
end
else
# Handle regular MikroTik polling results if needed
Logger.debug("Received non-backup MikroTik result", job_id: result.job_id)
@ -483,6 +561,21 @@ defmodule ToweropsWeb.AgentChannel do
{:noreply, socket}
end
# Builds all jobs for the agent and pushes them to the channel.
# Used both during join (synchronous) and for subsequent refreshes.
defp build_and_push_jobs(socket) do
jobs = build_jobs_for_agent(socket.assigns.agent_token_id)
job_list = %AgentJobList{jobs: jobs}
binary = AgentJobList.encode(job_list)
Logger.info("Sending #{length(jobs)} jobs to agent",
agent_token_id: socket.assigns.agent_token_id,
job_ids: inspect(Enum.map(jobs, & &1.job_id))
)
push(socket, "jobs", %{binary: Base.encode64(binary)})
end
@spec build_jobs_for_agent(Ecto.UUID.t()) :: [AgentJob.t()]
defp build_jobs_for_agent(agent_token_id) do
agent_token_id
@ -923,23 +1016,42 @@ defmodule ToweropsWeb.AgentChannel do
]
end
defp process_and_log_snmp_result(socket, result) do
case process_snmp_result(socket.assigns.organization_id, result, socket) do
:ok ->
:ok
{:error, reason} ->
Logger.error("Failed to process SNMP result",
device_id: result.device_id,
job_id: result.job_id,
reason: inspect(reason)
)
end
end
defp process_snmp_result(organization_id, result, socket) do
with {:ok, device} <- fetch_device(result.device_id),
:ok <- verify_device_organization(device, organization_id),
:ok <- verify_device_assignment(device, socket.assigns.agent_token_id) do
process_job_result(device, result, socket)
:ok
else
{:error, :device_not_found} ->
{:error, :device_not_found} = error ->
Logger.error("Device not found: #{result.device_id}")
error
{:error, :wrong_organization} ->
{:error, :wrong_organization} = error ->
Logger.error("Device #{result.device_id} not in agent's organization")
error
{:error, :device_reassigned} ->
{:error, :device_reassigned} = error ->
Logger.warning("Ignoring stale result for reassigned device",
device_id: result.device_id,
agent_token_id: socket.assigns.agent_token_id
)
error
end
end
@ -1358,9 +1470,11 @@ defmodule ToweropsWeb.AgentChannel do
case BackupRequests.get_request_by_job_id(result.job_id) do
nil ->
Logger.warning("Backup request not found for job_id", job_id: result.job_id)
{:error, :backup_request_not_found}
request ->
handle_backup_request(request, result)
:ok
end
end

View file

@ -989,4 +989,126 @@ defmodule Towerops.Snmp.Profiles.DynamicTest do
assert interfaces == []
end
end
describe "discover_sensors/2 with index templates" do
test "uses index_template from YAML to create unique sensor_index values" do
# This simulates MikroTik firewall connection count sensors which
# are scalar values (.0) but need unique indices
profile =
base_profile(%{
count_sensor_oids: [
%{
base_oid: "1.3.6.1.4.1.14988.1.1.22.1.1",
sensor_type: "count",
sensor_descr: "Total number of connections",
sensor_unit: "",
sensor_divisor: 1,
descr_oid: nil,
index_template: "mtxrCtTotalEntries.{{ $index }}"
},
%{
base_oid: "1.3.6.1.4.1.14988.1.1.22.1.2",
sensor_type: "count",
sensor_descr: "Total number of ipv4 connections",
sensor_unit: "",
sensor_divisor: 1,
descr_oid: nil,
index_template: "mtxrCtIP4Entries.{{ $index }}"
},
%{
base_oid: "1.3.6.1.4.1.14988.1.1.22.1.3",
sensor_type: "count",
sensor_descr: "Total number of ipv6 connections",
sensor_unit: "",
sensor_divisor: 1,
descr_oid: nil,
index_template: "mtxrCtIP6Entries.{{ $index }}"
}
]
})
stub(SnmpMock, :get, fn _, _, _ -> {:error, :no_such_object} end)
stub(SnmpMock, :walk, fn _, oid, _ ->
case oid do
"1.3.6.1.4.1.14988.1.1.22.1.1" ->
{:ok, [%{oid: "1.3.6.1.4.1.14988.1.1.22.1.1.0", value: {:integer, 1270}}]}
"1.3.6.1.4.1.14988.1.1.22.1.2" ->
{:ok, [%{oid: "1.3.6.1.4.1.14988.1.1.22.1.2.0", value: {:integer, 1270}}]}
"1.3.6.1.4.1.14988.1.1.22.1.3" ->
{:ok, [%{oid: "1.3.6.1.4.1.14988.1.1.22.1.3.0", value: {:integer, 0}}]}
_ ->
{:ok, []}
end
end)
assert {:ok, sensors} = Dynamic.discover_sensors(profile, @client_opts)
count_sensors = Enum.filter(sensors, &(&1.sensor_type == "count"))
assert length(count_sensors) == 3
# Each sensor should have a unique index based on the template
indices = count_sensors |> Enum.map(& &1.sensor_index) |> Enum.sort()
assert indices == ["mtxrCtIP4Entries.0", "mtxrCtIP6Entries.0", "mtxrCtTotalEntries.0"]
# Verify the values are correct
total = Enum.find(count_sensors, &(&1.sensor_index == "mtxrCtTotalEntries.0"))
assert total.last_value == 1270
assert total.sensor_descr == "Total number of connections"
ipv4 = Enum.find(count_sensors, &(&1.sensor_index == "mtxrCtIP4Entries.0"))
assert ipv4.last_value == 1270
assert ipv4.sensor_descr == "Total number of ipv4 connections"
ipv6 = Enum.find(count_sensors, &(&1.sensor_index == "mtxrCtIP6Entries.0"))
assert ipv6.last_value == 0
assert ipv6.sensor_descr == "Total number of ipv6 connections"
end
test "falls back to default index when index_template is nil" do
# When walking a single OID with multiple results, enumeration creates unique indices
profile =
base_profile(%{
count_sensor_oids: [
%{
base_oid: "1.3.6.1.4.1.99999.1.1",
sensor_type: "count",
sensor_descr: "Counter",
sensor_unit: "",
sensor_divisor: 1,
descr_oid: nil,
index_template: nil
}
]
})
stub(SnmpMock, :get, fn _, _, _ -> {:error, :no_such_object} end)
stub(SnmpMock, :walk, fn _, oid, _ ->
case oid do
"1.3.6.1.4.1.99999.1.1" ->
{:ok,
[
%{oid: "1.3.6.1.4.1.99999.1.1.1", value: {:integer, 42}},
%{oid: "1.3.6.1.4.1.99999.1.1.2", value: {:integer, 99}}
]}
_ ->
{:ok, []}
end
end)
assert {:ok, sensors} = Dynamic.discover_sensors(profile, @client_opts)
count_sensors = Enum.filter(sensors, &(&1.sensor_type == "count"))
assert length(count_sensors) == 2
# Should use default "count_1", "count_2" indices from enumeration
indices = count_sensors |> Enum.map(& &1.sensor_index) |> Enum.sort()
assert indices == ["count_1", "count_2"]
end
end
end

View file

@ -1946,8 +1946,9 @@ defmodule Towerops.SnmpTest do
}
]
count = Snmp.upsert_arp_entries(device.id, arp_data, [interface])
assert count == 2
{success, errors} = Snmp.upsert_arp_entries(device.id, arp_data, [interface])
assert success == 2
assert errors == 0
entries = Snmp.list_arp_entries(device.id)
assert length(entries) == 2
@ -1970,8 +1971,9 @@ defmodule Towerops.SnmpTest do
}
]
count = Snmp.upsert_arp_entries(device.id, arp_data, [])
assert count == 1
{success, errors} = Snmp.upsert_arp_entries(device.id, arp_data, [])
assert success == 1
assert errors == 0
entries = Snmp.list_arp_entries(device.id)
assert length(entries) == 1
@ -1980,6 +1982,72 @@ defmodule Towerops.SnmpTest do
end
end
describe "upsert_arp_entries/3 error tracking" do
test "returns success and error counts", %{device: device, snmp_device: snmp_device} do
interface =
%Interface{}
|> Interface.changeset(%{
snmp_device_id: snmp_device.id,
if_index: 1,
if_name: "eth0"
})
|> Repo.insert!()
arp_data = [
%{
ip_address: "192.168.1.100",
mac_address: "aa:bb:cc:dd:ee:ff",
if_index: 1,
entry_type: "dynamic",
last_seen_at: DateTime.utc_now()
},
%{
ip_address: "192.168.1.200",
mac_address: "11:22:33:44:55:66",
if_index: 1,
entry_type: "static",
last_seen_at: DateTime.utc_now()
}
]
{success, errors} = Snmp.upsert_arp_entries(device.id, arp_data, [interface])
assert success == 2
assert errors == 0
end
end
describe "upsert_mac_addresses/3 error tracking" do
test "returns success and error counts", %{device: device, snmp_device: snmp_device} do
interface =
%Interface{}
|> Interface.changeset(%{
snmp_device_id: snmp_device.id,
if_index: 1,
if_name: "eth0"
})
|> Repo.insert!()
mac_data = [
%{
mac_address: "aa:bb:cc:dd:ee:ff",
port_index: 1,
vlan_id: 1,
last_seen_at: DateTime.utc_now()
},
%{
mac_address: "11:22:33:44:55:66",
port_index: 1,
vlan_id: 1,
last_seen_at: DateTime.utc_now()
}
]
{success, errors} = Snmp.upsert_mac_addresses(device.id, mac_data, [interface])
assert success == 2
assert errors == 0
end
end
describe "list_discovered_devices_for_organization/1" do
test "merges data from neighbors, ARP, and MAC tables", %{
organization: organization,
@ -2507,4 +2575,132 @@ defmodule Towerops.SnmpTest do
assert length(neighbors) == 1
end
end
describe "delete_stale_and_upsert_neighbors/3" do
test "atomically deletes stale and upserts new neighbors", %{device: device, snmp_device: snmp_device} do
interface =
%Interface{}
|> Interface.changeset(%{snmp_device_id: snmp_device.id, if_index: 1, if_name: "eth0"})
|> Repo.insert!()
cutoff = DateTime.add(DateTime.utc_now(), -5, :minute)
# Create a stale neighbor (old timestamp)
{:ok, _stale} =
Snmp.upsert_neighbor(%{
device_id: device.id,
interface_id: interface.id,
remote_chassis_id: "stale-chassis",
remote_port_id: "ge0/0/1",
remote_system_name: "stale-switch",
protocol: "lldp",
last_discovered_at: DateTime.add(DateTime.utc_now(), -10, :minute)
})
# New neighbors to upsert
new_neighbors = [
%{
device_id: device.id,
interface_id: interface.id,
remote_chassis_id: "new-chassis",
remote_port_id: "ge0/0/2",
remote_system_name: "new-switch",
protocol: "lldp",
last_discovered_at: DateTime.utc_now()
}
]
assert {:ok, _} = Snmp.delete_stale_and_upsert_neighbors(device.id, new_neighbors, cutoff)
remaining = Repo.all(from n in Neighbor, where: n.device_id == ^device.id)
assert length(remaining) == 1
assert hd(remaining).remote_chassis_id == "new-chassis"
end
end
describe "delete_stale_and_upsert_arp_entries/4" do
test "atomically deletes stale and upserts new ARP entries", %{device: device, snmp_device: snmp_device} do
interface =
%Interface{}
|> Interface.changeset(%{snmp_device_id: snmp_device.id, if_index: 1, if_name: "eth0"})
|> Repo.insert!()
now = DateTime.truncate(DateTime.utc_now(), :second)
cutoff = DateTime.add(now, -5, :minute)
# Create a stale ARP entry
Repo.insert!(%ArpEntry{
device_id: device.id,
interface_id: interface.id,
ip_address: "10.0.0.99",
mac_address: "aa:bb:cc:dd:ee:ff",
last_seen_at: DateTime.add(now, -10, :minute)
})
# New entries to upsert
new_entries = [
%{
ip_address: "10.0.0.100",
mac_address: "11:22:33:44:55:66",
if_index: interface.if_index,
last_seen_at: now
}
]
assert {:ok, _} =
Snmp.delete_stale_and_upsert_arp_entries(
device.id,
new_entries,
[interface],
cutoff
)
remaining = Repo.all(from a in ArpEntry, where: a.device_id == ^device.id)
assert length(remaining) == 1
assert hd(remaining).ip_address == "10.0.0.100"
end
end
describe "delete_stale_and_upsert_mac_addresses/4" do
test "atomically deletes stale and upserts new MAC addresses", %{device: device, snmp_device: snmp_device} do
interface =
%Interface{}
|> Interface.changeset(%{snmp_device_id: snmp_device.id, if_index: 1, if_name: "eth0"})
|> Repo.insert!()
now = DateTime.truncate(DateTime.utc_now(), :second)
cutoff = DateTime.add(now, -5, :minute)
# Create a stale MAC entry
Repo.insert!(%MacAddress{
device_id: device.id,
interface_id: interface.id,
mac_address: "aa:bb:cc:dd:ee:ff",
vlan_id: 1,
last_seen_at: DateTime.add(now, -10, :minute)
})
# New entries to upsert
new_entries = [
%{
mac_address: "11:22:33:44:55:66",
port_index: interface.if_index,
vlan_id: 1,
last_seen_at: now
}
]
assert {:ok, _} =
Snmp.delete_stale_and_upsert_mac_addresses(
device.id,
new_entries,
[interface],
cutoff
)
remaining = Repo.all(from m in MacAddress, where: m.device_id == ^device.id)
assert length(remaining) == 1
assert hd(remaining).mac_address == "11:22:33:44:55:66"
end
end
end