- help_live/index.ex: 2,642→173 lines, 15 section modules + sidebar extracted MikroTik now real section, dead if false removed - agent_channel.ex: 2,506→1,751 lines, 3 helper modules extracted (heartbeat/subscriptions/job_builder), decode_and_process/4 eliminates 8 repeated handle_in patterns, guard-based size checks - antenna_catalog.ex: 1,174→47 lines, 107 specs → priv/antennas/catalog.json - topology.ex: unbounded query → batched loading (100/batch), all guard errors fixed, zero if/case/cond conditionals - proto: decoder_macros.ex + field_specs.ex infrastructure for macro-generated protobuf decoders
1746 lines
56 KiB
Elixir
1746 lines
56 KiB
Elixir
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.CheckList
|
|
alias Towerops.Agent.CheckResult, as: CheckResultProto
|
|
alias Towerops.Agent.CredentialTestResult
|
|
alias Towerops.Agent.LldpTopologyResult
|
|
alias Towerops.Agent.MikrotikResult
|
|
alias Towerops.Agent.MonitoringCheck, as: MonitoringCheckProto
|
|
alias Towerops.Agent.SnmpDevice
|
|
alias Towerops.Agent.SnmpResult
|
|
alias Towerops.Agents
|
|
alias Towerops.Alerts
|
|
alias Towerops.ConfigChanges
|
|
alias Towerops.Devices
|
|
alias Towerops.Devices.BackupRequests
|
|
alias Towerops.Devices.MikrotikBackups
|
|
alias Towerops.Monitoring
|
|
alias Towerops.Snmp
|
|
alias Towerops.Snmp.AgentDiscovery
|
|
alias Towerops.Snmp.ArpDiscovery
|
|
alias Towerops.Snmp.Discovery
|
|
alias Towerops.Snmp.MacDiscovery
|
|
alias Towerops.Snmp.NeighborDiscovery
|
|
alias Towerops.Snmp.Profiles.Base, as: SnmpProfiles
|
|
alias Towerops.Snmp.SensorChangeDetector
|
|
alias Towerops.Snmp.SensorScale
|
|
alias Towerops.Topology
|
|
alias ToweropsWeb.AgentChannel.Heartbeat
|
|
alias ToweropsWeb.AgentChannel.JobBuilder
|
|
alias ToweropsWeb.AgentChannel.Subscriptions, as: ChannelSubscriptions
|
|
alias ToweropsWeb.GraphQL.Subscriptions
|
|
alias ToweropsWeb.RemoteIp
|
|
|
|
require Logger
|
|
|
|
@typep socket :: Phoenix.Socket.t()
|
|
@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
|
|
# Maximum message size: 10MB (protobuf messages should be much smaller)
|
|
@max_message_size 10 * 1024 * 1024
|
|
|
|
@impl true
|
|
@spec join(String.t(), map(), Phoenix.Socket.t()) ::
|
|
{:ok, Phoenix.Socket.t()} | {:error, map()}
|
|
def join("agent:" <> agent_id, %{"token" => token} = payload, socket) do
|
|
Logger.info("Agent join attempt",
|
|
topic: "agent:#{agent_id}",
|
|
has_token: not is_nil(token),
|
|
token_length: if(is_binary(token), do: byte_size(token), else: 0),
|
|
payload_keys: Map.keys(payload)
|
|
)
|
|
|
|
do_join(Agents.verify_agent_token(token), socket)
|
|
end
|
|
|
|
def join("agent:" <> _agent_id, _payload, _socket) do
|
|
{:error, %{reason: "missing token"}}
|
|
end
|
|
|
|
defp do_join({:ok, agent_token}, socket) do
|
|
socket =
|
|
socket
|
|
|> assign(:agent_token_id, agent_token.id)
|
|
|> assign(:organization_id, agent_token.organization_id)
|
|
|> assign(:debug_enabled, agent_token.allow_remote_debug)
|
|
|
|
maybe_debug_log(socket, "Agent connected",
|
|
ip: get_remote_ip(socket),
|
|
organization_id: agent_token.organization_id
|
|
)
|
|
|
|
_ = ChannelSubscriptions.subscribe_all(agent_token)
|
|
|
|
now = DateTime.utc_now()
|
|
_ = Agents.update_agent_token_heartbeat(agent_token.id, get_remote_ip(socket), %{})
|
|
_ = ChannelSubscriptions.broadcast_connection(agent_token.id, agent_token.organization_id)
|
|
|
|
socket =
|
|
socket
|
|
|> assign(:last_heartbeat_at, now)
|
|
|> assign(:last_heartbeat_db_update, now)
|
|
|
|
Process.send_after(self(), :check_heartbeat, @heartbeat_check_interval_ms)
|
|
send(self(), :send_jobs)
|
|
|
|
{:ok, socket}
|
|
end
|
|
|
|
defp do_join({:error, _}, _socket) do
|
|
{:error, %{reason: "unauthorized"}}
|
|
end
|
|
|
|
@impl true
|
|
def terminate(_reason, %{assigns: %{agent_token_id: id, organization_id: org_id}}) do
|
|
_ = ChannelSubscriptions.broadcast_disconnection(id, org_id)
|
|
:ok
|
|
end
|
|
|
|
def terminate(_reason, _socket), do: :ok
|
|
|
|
@impl true
|
|
@spec handle_info(:send_jobs, Phoenix.Socket.t()) :: {:noreply, Phoenix.Socket.t()}
|
|
def handle_info(:send_jobs, socket) do
|
|
_ = cancel_poll_timer(socket)
|
|
|
|
updated_socket = build_and_push_jobs(socket)
|
|
build_and_push_check_jobs(updated_socket)
|
|
|
|
# Schedule next job dispatch cycle
|
|
interval_ms = poll_interval_ms()
|
|
timer = Process.send_after(self(), :send_jobs, interval_ms)
|
|
|
|
{:noreply, assign(updated_socket, :poll_timer, timer)}
|
|
end
|
|
|
|
# Poll a device immediately after discovery completes.
|
|
# Uses send(self(), ...) so discovery DB writes commit first.
|
|
def handle_info({:poll_after_discovery, device_id}, socket) do
|
|
case Devices.get_device_with_details(device_id) do
|
|
nil ->
|
|
Logger.error("Device not found for post-discovery poll", device_id: device_id)
|
|
|
|
device ->
|
|
jobs = [JobBuilder.build_polling_job(device)]
|
|
|
|
jobs =
|
|
if device.monitoring_enabled do
|
|
jobs ++ [JobBuilder.build_ping_job(device)]
|
|
else
|
|
jobs
|
|
end
|
|
|
|
job_list = %AgentJobList{jobs: jobs}
|
|
binary = AgentJobList.encode(job_list)
|
|
|
|
Logger.info("Sending immediate poll after discovery",
|
|
device_id: device_id,
|
|
device_name: device.name,
|
|
job_count: length(jobs)
|
|
)
|
|
|
|
push(socket, "jobs", %{binary: Base.encode64(binary)})
|
|
end
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
# 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, {:shutdown, :heartbeat_timeout}, socket}
|
|
else
|
|
Process.send_after(self(), :check_heartbeat, @heartbeat_check_interval_ms)
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
# Handle check changes — push updated check list to agent
|
|
def handle_info(:checks_changed, socket) do
|
|
build_and_push_check_jobs(socket)
|
|
{:noreply, socket}
|
|
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
|
|
)
|
|
|
|
{:stop, {:shutdown, :token_disabled}, socket}
|
|
end
|
|
|
|
# Handle restart request — push restart event to agent, then stop channel
|
|
def handle_info(:restart_requested, socket) do
|
|
Logger.info("Restart requested, notifying agent",
|
|
agent_token_id: socket.assigns.agent_token_id
|
|
)
|
|
|
|
push(socket, "restart", %{})
|
|
{:stop, {:shutdown, :restart_requested}, socket}
|
|
end
|
|
|
|
# Handle update request — push update event with download URL and checksum to agent
|
|
def handle_info({:update_requested, url, checksum}, socket) do
|
|
Logger.info("Update requested, notifying agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
url: url
|
|
)
|
|
|
|
push(socket, "update", %{url: url, checksum: checksum})
|
|
{:noreply, socket}
|
|
end
|
|
|
|
# Handle PubSub broadcast when device assignments change.
|
|
# Debounces rapid changes — cancels any pending refresh and schedules a new one.
|
|
def handle_info({:assignments_changed, event}, socket) do
|
|
Logger.info("Agent assignments changed, scheduling job refresh",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
event: event,
|
|
previous_device_ids: socket.assigns[:last_job_device_ids] || []
|
|
)
|
|
|
|
# Cancel any existing poll timer (regular cycle or pending debounce)
|
|
_ = cancel_poll_timer(socket)
|
|
|
|
# Schedule new refresh after debounce delay (configurable, defaults to 500ms in prod, 50ms in test)
|
|
debounce_ms = Application.get_env(:towerops, :agent_channel_debounce_ms, 500)
|
|
timer = Process.send_after(self(), :send_jobs, debounce_ms)
|
|
{:noreply, assign(socket, :poll_timer, timer)}
|
|
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 ->
|
|
job = JobBuilder.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
|
|
|
|
# Handle PubSub broadcast when backup is requested for a device
|
|
def handle_info({:backup_requested, device_id, job_id}, socket) do
|
|
Logger.info("Backup requested for device, assembling and sending backup job to agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
device_id: device_id,
|
|
job_id: job_id
|
|
)
|
|
|
|
case Devices.get_device_with_details(device_id) do
|
|
nil ->
|
|
Logger.warning("Backup requested but device not found",
|
|
device_id: device_id,
|
|
job_id: job_id
|
|
)
|
|
|
|
device ->
|
|
job = JobBuilder.build_backup_job(device, job_id)
|
|
job_list = %AgentJobList{jobs: [job]}
|
|
binary = AgentJobList.encode(job_list)
|
|
push(socket, "backup_job", %{binary: Base.encode64(binary)})
|
|
end
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
# Handle PubSub broadcast when credential test is requested
|
|
def handle_info({:credential_test_requested, test_id, snmp_config}, socket) do
|
|
Logger.info(
|
|
"AgentChannel sending credential test to agent: ip=#{snmp_config[:ip]} port=#{snmp_config[:port]} version=#{snmp_config[:version]} community_present=#{!is_nil(snmp_config[:community])} test_id=#{test_id}"
|
|
)
|
|
|
|
# Build SNMP device from config
|
|
snmp_device = %SnmpDevice{
|
|
ip: snmp_config[:ip],
|
|
port: snmp_config[:port],
|
|
version: snmp_config[:version],
|
|
community: snmp_config[:community],
|
|
v3_security_level: snmp_config[:v3_security_level],
|
|
v3_username: snmp_config[:v3_username],
|
|
v3_auth_protocol: snmp_config[:v3_auth_protocol],
|
|
v3_auth_password: snmp_config[:v3_auth_password],
|
|
v3_priv_protocol: snmp_config[:v3_priv_protocol],
|
|
v3_priv_password: snmp_config[:v3_priv_password]
|
|
}
|
|
|
|
# Build test credentials job
|
|
job = %AgentJob{
|
|
job_id: test_id,
|
|
job_type: :TEST_CREDENTIALS,
|
|
device_id: test_id,
|
|
snmp_device: snmp_device,
|
|
queries: [],
|
|
mikrotik_device: nil,
|
|
mikrotik_commands: []
|
|
}
|
|
|
|
job_list = %AgentJobList{jobs: [job]}
|
|
binary = AgentJobList.encode(job_list)
|
|
|
|
push(socket, "jobs", %{binary: Base.encode64(binary)})
|
|
{:noreply, socket}
|
|
end
|
|
|
|
# Handle PubSub broadcast when live polling is requested for a device
|
|
def handle_info({:live_poll_requested, device_id, sensor_oids, reply_topic}, socket) do
|
|
maybe_debug_log(socket, "Live poll requested for device",
|
|
device_id: device_id,
|
|
sensor_count: length(sensor_oids),
|
|
reply_topic: reply_topic
|
|
)
|
|
|
|
case Devices.get_device_with_details(device_id) do
|
|
nil ->
|
|
Logger.error("Device not found for live poll request", device_id: device_id)
|
|
{:noreply, socket}
|
|
|
|
device ->
|
|
job = JobBuilder.build_live_polling_job(device, sensor_oids, reply_topic)
|
|
job_list = %AgentJobList{jobs: [job]}
|
|
binary = AgentJobList.encode(job_list)
|
|
|
|
maybe_debug_log(socket, "Sending live poll job to agent",
|
|
device_id: device_id,
|
|
device_name: device.name,
|
|
oid_count: length(sensor_oids)
|
|
)
|
|
|
|
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
|
|
|
|
# Handle latency probe requests — send PING jobs for cloud-polled devices
|
|
def handle_info({:latency_probe_jobs, []}, socket) do
|
|
{:noreply, socket}
|
|
end
|
|
|
|
def handle_info({:latency_probe_jobs, devices}, socket) do
|
|
jobs =
|
|
Enum.flat_map(devices, fn device ->
|
|
for i <- 1..5 do
|
|
%AgentJob{
|
|
job_id: "probe:#{device.id}:#{i}",
|
|
job_type: :PING,
|
|
device_id: device.id,
|
|
snmp_device: %SnmpDevice{
|
|
ip: to_string(device.ip_address),
|
|
port: 0,
|
|
version: "",
|
|
community: ""
|
|
},
|
|
queries: []
|
|
}
|
|
end
|
|
end)
|
|
|
|
job_list = %AgentJobList{jobs: jobs}
|
|
binary = AgentJobList.encode(job_list)
|
|
|
|
Logger.info("Sending #{length(jobs)} latency probe pings to cloud poller",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
device_count: length(devices)
|
|
)
|
|
|
|
push(socket, "jobs", %{binary: Base.encode64(binary)})
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
@impl true
|
|
@spec handle_in(String.t(), map(), socket()) :: {:noreply, socket()}
|
|
|
|
def handle_in("result", %{"binary" => b}, socket) when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &SnmpResult.decode/1, &process_snmp_result_message/2, socket)
|
|
end
|
|
|
|
def handle_in("heartbeat", %{"binary" => b}, socket) when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &AgentHeartbeat.decode/1, &Heartbeat.process/2, socket)
|
|
end
|
|
|
|
def handle_in("error", %{"binary" => b}, socket) when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &AgentError.decode/1, &handle_agent_error/2, socket)
|
|
end
|
|
|
|
def handle_in("mikrotik_result", %{"binary" => b}, socket) when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &MikrotikResult.decode/1, &handle_validated_mikrotik_result/2, socket)
|
|
end
|
|
|
|
def handle_in("credential_test_result", %{"binary" => b}, socket)
|
|
when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &CredentialTestResult.decode/1, &handle_credential_test_result/2, socket)
|
|
end
|
|
|
|
def handle_in("monitoring_check", %{"binary" => b}, socket) when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &MonitoringCheckProto.decode/1, &handle_monitoring_check/2, socket)
|
|
end
|
|
|
|
def handle_in("check_result", %{"binary" => b}, socket) when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &CheckResultProto.decode/1, &handle_check_result/2, socket)
|
|
end
|
|
|
|
def handle_in("lldp_topology_result", %{"binary" => b}, socket)
|
|
when is_binary(b) and byte_size(b) <= @max_message_size do
|
|
decode_and_process(b, &LldpTopologyResult.decode/1, &handle_lldp_topology_result/2, socket)
|
|
end
|
|
|
|
# Oversized message handler — catches any binary that exceeds size limits
|
|
def handle_in(type, %{"binary" => b}, socket) when is_binary(b) do
|
|
Logger.warning("Rejected oversized #{type} from agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
size: byte_size(b),
|
|
max_size: @max_message_size
|
|
)
|
|
|
|
{:reply, {:error, %{reason: "Message too large (max #{@max_message_size} bytes)"}}, socket}
|
|
end
|
|
|
|
# Catch-all for unmatched events (diagnostic)
|
|
def handle_in(event, payload, socket) do
|
|
Logger.warning("Unhandled agent channel event",
|
|
event: event,
|
|
payload_keys: Map.keys(payload),
|
|
agent_token_id: socket.assigns.agent_token_id
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
## Decode-and-process helper
|
|
|
|
@spec decode_and_process(base64_string(), function(), function(), socket()) :: {:noreply, socket()}
|
|
defp decode_and_process(binary_b64, decode_fn, process_fn, socket) do
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, decoded} <- decode_fn.(binary) do
|
|
process_fn.(decoded, socket)
|
|
else
|
|
{:error, {type, message}} ->
|
|
Logger.error("Invalid message from agent: #{type} - #{message}",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
error_type: type,
|
|
error_message: message,
|
|
binary_size: byte_size(binary_b64)
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
## Message handler wrappers (normalize to (decoded, socket) -> {:noreply, socket})
|
|
|
|
defp handle_agent_error(error, socket) do
|
|
maybe_debug_log(socket, "Agent job error",
|
|
device_id: error.device_id,
|
|
error_message: error.message
|
|
)
|
|
|
|
Logger.error("Agent job error",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
device_id: error.device_id,
|
|
error: error.message
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
defp handle_credential_test_result(result, socket) do
|
|
Logger.info(
|
|
"AgentChannel received credential test result from agent, broadcasting to LiveView",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
test_id: result.test_id,
|
|
success: result.success,
|
|
has_error: result.error_message != "",
|
|
broadcast_topic: "credential_test:#{result.test_id}"
|
|
)
|
|
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"credential_test:#{result.test_id}",
|
|
{:credential_test_result, result}
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
defp handle_monitoring_check(check, socket) do
|
|
maybe_debug_log(socket, "Received ping result from agent",
|
|
device_id: check.device_id,
|
|
status: check.status,
|
|
response_time_ms: check.response_time_ms
|
|
)
|
|
|
|
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}
|
|
end
|
|
|
|
defp handle_check_result(result, socket) do
|
|
maybe_debug_log(socket, "Received check result from agent",
|
|
check_id: result.check_id,
|
|
status: result.status,
|
|
response_time_ms: result.response_time_ms
|
|
)
|
|
|
|
case store_check_result(result, socket) do
|
|
:ok ->
|
|
:ok
|
|
|
|
{:error, reason} ->
|
|
Logger.error("Failed to store check result",
|
|
check_id: result.check_id,
|
|
reason: inspect(reason)
|
|
)
|
|
end
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
defp handle_lldp_topology_result(%LldpTopologyResult{} = result, socket) do
|
|
maybe_debug_log(socket, "Received LLDP topology result from agent",
|
|
device_id: result.device_id,
|
|
job_id: result.job_id,
|
|
neighbor_count: length(result.neighbors),
|
|
local_system_name: result.local_system_name
|
|
)
|
|
|
|
case store_lldp_neighbors(result) do
|
|
{:ok, count} ->
|
|
Logger.info("Stored LLDP neighbors for device",
|
|
device_id: result.device_id,
|
|
neighbor_count: count
|
|
)
|
|
|
|
{:error, reason} ->
|
|
Logger.error("Failed to store LLDP neighbors",
|
|
device_id: result.device_id,
|
|
reason: inspect(reason)
|
|
)
|
|
end
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
## SNMP result processing
|
|
|
|
defp process_snmp_result_message(result, socket) do
|
|
maybe_debug_log(socket, "Received SNMP result from agent",
|
|
device_id: result.device_id,
|
|
job_type: result.job_type,
|
|
job_id: result.job_id,
|
|
oid_count: map_size(result.oid_values)
|
|
)
|
|
|
|
_ = do_process_snmp_result(result.job_id, result, socket)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
|
|
defp do_process_snmp_result("live_poll:" <> _rest, result, socket) do
|
|
handle_live_poll_result(result, socket)
|
|
end
|
|
|
|
defp do_process_snmp_result(_job_id, result, socket) do
|
|
process_and_log_snmp_result(socket, result)
|
|
end
|
|
|
|
# Handle live poll result — broadcast to reply topic instead of saving
|
|
defp handle_live_poll_result(result, socket) do
|
|
# Extract reply topic from job_id: "live_poll:device_id:reply_topic"
|
|
case String.split(result.job_id, ":", parts: 3) do
|
|
["live_poll", _device_id, reply_topic] ->
|
|
maybe_debug_log(socket, "Broadcasting live poll result",
|
|
reply_topic: reply_topic,
|
|
oid_count: map_size(result.oid_values)
|
|
)
|
|
|
|
# Broadcast OID values directly to the waiting LiveView
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
reply_topic,
|
|
{:live_poll_result, result.oid_values}
|
|
)
|
|
|
|
_ ->
|
|
Logger.warning("Invalid live poll job_id format: #{result.job_id}")
|
|
end
|
|
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 = JobBuilder.build_jobs_for_agent(socket.assigns.agent_token_id)
|
|
job_list = %AgentJobList{jobs: jobs}
|
|
binary = AgentJobList.encode(job_list)
|
|
|
|
# Extract device IDs for better debugging
|
|
device_ids = jobs |> Enum.map(& &1.device_id) |> Enum.uniq() |> Enum.sort()
|
|
|
|
Logger.info("Sending #{length(jobs)} jobs to agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
device_count: length(device_ids),
|
|
device_ids: inspect(device_ids),
|
|
job_ids: inspect(Enum.map(jobs, & &1.job_id))
|
|
)
|
|
|
|
# Store last job list to detect changes on next refresh
|
|
socket = assign(socket, :last_job_device_ids, device_ids)
|
|
|
|
push(socket, "jobs", %{binary: Base.encode64(binary)})
|
|
socket
|
|
end
|
|
|
|
# Builds service check jobs for the agent and pushes them via "check_jobs" event.
|
|
defp build_and_push_check_jobs(socket) do
|
|
agent_token_id = socket.assigns.agent_token_id
|
|
checks = JobBuilder.list_checks_for_agent(agent_token_id)
|
|
|
|
do_push_check_jobs(socket, checks)
|
|
end
|
|
|
|
defp do_push_check_jobs(_socket, []), do: :ok
|
|
|
|
defp do_push_check_jobs(socket, checks) do
|
|
check_list = JobBuilder.build_check_list(checks)
|
|
binary = CheckList.encode(check_list)
|
|
|
|
Logger.info("Sending #{length(check_list.checks)} check jobs to agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
check_ids: inspect(Enum.map(checks, & &1.id))
|
|
)
|
|
|
|
push(socket, "check_jobs", %{binary: Base.encode64(binary)})
|
|
end
|
|
|
|
@spec handle_validated_mikrotik_result(MikrotikResult.t(), socket()) :: {:noreply, socket()}
|
|
defp handle_validated_mikrotik_result(result, socket) do
|
|
_binary = MikrotikResult.encode(result)
|
|
|
|
maybe_debug_log(socket, "Received MikroTik result from agent",
|
|
device_id: result.device_id,
|
|
job_id: result.job_id,
|
|
has_error: result.error != "",
|
|
sentence_count: length(result.sentences)
|
|
)
|
|
|
|
# Check if this is a backup job by the job_id prefix
|
|
if String.starts_with?(result.job_id, "backup:") do
|
|
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)
|
|
end
|
|
|
|
{:noreply, socket}
|
|
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 ->
|
|
Logger.error("Device not found: #{result.device_id}")
|
|
error
|
|
|
|
{:error, :wrong_organization} = error ->
|
|
Logger.error("Device #{result.device_id} not in agent's organization")
|
|
error
|
|
|
|
{: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
|
|
|
|
# Verify device is still assigned to this agent before processing results
|
|
defp verify_device_assignment(device, agent_token_id) do
|
|
# Get effective agent token (resolves inheritance from site/org)
|
|
effective_agent_token_id = Agents.get_effective_agent_token(device)
|
|
|
|
# Check if it matches the requesting agent
|
|
if effective_agent_token_id == agent_token_id do
|
|
# Also check if there's an explicit device assignment that's disabled
|
|
case Agents.get_device_assignment(device.id) do
|
|
%{enabled: false} ->
|
|
{:error, :device_reassigned}
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
else
|
|
{:error, :device_reassigned}
|
|
end
|
|
end
|
|
|
|
defp store_monitoring_check(check, socket) do
|
|
organization_id = socket.assigns.organization_id
|
|
agent_token_id = socket.assigns.agent_token_id
|
|
|
|
with {:ok, device} <- fetch_device(check.device_id),
|
|
:ok <- verify_device_organization(device, organization_id) do
|
|
# Use server time to avoid clock skew issues
|
|
# Agent timestamps can be unreliable due to clock drift
|
|
checked_at = Towerops.Time.now()
|
|
|
|
# Keep status as string for storage, convert to atom for device status update
|
|
status_string = check.status
|
|
status_atom = String.to_existing_atom(status_string)
|
|
|
|
attrs = %{
|
|
device_id: check.device_id,
|
|
agent_token_id: agent_token_id,
|
|
status: status_string,
|
|
response_time_ms: check.response_time_ms,
|
|
checked_at: checked_at
|
|
}
|
|
|
|
case Monitoring.create_monitoring_check(attrs) do
|
|
{:ok, _check} ->
|
|
update_device_status_from_check(device, status_atom)
|
|
:ok
|
|
|
|
{:error, changeset} ->
|
|
Logger.error("Failed to store ping result",
|
|
device_id: check.device_id,
|
|
errors: inspect(changeset.errors)
|
|
)
|
|
|
|
{:error, :storage_failed}
|
|
end
|
|
else
|
|
{:error, :device_not_found} ->
|
|
Logger.error("Device not found for ping result: #{check.device_id}")
|
|
{:error, :device_not_found}
|
|
|
|
{:error, :wrong_organization} ->
|
|
Logger.error("Device #{check.device_id} not in agent's organization")
|
|
{:error, :wrong_organization}
|
|
end
|
|
end
|
|
|
|
defp store_check_result(%CheckResultProto{} = result, socket) do
|
|
organization_id = socket.assigns.organization_id
|
|
agent_token_id = socket.assigns.agent_token_id
|
|
|
|
with check when not is_nil(check) <- Monitoring.get_check(result.check_id),
|
|
true <- check.organization_id == organization_id do
|
|
record_and_process_check_result(check, result, organization_id, agent_token_id)
|
|
else
|
|
nil ->
|
|
Logger.error("Check not found for check result: #{result.check_id}")
|
|
{:error, :check_not_found}
|
|
|
|
false ->
|
|
Logger.error("Check #{result.check_id} not in agent's organization")
|
|
{:error, :wrong_organization}
|
|
end
|
|
end
|
|
|
|
defp record_and_process_check_result(check, result, organization_id, agent_token_id) do
|
|
now = Towerops.Time.now()
|
|
|
|
case Monitoring.create_check_result(%{
|
|
check_id: check.id,
|
|
organization_id: organization_id,
|
|
status: result.status,
|
|
output: result.output,
|
|
response_time_ms: result.response_time_ms,
|
|
checked_at: now,
|
|
agent_token_id: agent_token_id
|
|
}) do
|
|
{:ok, _} ->
|
|
old_state = check.current_state
|
|
{:ok, updated_check} = Monitoring.update_check_state(check, result.status, result.output)
|
|
handle_check_state_change(old_state, updated_check, result.output)
|
|
_ = broadcast_check_update(check.device_id)
|
|
:ok
|
|
|
|
{:error, changeset} ->
|
|
Logger.error("Failed to store check result",
|
|
check_id: check.id,
|
|
errors: inspect(changeset.errors)
|
|
)
|
|
|
|
{:error, :storage_failed}
|
|
end
|
|
end
|
|
|
|
defp broadcast_check_update(nil), do: :ok
|
|
|
|
defp broadcast_check_update(device_id) do
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"device:#{device_id}",
|
|
{:monitoring_check_updated, device_id}
|
|
)
|
|
end
|
|
|
|
defp handle_check_state_change(old_state, check, output) do
|
|
# State changed from OK to problem and is now hard state
|
|
_ =
|
|
if old_state == 0 and check.current_state in [1, 2] and check.current_state_type == "hard" do
|
|
maybe_create_check_alert(check, output)
|
|
end
|
|
|
|
# State changed from problem to OK
|
|
if old_state in [1, 2] and check.current_state == 0 do
|
|
Alerts.resolve_check_alerts(check.id)
|
|
end
|
|
end
|
|
|
|
defp maybe_create_check_alert(check, output) do
|
|
if !Alerts.has_active_check_alert?(check.id) do
|
|
severity = if check.current_state == 1, do: 1, else: 2
|
|
|
|
Alerts.create_alert(%{
|
|
check_id: check.id,
|
|
device_id: check.device_id,
|
|
organization_id: check.organization_id,
|
|
alert_type: "check_#{check.check_type}",
|
|
severity: severity,
|
|
message: "Check '#{check.name}' is #{Monitoring.Check.state_label(check.current_state)}",
|
|
output: output,
|
|
triggered_at: DateTime.utc_now()
|
|
})
|
|
end
|
|
end
|
|
|
|
defp store_lldp_neighbors(%LldpTopologyResult{} = result) do
|
|
now = DateTime.utc_now()
|
|
device_id = result.device_id
|
|
|
|
# Store each neighbor using Topology module's upsert function
|
|
results =
|
|
Enum.map(result.neighbors, fn neighbor ->
|
|
# Call the private upsert function via the public Topology API
|
|
# The neighbor struct from protobuf matches the expected format
|
|
Topology.upsert_neighbor(device_id, neighbor, now)
|
|
end)
|
|
|
|
success_count = Enum.count(results, &match?({:ok, _}, &1))
|
|
{:ok, success_count}
|
|
rescue
|
|
e ->
|
|
Logger.error("Exception storing LLDP neighbors",
|
|
device_id: result.device_id,
|
|
error: Exception.message(e)
|
|
)
|
|
|
|
{:error, :storage_exception}
|
|
end
|
|
|
|
defp update_device_status_from_check(device, check_status) do
|
|
new_status = if check_status in [:success, :up, :ok], do: :up, else: :down
|
|
old_status = device.status
|
|
|
|
Devices.update_device_status(device, new_status)
|
|
|
|
_ =
|
|
if old_status != new_status do
|
|
handle_agent_status_change(device, old_status, new_status)
|
|
end
|
|
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"device:#{device.id}",
|
|
{:device_status_changed, device.id, new_status, nil}
|
|
)
|
|
|
|
Subscriptions.publish_device_status(
|
|
device.id,
|
|
device.organization_id,
|
|
%{device_name: device.name, status: to_string(new_status)}
|
|
)
|
|
end
|
|
|
|
defp handle_agent_status_change(device, _old_status, :down) do
|
|
if !Alerts.has_active_alert?(device.id, :device_down) do
|
|
now = Towerops.Time.now()
|
|
|
|
message =
|
|
if device.snmp_enabled,
|
|
do: "Device is not responding to SNMP",
|
|
else: "Device is not responding to ping"
|
|
|
|
_ =
|
|
Alerts.create_alert(%{
|
|
device_id: device.id,
|
|
alert_type: :device_down,
|
|
triggered_at: now,
|
|
message: message
|
|
})
|
|
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"alerts:org:#{device.organization_id}:new",
|
|
{:new_alert, device.id, :device_down}
|
|
)
|
|
end
|
|
end
|
|
|
|
defp handle_agent_status_change(device, _old_status, :up) do
|
|
now = Towerops.Time.now()
|
|
|
|
message =
|
|
if device.snmp_enabled,
|
|
do: "Device is now responding to SNMP",
|
|
else: "Device is now responding to ping"
|
|
|
|
_ =
|
|
Alerts.create_alert(%{
|
|
device_id: device.id,
|
|
alert_type: :device_up,
|
|
triggered_at: now,
|
|
resolved_at: now,
|
|
message: message
|
|
})
|
|
|
|
case Alerts.get_active_alert(device.id, :device_down) do
|
|
nil -> :ok
|
|
alert -> Alerts.resolve_alert(alert)
|
|
end
|
|
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"alerts:org:#{device.organization_id}:resolved",
|
|
{:alert_resolved, device.id, :device_down}
|
|
)
|
|
end
|
|
|
|
# Resolve any stuck alerts for a device that's currently up
|
|
# This handles edge cases where device is UP but alerts didn't resolve (e.g., due to bugs)
|
|
defp resolve_any_stuck_down_alerts(device) do
|
|
now = Towerops.Time.now()
|
|
|
|
# Resolve stuck device_down alerts
|
|
_ =
|
|
case Alerts.get_active_alert(device.id, :device_down) do
|
|
nil ->
|
|
:ok
|
|
|
|
alert ->
|
|
Logger.info(
|
|
"Resolving stuck device_down alert for #{device.name} - device is UP but alert was not resolved",
|
|
device_id: device.id,
|
|
alert_id: alert.id
|
|
)
|
|
|
|
# Create device_up alert (informational)
|
|
message =
|
|
if device.snmp_enabled,
|
|
do: "Device is now responding to SNMP",
|
|
else: "Device is now responding to ping"
|
|
|
|
_ =
|
|
Alerts.create_alert(%{
|
|
device_id: device.id,
|
|
alert_type: :device_up,
|
|
triggered_at: now,
|
|
resolved_at: now,
|
|
message: message
|
|
})
|
|
|
|
# Resolve the stuck alert
|
|
case Alerts.resolve_alert(alert) do
|
|
{:ok, _} ->
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"alerts:org:#{device.organization_id}:resolved",
|
|
{:alert_resolved, device.id, :device_down}
|
|
)
|
|
|
|
{:error, reason} ->
|
|
Logger.error(
|
|
"Failed to resolve stuck device_down alert #{alert.id} for #{device.name}: #{inspect(reason)}",
|
|
device_id: device.id
|
|
)
|
|
end
|
|
end
|
|
|
|
# Also resolve any stuck device_up alerts (should have been created with resolved_at)
|
|
case Alerts.get_active_alert(device.id, :device_up) do
|
|
nil ->
|
|
:ok
|
|
|
|
alert ->
|
|
Logger.info(
|
|
"Resolving stuck device_up alert for #{device.name} - should have been auto-resolved",
|
|
device_id: device.id,
|
|
alert_id: alert.id
|
|
)
|
|
|
|
Alerts.resolve_alert(alert)
|
|
end
|
|
|
|
:ok
|
|
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.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_job_result(device, %{job_type: :PING} = result, socket) do
|
|
# Backward compatibility: older agents send PING results as SnmpResult instead of MonitoringCheck
|
|
# Newer agents should use the "monitoring_check" event instead
|
|
Logger.warning("Agent sent PING result as SnmpResult (deprecated). Agent should be upgraded.",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
device_id: device.id,
|
|
job_id: result.job_id
|
|
)
|
|
|
|
# PING jobs sent as SnmpResult have empty oid_values
|
|
# If we receive a result (not an error), assume success
|
|
check = %{
|
|
device_id: device.id,
|
|
status: "success",
|
|
# Unknown - older agents don't send this
|
|
response_time_ms: 0.0,
|
|
timestamp: result.timestamp
|
|
}
|
|
|
|
store_monitoring_check(check, socket)
|
|
end
|
|
|
|
defp process_discovery_result(device, result, socket) do
|
|
Logger.info("Discovery results received from agent for #{device.name}, processing data")
|
|
|
|
# Process discovery data first, only update last_discovery_at on success.
|
|
# This ensures failed discoveries are retried on the next poll cycle
|
|
# instead of being blocked for 24 hours.
|
|
process_discovery_data(device, result, socket)
|
|
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} ->
|
|
# Only mark discovery timestamp on success so failures are retried
|
|
# on the next poll cycle instead of waiting 24 hours
|
|
_ = Devices.update_device(device, %{last_discovery_at: DateTime.utc_now()})
|
|
|
|
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}
|
|
)
|
|
|
|
# Poll the device immediately so it doesn't sit idle until the next cycle
|
|
send(self(), {:poll_after_discovery, 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)
|
|
)
|
|
|
|
{:error, reason}
|
|
end
|
|
end
|
|
|
|
defp process_polling_result(device, result, _socket) do
|
|
Devices.update_snmp_poll_time(device)
|
|
|
|
# Successful SNMP poll means device is up — update status and resolve any stuck alerts
|
|
old_status = device.status
|
|
new_status = :up
|
|
|
|
# Update device status if needed
|
|
if old_status != new_status do
|
|
Devices.update_device_status(device, new_status)
|
|
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"device:#{device.id}",
|
|
{:device_status_changed, device.id, new_status, nil}
|
|
)
|
|
|
|
Subscriptions.publish_device_status(
|
|
device.id,
|
|
device.organization_id,
|
|
%{device_name: device.name, status: to_string(new_status)}
|
|
)
|
|
end
|
|
|
|
# Always check for and resolve stuck device_down alerts after successful poll
|
|
# This handles edge cases where alerts got stuck (device already UP but alert not resolved)
|
|
resolve_any_stuck_down_alerts(device)
|
|
|
|
snmp_device = device.snmp_device
|
|
oid_values = normalize_oid_keys(result.oid_values)
|
|
# Use server time to avoid clock drift issues between agents
|
|
timestamp = Towerops.Time.now()
|
|
|
|
Logger.debug(
|
|
"Processing polling result for #{device.name}: #{map_size(oid_values)} OIDs, #{length(snmp_device.sensors)} sensors, #{length(snmp_device.interfaces)} interfaces"
|
|
)
|
|
|
|
# Batch insert sensor readings and process metadata updates
|
|
process_sensor_readings_batch(snmp_device.sensors, oid_values, timestamp)
|
|
|
|
# Publish sensor readings to GraphQL subscriptions
|
|
publish_sensor_readings_to_graphql(device, snmp_device.sensors, oid_values, timestamp)
|
|
|
|
# Batch insert interface stats
|
|
_ = process_interface_stats_batch(snmp_device.interfaces, oid_values, timestamp)
|
|
|
|
# Process complex SNMP data (neighbors, ARP, MAC, IP addresses, processors, storage)
|
|
# using Replay adapter to reuse existing parsing logic
|
|
_ =
|
|
if map_size(oid_values) > length(snmp_device.sensors) + length(snmp_device.interfaces) * 6 do
|
|
process_additional_polling_data(device, oid_values)
|
|
end
|
|
|
|
# Notify LiveViews that new sensor/interface data is available for graphing
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"device:#{device.id}",
|
|
{:state_sensors_updated, device.id}
|
|
)
|
|
end
|
|
|
|
# Process neighbors, ARP, MAC, IP addresses, processors, storage from agent polling
|
|
defp process_additional_polling_data(device, oid_values) do
|
|
# Build Replay adapter opts to reuse existing parsing logic
|
|
client_opts = [
|
|
ip: device.ip_address,
|
|
adapter: Towerops.Snmp.Adapters.Replay,
|
|
oid_map: oid_values
|
|
]
|
|
|
|
# Reuse existing polling functions with Replay adapter.
|
|
# Runs in a spawned process to avoid blocking the channel.
|
|
device_id = device.id
|
|
device_name = device.name
|
|
snmp_device_id = device.snmp_device.id
|
|
|
|
Task.Supervisor.start_child(Towerops.TaskSupervisor, fn ->
|
|
try do
|
|
# Fetch interfaces once for reuse across multiple operations
|
|
interfaces = Snmp.list_interfaces(snmp_device_id)
|
|
|
|
# Add device_id to interfaces for neighbor discovery (required by build_neighbor_record)
|
|
interfaces_with_device_id = Enum.map(interfaces, &Map.put(&1, :device_id, device_id))
|
|
|
|
# Process neighbors (LLDP/CDP)
|
|
process_neighbors(client_opts, device_id, interfaces_with_device_id)
|
|
|
|
# Process ARP entries
|
|
process_arp_entries(client_opts, device_id, interfaces)
|
|
|
|
# Process MAC addresses
|
|
_ = process_mac_addresses(client_opts, device_id, interfaces)
|
|
|
|
# Process IP addresses
|
|
process_ip_addresses(client_opts, device_id, interfaces)
|
|
|
|
# Process processors
|
|
process_processors(client_opts, snmp_device_id)
|
|
|
|
# Process storage
|
|
process_storage(client_opts, snmp_device_id)
|
|
rescue
|
|
e ->
|
|
Logger.error(
|
|
"Polling data processing crashed for device #{device_name}: #{Exception.message(e)}",
|
|
device_id: device_id,
|
|
error: Exception.format(:error, e, __STACKTRACE__)
|
|
)
|
|
end
|
|
end)
|
|
end
|
|
|
|
defp process_neighbors(client_opts, device_id, interfaces_with_device_id) do
|
|
case NeighborDiscovery.discover_neighbors(client_opts, interfaces_with_device_id) do
|
|
{:ok, neighbors} when neighbors != [] ->
|
|
Discovery.save_neighbors(device_id, neighbors)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp process_arp_entries(client_opts, device_id, interfaces) do
|
|
case ArpDiscovery.discover_arp_table(client_opts) do
|
|
{:ok, arp_entries} when arp_entries != [] ->
|
|
Discovery.save_arp_entries(device_id, arp_entries, interfaces)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp process_mac_addresses(client_opts, device_id, interfaces) do
|
|
case MacDiscovery.discover_mac_table(client_opts) do
|
|
{:ok, mac_addresses} when mac_addresses != [] ->
|
|
Snmp.upsert_mac_addresses(device_id, mac_addresses, interfaces)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp process_ip_addresses(client_opts, device_id, interfaces) do
|
|
case SnmpProfiles.discover_all_ip_addresses(client_opts) do
|
|
{:ok, ip_addresses} when ip_addresses != [] ->
|
|
discovered_device = %{
|
|
device_id: device_id,
|
|
interfaces: interfaces
|
|
}
|
|
|
|
Discovery.sync_ip_addresses(discovered_device, ip_addresses)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp process_processors(client_opts, snmp_device_id) do
|
|
case SnmpProfiles.discover_processors(client_opts) do
|
|
{:ok, processors} when processors != [] ->
|
|
discovered_device = %{id: snmp_device_id}
|
|
Discovery.sync_processors(discovered_device, processors)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp process_storage(client_opts, snmp_device_id) do
|
|
case SnmpProfiles.discover_storage(client_opts) do
|
|
{:ok, storage} when storage != [] ->
|
|
discovered_device = %{id: snmp_device_id}
|
|
Discovery.sync_storage(discovered_device, storage)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp process_sensor_readings_batch(sensors, oid_values, timestamp) do
|
|
# Build reading entries and collect sensor updates in one pass
|
|
{reading_entries, sensor_updates} =
|
|
Enum.reduce(sensors, {[], []}, fn sensor, {readings, updates} ->
|
|
case resolve_sensor_value(sensor, oid_values) do
|
|
nil ->
|
|
{readings, updates}
|
|
|
|
final_value ->
|
|
reading = %{
|
|
sensor_id: sensor.id,
|
|
value: final_value,
|
|
status: "ok",
|
|
checked_at: timestamp
|
|
}
|
|
|
|
{[reading | readings], [{sensor, final_value} | updates]}
|
|
end
|
|
end)
|
|
|
|
# Batch insert all sensor readings in one SQL statement
|
|
_ = Snmp.create_sensor_readings_batch(reading_entries)
|
|
|
|
# Process change detection and metadata updates individually
|
|
Enum.each(sensor_updates, fn {sensor, final_value} ->
|
|
SensorChangeDetector.detect_and_broadcast(sensor, final_value, timestamp)
|
|
|
|
Snmp.update_sensor(sensor, %{
|
|
last_value: final_value,
|
|
last_checked_at: timestamp
|
|
})
|
|
end)
|
|
end
|
|
|
|
defp publish_sensor_readings_to_graphql(device, sensors, oid_values, timestamp) do
|
|
readings = do_build_readings(sensors, oid_values, timestamp)
|
|
do_publish_readings(device, readings)
|
|
end
|
|
|
|
defp do_build_readings(sensors, oid_values, timestamp) do
|
|
sensors
|
|
|> Enum.filter(fn sensor ->
|
|
normalized_oid = String.trim_leading(sensor.sensor_oid, ".")
|
|
Map.has_key?(oid_values, normalized_oid)
|
|
end)
|
|
|> Enum.map(fn sensor ->
|
|
normalized_oid = String.trim_leading(sensor.sensor_oid, ".")
|
|
raw_value = Map.get(oid_values, normalized_oid)
|
|
parsed = parse_float(raw_value)
|
|
|
|
value =
|
|
if parsed do
|
|
divided = parsed / sensor.sensor_divisor
|
|
scaled = SensorScale.normalize(sensor.sensor_type, divided)
|
|
Float.round(scaled / 1.0, 2)
|
|
end
|
|
|
|
%{
|
|
sensor_id: sensor.id,
|
|
sensor_type: sensor.sensor_type,
|
|
sensor_descr: sensor.sensor_descr,
|
|
sensor_unit: sensor.sensor_unit,
|
|
value: value,
|
|
checked_at: to_string(timestamp)
|
|
}
|
|
end)
|
|
|> Enum.reject(fn r -> is_nil(r.value) end)
|
|
end
|
|
|
|
defp do_publish_readings(_device, []), do: :ok
|
|
|
|
defp do_publish_readings(device, readings) do
|
|
Subscriptions.publish_sensor_readings(
|
|
device.id,
|
|
device.organization_id,
|
|
device.name,
|
|
readings
|
|
)
|
|
end
|
|
|
|
defp resolve_sensor_value(sensor, oid_values) do
|
|
normalized_oid = String.trim_leading(sensor.sensor_oid, ".")
|
|
|
|
with value when not is_nil(value) <- Map.get(oid_values, normalized_oid),
|
|
parsed when not is_nil(parsed) <- parse_float(value) do
|
|
divided = parsed / sensor.sensor_divisor
|
|
normalized = SensorScale.normalize(sensor.sensor_type, divided)
|
|
Float.round(normalized / 1.0, 2)
|
|
else
|
|
_ -> nil
|
|
end
|
|
end
|
|
|
|
defp process_interface_stats_batch(interfaces, oid_values, timestamp) do
|
|
entries =
|
|
Enum.map(interfaces, fn iface ->
|
|
idx = iface.if_index
|
|
|
|
# Prefer 64-bit HC counters; fall back to 32-bit for devices that don't support ifXTable
|
|
in_octets =
|
|
parse_integer(Map.get(oid_values, "1.3.6.1.2.1.31.1.1.1.6.#{idx}")) ||
|
|
parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.10.#{idx}"))
|
|
|
|
out_octets =
|
|
parse_integer(Map.get(oid_values, "1.3.6.1.2.1.31.1.1.1.10.#{idx}")) ||
|
|
parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.16.#{idx}"))
|
|
|
|
%{
|
|
interface_id: iface.id,
|
|
if_in_octets: in_octets,
|
|
if_out_octets: out_octets,
|
|
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)
|
|
|
|
Snmp.create_interface_stats_batch(entries)
|
|
end
|
|
|
|
defp parse_integer(value) do
|
|
case Towerops.Numeric.parse_integer(value) do
|
|
{:ok, i} -> i
|
|
:error -> nil
|
|
end
|
|
end
|
|
|
|
# Normalize OID keys to strip leading dots.
|
|
# The agent (via gosnmp/net-snmp) returns OIDs with leading dots
|
|
# (e.g., ".1.3.6.1.2.1.1.1.0") but sensor OIDs and hardcoded interface
|
|
# OIDs use the format without (e.g., "1.3.6.1.2.1.1.1.0").
|
|
defp normalize_oid_keys(oid_values) do
|
|
Map.new(oid_values, fn {oid, value} ->
|
|
{String.trim_leading(oid, "."), value}
|
|
end)
|
|
end
|
|
|
|
defp parse_float(value) do
|
|
case Towerops.Numeric.parse_float(value) do
|
|
{:ok, f} -> f
|
|
:error -> nil
|
|
end
|
|
end
|
|
|
|
defp get_remote_ip(socket), do: RemoteIp.from_socket(socket)
|
|
|
|
# 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
|
|
|
|
defp process_backup_result(result) 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
|
|
|
|
defp handle_backup_request(request, %{error: ""} = result) do
|
|
case extract_config_from_result(result) do
|
|
{:ok, config_text} ->
|
|
# Manual backups always create a snapshot, even if config unchanged
|
|
# Automated backups (daily_cron) use deduplication to save storage
|
|
should_create_backup =
|
|
request.source == "manual" || backup_needed?(result.device_id, config_text)
|
|
|
|
if should_create_backup do
|
|
create_and_save_backup(request, result, config_text)
|
|
else
|
|
mark_backup_skipped(request, result)
|
|
end
|
|
|
|
{:error, reason} ->
|
|
mark_backup_failed(request, result, "Failed to extract config: #{reason}")
|
|
end
|
|
end
|
|
|
|
defp handle_backup_request(request, result) do
|
|
mark_backup_failed(request, result, result.error)
|
|
end
|
|
|
|
defp create_and_save_backup(request, result, config_text) do
|
|
case MikrotikBackups.create_backup(result.device_id, config_text, request.source) do
|
|
{:ok, new_backup} ->
|
|
BackupRequests.update_request_status(request.id, "success", nil)
|
|
Logger.info("Backup created", device_id: result.device_id, source: request.source)
|
|
|
|
# Emit config change event if there was a previous backup
|
|
maybe_emit_config_change_event(result.device_id, new_backup)
|
|
|
|
{:error, reason} ->
|
|
mark_backup_failed(request, result, "Failed to create backup: #{inspect(reason)}")
|
|
end
|
|
end
|
|
|
|
defp maybe_emit_config_change_event(device_id, new_backup) do
|
|
# Find the previous backup (second most recent)
|
|
backups = MikrotikBackups.list_device_backups(device_id, limit: 2)
|
|
|
|
case backups do
|
|
[_new, previous | _] ->
|
|
maybe_record_change(device_id, previous, new_backup)
|
|
|
|
_ ->
|
|
:ok
|
|
end
|
|
end
|
|
|
|
defp maybe_record_change(_device_id, %{config_hash_normalized: hash}, %{config_hash_normalized: hash}), do: :ok
|
|
|
|
defp maybe_record_change(device_id, previous, new_backup) do
|
|
case ConfigChanges.record_change_event(device_id, previous, new_backup) do
|
|
{:ok, %{id: event_id} = event} ->
|
|
Logger.info("Config change event recorded", device_id: device_id, event_id: event_id)
|
|
|
|
Task.Supervisor.start_child(Towerops.TaskSupervisor, fn ->
|
|
ConfigChanges.Correlator.correlate(event)
|
|
end)
|
|
|
|
{:ok, :no_changes} ->
|
|
:ok
|
|
|
|
{:error, reason} ->
|
|
Logger.warning("Failed to record config change event: #{inspect(reason)}",
|
|
device_id: device_id
|
|
)
|
|
end
|
|
end
|
|
|
|
defp backup_needed?(device_id, config_text) do
|
|
{:ok, needs_backup?} = MikrotikBackups.needs_backup?(device_id, config_text)
|
|
needs_backup?
|
|
end
|
|
|
|
defp mark_backup_skipped(request, result) do
|
|
BackupRequests.update_request_status(request.id, "success", nil)
|
|
Logger.info("Backup skipped (no changes)", device_id: result.device_id)
|
|
end
|
|
|
|
defp mark_backup_failed(request, result, error_message) do
|
|
BackupRequests.update_request_status(request.id, "failed", error_message)
|
|
Logger.error("Backup failed", device_id: result.device_id, error: error_message)
|
|
end
|
|
|
|
defp extract_config_from_result(result) do
|
|
# SSH backups return a single sentence with "config" attribute (new method)
|
|
# API backups use /file/read which returns chunks with "data" attribute (old method)
|
|
|
|
# Try SSH backup format first (single sentence with "config" attribute)
|
|
config_direct =
|
|
Enum.find_value(result.sentences, fn sentence -> Map.get(sentence.attributes, "config") end)
|
|
|
|
case config_direct do
|
|
nil ->
|
|
# Try chunked API backup format
|
|
extract_config_from_chunked_data(result)
|
|
|
|
config_text when is_binary(config_text) and config_text != "" ->
|
|
{:ok, config_text}
|
|
|
|
_ ->
|
|
{:error, "Config attribute found but empty or invalid"}
|
|
end
|
|
end
|
|
|
|
defp extract_config_from_chunked_data(result) do
|
|
# The old API backup workflow uses /file/read which returns chunks with "data" attribute
|
|
data_chunks =
|
|
result.sentences
|
|
|> Enum.filter(fn sentence -> Map.has_key?(sentence.attributes, "data") end)
|
|
|> Enum.map(fn sentence -> Map.get(sentence.attributes, "data") end)
|
|
|
|
case data_chunks do
|
|
[] ->
|
|
# Fallback: try other methods
|
|
extract_config_from_fallback(result)
|
|
|
|
chunks ->
|
|
# Combine all chunks in order
|
|
config_text = Enum.join(chunks, "")
|
|
|
|
if config_text == "" do
|
|
{:error, "No config data in result - all chunks were empty"}
|
|
else
|
|
{:ok, config_text}
|
|
end
|
|
end
|
|
end
|
|
|
|
defp extract_config_from_fallback(result) do
|
|
# Fallback: try to find any sentence with configuration-like content
|
|
# Look for sentences with common RouterOS config markers
|
|
config_text =
|
|
Enum.find_value(result.sentences, fn sentence ->
|
|
find_config_in_attributes(sentence.attributes)
|
|
end)
|
|
|
|
case config_text do
|
|
nil -> {:error, "No config data in result - received #{length(result.sentences)} sentences"}
|
|
text -> {:ok, text}
|
|
end
|
|
end
|
|
|
|
defp find_config_in_attributes(attributes) do
|
|
Enum.find_value(attributes, fn {_key, value} ->
|
|
if routeros_config?(value), do: value
|
|
end)
|
|
end
|
|
|
|
defp routeros_config?(value) when is_binary(value) do
|
|
String.contains?(value, ["# ", "/", "set "]) and String.length(value) > 100
|
|
end
|
|
|
|
defp routeros_config?(_), do: false
|
|
|
|
defp cancel_poll_timer(socket) do
|
|
if timer = socket.assigns[:poll_timer] do
|
|
Process.cancel_timer(timer)
|
|
end
|
|
end
|
|
|
|
defp poll_interval_ms do
|
|
Application.get_env(:towerops, :agent_poll_interval_ms, 60_000)
|
|
end
|
|
|
|
## Safe Encoding/Decoding Helpers
|
|
|
|
# Safely decode Base64-encoded protobuf binary.
|
|
# Returns {:ok, binary} on success or {:error, {type, message}} on failure.
|
|
@spec safe_base64_decode(base64_string()) :: {:ok, protobuf_binary()} | {:error, {atom(), String.t()}}
|
|
defp safe_base64_decode(base64_string) when is_binary(base64_string) do
|
|
case Base.decode64(base64_string) do
|
|
{:ok, binary} ->
|
|
{:ok, binary}
|
|
|
|
:error ->
|
|
{:error, {:invalid_base64, "Message is not valid Base64-encoded data"}}
|
|
end
|
|
rescue
|
|
e in ArgumentError ->
|
|
Logger.error("Base64 decode error: #{Exception.message(e)}")
|
|
{:error, {:invalid_base64, "Failed to decode Base64 data"}}
|
|
end
|
|
end
|