defmodule ToweropsWeb.AgentChannel do @moduledoc """ Phoenix channel for bidirectional agent communication. Handles WebSocket connections from remote SNMP polling agents, distributing polling jobs and collecting metrics via Protocol Buffers binary encoding. ## Messages from Server → Agent (binary protobuf) - "jobs" - List of SNMP jobs to execute - "job" - Single new job (pushed when device added) ## Messages from Agent → Server (binary protobuf) - "result" - SNMP query results - "heartbeat" - Agent metadata (version, hostname, uptime) - "error" - Job execution errors All messages use Protocol Buffers binary encoding for efficiency. """ use ToweropsWeb, :channel alias Towerops.Agent.AgentError alias Towerops.Agent.AgentHeartbeat alias Towerops.Agent.AgentJob alias Towerops.Agent.AgentJobList alias Towerops.Agent.MikrotikCommand alias Towerops.Agent.MikrotikDevice alias Towerops.Agent.SnmpDevice alias Towerops.Agent.SnmpQuery alias Towerops.Agent.SnmpResult alias Towerops.Agents alias Towerops.Devices alias Towerops.Snmp alias Towerops.Snmp.AgentDiscovery require Logger @impl true @spec join(String.t(), map(), Phoenix.Socket.t()) :: {:ok, Phoenix.Socket.t()} | {:error, map()} def join("agent:" <> _agent_id, %{"token" => token}, socket) do # Verify agent token from join payload case Agents.verify_agent_token(token) do {:ok, agent_token} -> socket = socket |> assign(:agent_token_id, agent_token.id) |> assign(:organization_id, agent_token.organization_id) |> assign(:debug_enabled, agent_token.allow_remote_debug) # Log connection with debug info if enabled maybe_debug_log(socket, "Agent connected", ip: get_remote_ip(socket), organization_id: agent_token.organization_id ) # Subscribe to assignment changes for this agent _ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:assignments") # Subscribe to discovery requests for this agent _ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:discovery") # Update last_seen_at and IP on join remote_ip = get_remote_ip(socket) _ = Agents.update_agent_token_heartbeat(agent_token.id, remote_ip, %{}) # Send initial job list send(self(), :send_jobs) {:ok, socket} {:error, _} -> {:error, %{reason: "unauthorized"}} end end def join("agent:" <> _agent_id, _payload, _socket) do {:error, %{reason: "missing token"}} end @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) maybe_debug_log(socket, "Sending jobs to agent", job_count: length(jobs), device_ids: Enum.map(jobs, & &1.device_id), job_types: Enum.map(jobs, & &1.job_type) ) push(socket, "jobs", %{binary: Base.encode64(binary)}) {: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", agent_token_id: socket.assigns.agent_token_id ) send(self(), :send_jobs) {:noreply, socket} end # Handle PubSub broadcast when discovery is requested for a device def handle_info({:discovery_requested, device_id}, socket) do Logger.info("Discovery requested for device, sending discovery job to agent", agent_token_id: socket.assigns.agent_token_id, device_id: device_id ) case Devices.get_device_with_details(device_id) do nil -> Logger.error("Device not found for discovery request", device_id: device_id) {:noreply, socket} device -> # Build and send discovery job to agent job = build_discovery_job(device) job_list = %AgentJobList{jobs: [job]} binary = AgentJobList.encode(job_list) push(socket, "discovery_job", %{binary: Base.encode64(binary)}) {:noreply, socket} end end @impl true @spec handle_in(String.t(), map(), Phoenix.Socket.t()) :: {:noreply, Phoenix.Socket.t()} def handle_in("result", %{"binary" => binary_b64}, socket) do binary = Base.decode64!(binary_b64) result = SnmpResult.decode(binary) maybe_debug_log(socket, "Received SNMP result from agent", device_id: result.device_id, job_type: result.job_type, binary_size: byte_size(binary_b64), oid_count: map_size(result.oid_values) ) _ = process_snmp_result(socket.assigns.organization_id, result, socket) {:noreply, socket} end def handle_in("heartbeat", %{"binary" => binary_b64}, socket) do binary = Base.decode64!(binary_b64) heartbeat = AgentHeartbeat.decode(binary) metadata = %{ "version" => heartbeat.version, "hostname" => heartbeat.hostname, "uptime_seconds" => heartbeat.uptime_seconds } _ = Agents.update_agent_token_heartbeat( socket.assigns.agent_token_id, heartbeat.ip_address, metadata ) {:noreply, socket} end @impl true def handle_in("error", %{"binary" => binary_b64}, socket) do binary = Base.decode64!(binary_b64) error = AgentError.decode(binary) maybe_debug_log(socket, "Agent job error", device_id: error.device_id, error_message: error.message, binary_size: byte_size(binary_b64) ) Logger.error("Agent job error", agent_token_id: socket.assigns.agent_token_id, device_id: error.device_id, error: error.message ) {:noreply, socket} end # Private helpers @spec build_jobs_for_agent(Ecto.UUID.t()) :: [AgentJob.t()] defp build_jobs_for_agent(agent_token_id) do agent_token_id |> Agents.list_agent_polling_targets() |> Enum.flat_map(&build_jobs_for_device/1) end defp build_jobs_for_device(device) do # Build SNMP job (discovery or polling) snmp_job = if needs_discovery?(device) do build_discovery_job(device) else build_polling_job(device) end # Build MikroTik job if device is MikroTik and MikroTik API is enabled mikrotik_job = if mikrotik_device?(device) do build_mikrotik_job(device) end # Return list of jobs (filter out nil) Enum.reject([snmp_job, mikrotik_job], &is_nil/1) end defp needs_discovery?(device) do # Need discovery if no SNMP device or last discovery was >24 hours ago is_nil(device.snmp_device) or is_nil(device.last_discovery_at) or DateTime.diff(DateTime.utc_now(), device.last_discovery_at, :hour) > 24 end defp mikrotik_device?(device) do # Check if device is MikroTik based on SNMP discovery data device.snmp_device && (String.contains?(device.snmp_device.manufacturer || "", "MikroTik") || String.contains?(device.snmp_device.sys_descr || "", "RouterOS")) end defp build_discovery_job(device) do %AgentJob{ job_id: "discover:#{device.id}", job_type: :DISCOVER, device_id: device.id, snmp_device: %SnmpDevice{ ip: device.ip_address, community: Devices.resolve_snmp_community(device), version: device.snmp_version, port: device.snmp_port || 161 }, queries: build_discovery_queries() } end defp build_polling_job(device) do _snmp_device = device.snmp_device %AgentJob{ job_id: "poll:#{device.id}", job_type: :POLL, device_id: device.id, snmp_device: %SnmpDevice{ ip: device.ip_address, community: Devices.resolve_snmp_community(device), version: device.snmp_version, port: device.snmp_port || 161 }, queries: build_polling_queries(device) } end defp build_discovery_queries do [ # System info (GET) %SnmpQuery{ query_type: :GET, oids: [ # sysDescr "1.3.6.1.2.1.1.1.0", # sysObjectID "1.3.6.1.2.1.1.2.0", # sysUpTime "1.3.6.1.2.1.1.3.0", # sysContact "1.3.6.1.2.1.1.4.0", # sysName "1.3.6.1.2.1.1.5.0", # sysLocation "1.3.6.1.2.1.1.6.0" ] }, # Interface table (WALK) %SnmpQuery{ query_type: :WALK, # ifTable oids: ["1.3.6.1.2.1.2.2.1"] }, # Interface extensions (WALK) %SnmpQuery{ query_type: :WALK, # ifXTable oids: ["1.3.6.1.2.1.31.1.1.1"] }, # Sensor discovery - ENTITY-SENSOR-MIB (WALK) %SnmpQuery{ query_type: :WALK, oids: [ # ENTITY-SENSOR-MIB - sensor tables "1.3.6.1.2.1.99.1.1.1", # ENTITY-MIB - physical entity descriptions (needed for sensor names) "1.3.6.1.2.1.47.1.1.1.1.2", # entPhysicalName "1.3.6.1.2.1.47.1.1.1.1.7", # entPhysicalClass "1.3.6.1.2.1.47.1.1.1.1.5" ] }, # Vendor-specific MIBs (WALK) %SnmpQuery{ query_type: :WALK, oids: [ # Cisco "1.3.6.1.4.1.9", # MikroTik "1.3.6.1.4.1.14988", # Ubiquiti "1.3.6.1.4.1.41112" ] }, # Neighbor discovery (WALK) %SnmpQuery{ query_type: :WALK, oids: [ # LLDP "1.0.8802.1.1.2.1.4.1.1", # Cisco CDP "1.3.6.1.4.1.9.9.23" ] } ] end defp build_polling_queries(device) do sensor_oids = Enum.map(device.snmp_device.sensors, & &1.sensor_oid) interface_oids = Enum.flat_map(device.snmp_device.interfaces, fn iface -> idx = iface.if_index [ # ifInOctets "1.3.6.1.2.1.2.2.1.10.#{idx}", # ifOutOctets "1.3.6.1.2.1.2.2.1.16.#{idx}", # ifInErrors "1.3.6.1.2.1.2.2.1.14.#{idx}", # ifOutErrors "1.3.6.1.2.1.2.2.1.20.#{idx}", # ifInDiscards "1.3.6.1.2.1.2.2.1.13.#{idx}", # ifOutDiscards "1.3.6.1.2.1.2.2.1.19.#{idx}" ] end) # Neighbor discovery OIDs neighbor_oids = [ # LLDP "1.0.8802.1.1.2.1.4.1.1", # Cisco CDP "1.3.6.1.4.1.9.9.23" ] [ %SnmpQuery{ query_type: :GET, oids: sensor_oids ++ interface_oids }, %SnmpQuery{ query_type: :WALK, oids: neighbor_oids } ] end defp build_mikrotik_job(device) do mikrotik_config = Devices.get_mikrotik_config(device) # Only build job if MikroTik API is enabled and has credentials if mikrotik_config.enabled && mikrotik_config.username do %AgentJob{ job_id: "mikrotik:#{device.id}", job_type: :MIKROTIK, device_id: device.id, mikrotik_device: %MikrotikDevice{ ip: device.ip_address, username: mikrotik_config.username, password: mikrotik_config.password, port: mikrotik_config.port, use_ssl: mikrotik_config.use_ssl }, mikrotik_commands: build_mikrotik_commands() } end end defp build_mikrotik_commands do [ # System identity %MikrotikCommand{ command: "/system/identity/print", args: %{} }, # System resources %MikrotikCommand{ command: "/system/resource/print", args: %{} } ] 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) do process_job_result(device, result, socket) else {:error, :device_not_found} -> Logger.error("Device not found: #{result.device_id}") {:error, :wrong_organization} -> Logger.error("Device #{result.device_id} not in agent's organization") end end defp fetch_device(device_id) do case Devices.get_device_with_details(device_id) do nil -> {:error, :device_not_found} device -> {:ok, device} end end defp verify_device_organization(device, organization_id) do # Cloud pollers (organization_id = nil) can access any device if is_nil(organization_id) or device.site.organization_id == organization_id do :ok else {:error, :wrong_organization} end end defp process_job_result(device, %{job_type: :DISCOVER} = result, socket) do process_discovery_result(device, result, socket) end defp process_job_result(device, %{job_type: :POLL} = result, socket) do process_polling_result(device, result, socket) end defp process_discovery_result(device, result, socket) do # Agent has completed discovery and sent back SNMP data # Update last_discovery_at to signal completion to DiscoveryWorker Logger.info("Discovery results received from agent for #{device.name}, processing data") # Update device timestamp to mark discovery as complete case Devices.update_device(device, %{last_discovery_at: DateTime.utc_now()}) do {:ok, updated_device} -> Logger.info("Updated last_discovery_at for device #{device.name}") # Process the SNMP data from agent's discovery queries # The agent sends back oid_values map with all the collected data process_discovery_data(updated_device, result, socket) {:error, changeset} -> Logger.error("Failed to update last_discovery_at for device #{device.name}: #{inspect(changeset)}") end end defp process_discovery_data(device, result, socket) do oid_values = Map.new(result.oid_values) Logger.info("Processing discovery with #{map_size(oid_values)} OIDs", device_id: device.id, device_name: device.name ) # Debug log full OID values if debug is enabled maybe_debug_log(socket, "Discovery OID values received", device_id: device.id, device_name: device.name, oid_count: map_size(oid_values), sample_oids: oid_values |> Enum.take(10) |> Map.new(), all_oids: if(socket.assigns[:debug_enabled], do: oid_values, else: :redacted) ) # Process full discovery using agent data case AgentDiscovery.process_agent_discovery(device, oid_values) do {:ok, _discovered_device} -> Logger.info("Agent discovery completed", device_id: device.id, device_name: device.name ) maybe_debug_log(socket, "Discovery processing succeeded", device_id: device.id, device_name: device.name ) # Broadcast for real-time UI updates _ = Phoenix.PubSub.broadcast( Towerops.PubSub, "device:#{device.id}", {:discovery_completed, device.id} ) :ok {:error, reason} -> Logger.error("Agent discovery failed: #{inspect(reason)}", device_id: device.id, device_name: device.name ) maybe_debug_log(socket, "Discovery processing failed", device_id: device.id, device_name: device.name, error_reason: inspect(reason) ) # Don't update last_discovery_at - DiscoveryWorker will retry # with direct discovery as fallback {:error, reason} end end defp process_polling_result(device, result, _socket) do snmp_device = device.snmp_device oid_values = Map.new(result.oid_values) timestamp = DateTime.from_unix!(result.timestamp, :second) # Process sensor readings Enum.each(snmp_device.sensors, fn sensor -> if value = Map.get(oid_values, sensor.sensor_oid) do Snmp.create_sensor_reading(%{ sensor_id: sensor.id, value: parse_float(value) / sensor.sensor_divisor, status: "ok", checked_at: timestamp }) end end) # Process interface stats Enum.each(device.snmp_device.interfaces, fn iface -> idx = iface.if_index Snmp.create_interface_stat(%{ interface_id: iface.id, if_in_octets: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.10.#{idx}")), if_out_octets: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.16.#{idx}")), if_in_errors: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.14.#{idx}")), if_out_errors: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.20.#{idx}")), if_in_discards: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.13.#{idx}")), if_out_discards: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.19.#{idx}")), checked_at: timestamp }) end) # Note: Neighbor discovery is handled by separate polling worker # since it requires complex LLDP/CDP parsing end defp parse_integer(nil), do: nil defp parse_integer(value) when is_integer(value), do: value defp parse_integer(value) when is_binary(value) do case Integer.parse(value) do {int, _} -> int :error -> nil end end defp parse_integer(_), do: nil defp parse_float(value) when is_float(value), do: value defp parse_float(value) when is_integer(value), do: value / 1.0 defp parse_float(value) when is_binary(value) do case Float.parse(value) do {float, _} -> float :error -> nil end end defp parse_float(_), do: nil defp get_remote_ip(socket) do case socket.transport_pid do nil -> nil pid when is_pid(pid) -> case :ranch.get_addr(pid) do {ip, _port} when is_tuple(ip) -> format_ip(ip) _ -> nil end end rescue _ -> nil end defp format_ip({a, b, c, d}), do: "#{a}.#{b}.#{c}.#{d}" defp format_ip({a, b, c, d, e, f, g, h}), do: "#{Integer.to_string(a, 16)}:#{Integer.to_string(b, 16)}:#{Integer.to_string(c, 16)}:#{Integer.to_string(d, 16)}:#{Integer.to_string(e, 16)}:#{Integer.to_string(f, 16)}:#{Integer.to_string(g, 16)}:#{Integer.to_string(h, 16)}" # Debug logging helper - only logs when debug is enabled for this agent # Uses :info level since production logger filters out :debug defp maybe_debug_log(socket, message, metadata) do if socket.assigns[:debug_enabled] do Logger.info( message, Keyword.merge( [ agent_token_id: socket.assigns.agent_token_id, debug_enabled: true ], metadata ) ) end end end