diff --git a/lib/towerops/agents.ex b/lib/towerops/agents.ex index ea64daa8..80412974 100644 --- a/lib/towerops/agents.ex +++ b/lib/towerops/agents.ex @@ -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) diff --git a/lib/towerops/profiles/yaml_profiles.ex b/lib/towerops/profiles/yaml_profiles.ex index 15e6d981..31088f80 100644 --- a/lib/towerops/profiles/yaml_profiles.ex +++ b/lib/towerops/profiles/yaml_profiles.ex @@ -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 diff --git a/lib/towerops/snmp.ex b/lib/towerops/snmp.ex index 5f2f75ff..eaff81d8 100644 --- a/lib/towerops/snmp.ex +++ b/lib/towerops/snmp.ex @@ -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 diff --git a/lib/towerops/snmp/discovery.ex b/lib/towerops/snmp/discovery.ex index 25fb35dd..5f3c23ee 100644 --- a/lib/towerops/snmp/discovery.ex +++ b/lib/towerops/snmp/discovery.ex @@ -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 diff --git a/lib/towerops/snmp/profiles/dynamic.ex b/lib/towerops/snmp/profiles/dynamic.ex index 982f54ad..1e8b4ff5 100644 --- a/lib/towerops/snmp/profiles/dynamic.ex +++ b/lib/towerops/snmp/profiles/dynamic.ex @@ -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 diff --git a/lib/towerops/workers/device_poller_worker.ex b/lib/towerops/workers/device_poller_worker.ex index 7f837560..70d93489 100644 --- a/lib/towerops/workers/device_poller_worker.ex +++ b/lib/towerops/workers/device_poller_worker.ex @@ -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}") diff --git a/lib/towerops_web/channels/agent_channel.ex b/lib/towerops_web/channels/agent_channel.ex index fe90e240..831ab8a2 100644 --- a/lib/towerops_web/channels/agent_channel.ex +++ b/lib/towerops_web/channels/agent_channel.ex @@ -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 diff --git a/test/towerops/snmp/profiles/dynamic_test.exs b/test/towerops/snmp/profiles/dynamic_test.exs index 0b773208..ebdd0b6a 100644 --- a/test/towerops/snmp/profiles/dynamic_test.exs +++ b/test/towerops/snmp/profiles/dynamic_test.exs @@ -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 diff --git a/test/towerops/snmp_test.exs b/test/towerops/snmp_test.exs index 97bdeea6..f212c64d 100644 --- a/test/towerops/snmp_test.exs +++ b/test/towerops/snmp_test.exs @@ -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