Replace individual Repo.insert() with Repo.insert_all() for sensor readings, interface stats, processor readings, and storage readings in agent channel and device poller worker. Add partial/composite database indexes for common query patterns. Dashboard now uses GROUP BY aggregation instead of loading all devices. DeviceLive.Show only reloads data for the active tab on PubSub events. DeviceLive.Index debounces status changes and skips quota queries on updates. Agent polling preloads filtered to monitored sensors/interfaces only.
1866 lines
58 KiB
Elixir
1866 lines
58 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.AgentJob
|
|
alias Towerops.Agent.AgentJobList
|
|
alias Towerops.Agent.MikrotikCommand
|
|
alias Towerops.Agent.MikrotikDevice
|
|
alias Towerops.Agent.MikrotikResult
|
|
alias Towerops.Agent.SnmpDevice
|
|
alias Towerops.Agent.SnmpQuery
|
|
alias Towerops.Agent.Validator
|
|
alias Towerops.Agents
|
|
alias Towerops.Alerts
|
|
alias Towerops.Devices
|
|
alias Towerops.Devices.BackupRequests
|
|
alias Towerops.Devices.MikrotikBackups
|
|
alias Towerops.Monitoring
|
|
alias Towerops.Snmp
|
|
alias Towerops.Snmp.AgentDiscovery
|
|
alias Towerops.Snmp.Discovery
|
|
alias Towerops.Snmp.SensorChangeDetector
|
|
|
|
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
|
|
|
|
@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)
|
|
)
|
|
|
|
# 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")
|
|
|
|
# Subscribe to backup requests for this agent
|
|
_ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:backup")
|
|
|
|
# Subscribe to credential test requests for this agent
|
|
_ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:credential_test")
|
|
|
|
# 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, %{})
|
|
|
|
# Broadcast agent connection for real-time UI updates
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"agents:health",
|
|
{:agent_connected, agent_token.id, agent_token.organization_id}
|
|
)
|
|
|
|
# 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 after join completes (push/3 cannot be called during join)
|
|
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
|
|
def terminate(_reason, socket) do
|
|
# Broadcast agent disconnection for real-time UI updates
|
|
if Map.has_key?(socket.assigns, :agent_token_id) do
|
|
_ =
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"agents:health",
|
|
{:agent_disconnected, socket.assigns.agent_token_id, socket.assigns.organization_id}
|
|
)
|
|
end
|
|
|
|
:ok
|
|
end
|
|
|
|
@impl true
|
|
@spec handle_info(:send_jobs, Phoenix.Socket.t()) :: {:noreply, Phoenix.Socket.t()}
|
|
def handle_info(:send_jobs, socket) do
|
|
build_and_push_jobs(socket)
|
|
{:noreply, socket}
|
|
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 = [build_polling_job(device)]
|
|
|
|
jobs =
|
|
if device.monitoring_enabled do
|
|
jobs ++ [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, :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
|
|
)
|
|
|
|
{:stop, :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, :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 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
|
|
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
|
|
|
|
# Handle PubSub broadcast when backup is requested for a device
|
|
def handle_info({:backup_requested, job}, socket) do
|
|
Logger.info("Backup requested for device, sending backup job to agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
device_id: job.device_id,
|
|
job_id: job.job_id
|
|
)
|
|
|
|
job_list = %AgentJobList{jobs: [job]}
|
|
binary = AgentJobList.encode(job_list)
|
|
|
|
push(socket, "backup_job", %{binary: Base.encode64(binary)})
|
|
{: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=#{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 ->
|
|
# Build live polling job with only requested sensor OIDs
|
|
job = 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
|
|
|
|
@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
|
|
Logger.info("Received SNMP result from agent (binary size: #{byte_size(binary_b64)})")
|
|
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, result} <- Validator.validate_snmp_result(binary) 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,
|
|
binary_size: byte_size(binary_b64),
|
|
oid_count: map_size(result.oid_values)
|
|
)
|
|
|
|
# Check if this is a live poll result
|
|
if String.starts_with?(result.job_id, "live_poll:") do
|
|
handle_live_poll_result(result, socket)
|
|
else
|
|
process_and_log_snmp_result(socket, result)
|
|
end
|
|
|
|
{:noreply, socket}
|
|
else
|
|
{:error, {type, message}} ->
|
|
Logger.error("Invalid SNMP result 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}
|
|
|
|
{:error, :base64_decode_failed} ->
|
|
Logger.error("Failed to decode SNMP result (invalid base64)",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
binary_size: byte_size(binary_b64)
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
@spec handle_in(String.t(), %{required(String.t()) => base64_string()}, socket()) ::
|
|
{:noreply, socket()}
|
|
def handle_in("heartbeat", %{"binary" => binary_b64}, socket) when is_binary(binary_b64) do
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, heartbeat} <- Validator.validate_heartbeat(binary) do
|
|
metadata = %{
|
|
"version" => heartbeat.version,
|
|
"uptime_seconds" => heartbeat.uptime_seconds,
|
|
"arch" => heartbeat.arch
|
|
}
|
|
|
|
_ =
|
|
Agents.update_agent_token_heartbeat(
|
|
socket.assigns.agent_token_id,
|
|
get_remote_ip(socket),
|
|
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(
|
|
Towerops.PubSub,
|
|
"agents:health",
|
|
{:agent_heartbeat, socket.assigns.agent_token_id, socket.assigns.organization_id}
|
|
)
|
|
|
|
{:noreply, socket}
|
|
else
|
|
{:error, {type, message}} ->
|
|
Logger.error("Invalid heartbeat from agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
error_type: type,
|
|
error_message: message
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
def handle_in("error", %{"binary" => binary_b64}, socket) when is_binary(binary_b64) do
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, error} <- Validator.validate_agent_error(binary) do
|
|
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}
|
|
else
|
|
{:error, {type, message}} ->
|
|
Logger.error("Invalid error message from agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
error_type: type,
|
|
error_message: message
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
def handle_in("mikrotik_result", %{"binary" => binary_b64}, socket) when is_binary(binary_b64) do
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, result} <- Validator.validate_mikrotik_result(binary) do
|
|
handle_validated_mikrotik_result(result, socket)
|
|
else
|
|
{:error, {type, message}} ->
|
|
Logger.error("Invalid MikroTik result from agent",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
error_type: type,
|
|
error_message: message
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
def handle_in("credential_test_result", %{"binary" => binary_b64}, socket) when is_binary(binary_b64) do
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, result} <- Validator.validate_credential_test_result(binary) 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}"
|
|
)
|
|
|
|
# Broadcast result to the requesting LiveView via PubSub
|
|
# The test_id contains the device_id, so we can broadcast to that topic
|
|
Phoenix.PubSub.broadcast(
|
|
Towerops.PubSub,
|
|
"credential_test:#{result.test_id}",
|
|
{:credential_test_result, result}
|
|
)
|
|
|
|
{:noreply, socket}
|
|
else
|
|
{:error, {type, message}} ->
|
|
Logger.error(
|
|
"Invalid credential test result 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}
|
|
|
|
{:error, :base64_decode_failed} ->
|
|
Logger.error("Failed to decode credential test result (invalid base64)",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
binary_size: byte_size(binary_b64)
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
end
|
|
|
|
@spec handle_in(String.t(), %{required(String.t()) => base64_string()}, socket()) ::
|
|
{:noreply, socket()}
|
|
def handle_in("monitoring_check", %{"binary" => binary_b64}, socket) when is_binary(binary_b64) do
|
|
with {:ok, binary} <- safe_base64_decode(binary_b64),
|
|
{:ok, check} <- Validator.validate_monitoring_check_message(binary) 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,
|
|
binary_size: byte_size(binary_b64)
|
|
)
|
|
|
|
# Store ping result in database
|
|
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
|
|
{:error, {type, message}} ->
|
|
Logger.error("Invalid ping result 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}
|
|
|
|
{:error, :base64_decode_failed} ->
|
|
Logger.error("Failed to decode ping result (invalid base64)",
|
|
agent_token_id: socket.assigns.agent_token_id,
|
|
binary_size: byte_size(binary_b64)
|
|
)
|
|
|
|
{:noreply, socket}
|
|
end
|
|
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
|
|
|
|
# Private helpers
|
|
|
|
@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
|
|
|
|
# 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
|
|
|> 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 only during discovery (not during regular polling)
|
|
# MikroTik commands are discovery-type operations that should run once per 24h
|
|
mikrotik_job =
|
|
if needs_discovery?(device) and mikrotik_device?(device) do
|
|
build_mikrotik_job(device)
|
|
end
|
|
|
|
# Build ping job if monitoring is enabled
|
|
ping_job =
|
|
if device.monitoring_enabled do
|
|
build_ping_job(device)
|
|
end
|
|
|
|
# Return list of jobs (filter out nil)
|
|
Enum.reject([snmp_job, mikrotik_job, ping_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
|
|
|
|
# Resolves SNMP credentials for a device.
|
|
# For SNMPv3, returns a map with resolved credentials via cascade.
|
|
# For v1/v2c, returns a simple map with community string and version.
|
|
defp resolve_snmp_credentials(device) do
|
|
if device.snmp_version == "3" do
|
|
Devices.get_snmpv3_config(device)
|
|
else
|
|
%{
|
|
community: Devices.resolve_snmp_community(device),
|
|
version: device.snmp_version
|
|
}
|
|
end
|
|
end
|
|
|
|
# Checks if SNMP credentials are present (not nil/empty).
|
|
defp credentials_present?(%{community: community}) when is_binary(community) do
|
|
community != ""
|
|
end
|
|
|
|
defp credentials_present?(%{username: username}) when is_binary(username) do
|
|
username != ""
|
|
end
|
|
|
|
defp credentials_present?(_), do: false
|
|
|
|
# Builds the SnmpDevice protobuf message with appropriate credentials.
|
|
# For v1/v2c: Uses community string
|
|
# For v3: Uses SNMPv3 credentials (security_level, username, auth/priv protocols/passwords)
|
|
defp build_snmp_device_message(device, snmp_config) do
|
|
if device.snmp_version == "3" do
|
|
build_v3_snmp_device(device, snmp_config)
|
|
else
|
|
build_v2c_snmp_device(device, snmp_config)
|
|
end
|
|
end
|
|
|
|
defp build_v3_snmp_device(device, snmp_config) do
|
|
log_v3_device_build(device, snmp_config)
|
|
|
|
%SnmpDevice{
|
|
ip: device.ip_address,
|
|
version: device.snmp_version,
|
|
port: device.snmp_port || 161,
|
|
community: "",
|
|
v3_security_level: snmp_config.security_level || "",
|
|
v3_username: snmp_config.username || "",
|
|
v3_auth_protocol: snmp_config.auth_protocol || "",
|
|
v3_auth_password: snmp_config.auth_password || "",
|
|
v3_priv_protocol: snmp_config.priv_protocol || "",
|
|
v3_priv_password: snmp_config.priv_password || ""
|
|
}
|
|
end
|
|
|
|
defp log_v3_device_build(device, snmp_config) do
|
|
auth_password_present = !is_nil(snmp_config.auth_password) && snmp_config.auth_password != ""
|
|
priv_password_present = !is_nil(snmp_config.priv_password) && snmp_config.priv_password != ""
|
|
|
|
Logger.info(
|
|
"Building SNMPv3 device message",
|
|
device_id: device.id,
|
|
device_name: device.name,
|
|
device_ip: device.ip_address,
|
|
security_level: snmp_config.security_level || "",
|
|
username: snmp_config.username || "",
|
|
auth_protocol: snmp_config.auth_protocol || "",
|
|
auth_password_present: auth_password_present,
|
|
priv_protocol: snmp_config.priv_protocol || "",
|
|
priv_password_present: priv_password_present
|
|
)
|
|
end
|
|
|
|
defp build_v2c_snmp_device(device, snmp_config) do
|
|
community = snmp_config.community || ""
|
|
|
|
%SnmpDevice{
|
|
ip: device.ip_address,
|
|
version: device.snmp_version,
|
|
port: device.snmp_port || 161,
|
|
community: community
|
|
}
|
|
end
|
|
|
|
defp build_discovery_job(device) do
|
|
snmp_credentials = resolve_snmp_credentials(device)
|
|
|
|
maybe_debug_log(%{assigns: %{agent_token_id: "system"}}, "Building discovery job",
|
|
device_id: device.id,
|
|
device_name: device.name,
|
|
snmp_version: device.snmp_version,
|
|
credentials_present: credentials_present?(snmp_credentials),
|
|
credential_source:
|
|
if(device.snmp_version == "3",
|
|
do: Map.get(snmp_credentials, :source),
|
|
else: device.snmp_community_source
|
|
)
|
|
)
|
|
|
|
%AgentJob{
|
|
job_id: "discover:#{device.id}",
|
|
job_type: :DISCOVER,
|
|
device_id: device.id,
|
|
snmp_device: build_snmp_device_message(device, snmp_credentials),
|
|
queries: build_discovery_queries()
|
|
}
|
|
end
|
|
|
|
defp build_polling_job(device) do
|
|
_snmp_device = device.snmp_device
|
|
snmp_credentials = resolve_snmp_credentials(device)
|
|
|
|
maybe_debug_log(%{assigns: %{agent_token_id: "system"}}, "Building polling job",
|
|
device_id: device.id,
|
|
device_name: device.name,
|
|
snmp_version: device.snmp_version,
|
|
credentials_present: credentials_present?(snmp_credentials),
|
|
credential_source:
|
|
if(device.snmp_version == "3",
|
|
do: Map.get(snmp_credentials, :source),
|
|
else: device.snmp_community_source
|
|
),
|
|
site_id: device.site_id,
|
|
site_loaded: not is_nil(Map.get(device, :site))
|
|
)
|
|
|
|
%AgentJob{
|
|
job_id: "poll:#{device.id}",
|
|
job_type: :POLL,
|
|
device_id: device.id,
|
|
snmp_device: build_snmp_device_message(device, snmp_credentials),
|
|
queries: build_polling_queries(device)
|
|
}
|
|
end
|
|
|
|
defp build_ping_job(device) do
|
|
%AgentJob{
|
|
job_id: "ping:#{device.id}",
|
|
job_type: :PING,
|
|
device_id: device.id,
|
|
# Ping only needs IP address - use minimal SnmpDevice message
|
|
snmp_device: %SnmpDevice{
|
|
ip: to_string(device.ip_address),
|
|
port: 0,
|
|
version: "",
|
|
community: ""
|
|
},
|
|
queries: []
|
|
}
|
|
end
|
|
|
|
defp build_live_polling_job(device, sensor_oids, reply_topic) do
|
|
snmp_credentials = resolve_snmp_credentials(device)
|
|
|
|
%AgentJob{
|
|
job_id: "live_poll:#{device.id}:#{reply_topic}",
|
|
job_type: :POLL,
|
|
device_id: device.id,
|
|
snmp_device: build_snmp_device_message(device, snmp_credentials),
|
|
queries: [
|
|
%SnmpQuery{
|
|
query_type: :GET,
|
|
oids: sensor_oids
|
|
}
|
|
]
|
|
}
|
|
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
|
|
|
|
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"
|
|
]
|
|
},
|
|
# State sensors - ENTITY-STATE-MIB (WALK)
|
|
%SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: ["1.3.6.1.2.1.131.1.1.1.1"]
|
|
},
|
|
# HOST-RESOURCES-MIB (WALK) - processors, storage, devices
|
|
%SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# hrProcessorTable
|
|
"1.3.6.1.2.1.25.3.3",
|
|
# hrStorageTable
|
|
"1.3.6.1.2.1.25.2.3",
|
|
# hrDeviceTable (for processor descriptions)
|
|
"1.3.6.1.2.1.25.3.2"
|
|
]
|
|
},
|
|
# UCD-SNMP-MIB (GET) - Linux CPU stats fallback
|
|
%SnmpQuery{
|
|
query_type: :GET,
|
|
oids: [
|
|
# ssCpuUser
|
|
"1.3.6.1.4.1.2021.11.9.0",
|
|
# ssCpuSystem
|
|
"1.3.6.1.4.1.2021.11.10.0",
|
|
# ssCpuIdle
|
|
"1.3.6.1.4.1.2021.11.11.0"
|
|
]
|
|
},
|
|
# Vendor-specific MIBs (WALK) - targeted subtrees only
|
|
%SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# Cisco PROCESS-MIB (CPU)
|
|
"1.3.6.1.4.1.9.9.109",
|
|
# Cisco ENVMON-MIB (temp/voltage/fan/PSU)
|
|
"1.3.6.1.4.1.9.9.13",
|
|
# Cisco CDP
|
|
"1.3.6.1.4.1.9.9.23",
|
|
# Cisco WAP
|
|
"1.3.6.1.4.1.9.9.618",
|
|
# MikroTik
|
|
"1.3.6.1.4.1.14988",
|
|
# Ubiquiti
|
|
"1.3.6.1.4.1.41112",
|
|
# Juniper
|
|
"1.3.6.1.4.1.2636",
|
|
# HP/H3C/Comware
|
|
"1.3.6.1.4.1.25506",
|
|
# Fortinet
|
|
"1.3.6.1.4.1.12356",
|
|
# Cambium
|
|
"1.3.6.1.4.1.17713"
|
|
]
|
|
},
|
|
# 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"
|
|
]
|
|
},
|
|
# IP address discovery (WALK)
|
|
%SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# IP-MIB ipAddrTable (IPv4 - deprecated but widely supported)
|
|
"1.3.6.1.2.1.4.20",
|
|
# IP-MIB ipAddressTable (RFC 4293 - unified IPv4/IPv6)
|
|
"1.3.6.1.2.1.4.34"
|
|
]
|
|
}
|
|
]
|
|
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
|
|
|
|
[
|
|
# ifHCInOctets (64-bit counter)
|
|
"1.3.6.1.2.1.31.1.1.1.6.#{idx}",
|
|
# ifHCOutOctets (64-bit counter)
|
|
"1.3.6.1.2.1.31.1.1.1.10.#{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)
|
|
|
|
# Build separate WALK queries to isolate failures
|
|
# If one table doesn't exist, others can still succeed
|
|
base_queries = [
|
|
%SnmpQuery{
|
|
query_type: :GET,
|
|
oids: sensor_oids ++ interface_oids
|
|
}
|
|
]
|
|
|
|
# Neighbor discovery (LLDP/CDP) - separate query
|
|
neighbor_query = %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"
|
|
]
|
|
}
|
|
|
|
# ARP table - separate query (may not be supported on all devices)
|
|
arp_query = %SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# IP-MIB ipNetToMediaTable (IPv4 ARP)
|
|
"1.3.6.1.2.1.4.22",
|
|
# IP-MIB ipNetToPhysicalTable (RFC 4293 - unified IPv4/IPv6)
|
|
"1.3.6.1.2.1.4.35"
|
|
]
|
|
}
|
|
|
|
# MAC address table - separate query (only on switches)
|
|
mac_query = %SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# BRIDGE-MIB dot1dTpFdbTable
|
|
"1.3.6.1.2.1.17.4.3"
|
|
]
|
|
}
|
|
|
|
# IP address tables - separate query
|
|
ip_query = %SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# IP-MIB ipAddrTable (IPv4)
|
|
"1.3.6.1.2.1.4.20",
|
|
# IP-MIB ipAddressTable (IPv6)
|
|
"1.3.6.1.2.1.4.34"
|
|
]
|
|
}
|
|
|
|
# HOST-RESOURCES-MIB - separate query (may not be supported)
|
|
host_resources_query = %SnmpQuery{
|
|
query_type: :WALK,
|
|
oids: [
|
|
# hrProcessorTable
|
|
"1.3.6.1.2.1.25.3.3",
|
|
# hrStorageTable
|
|
"1.3.6.1.2.1.25.2.3"
|
|
]
|
|
}
|
|
|
|
base_queries ++ [neighbor_query, arp_query, mac_query, ip_query, host_resources_query]
|
|
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_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 = DateTime.truncate(DateTime.utc_now(), :second)
|
|
|
|
# Convert status string to atom
|
|
status = String.to_existing_atom(check.status)
|
|
|
|
attrs = %{
|
|
device_id: check.device_id,
|
|
agent_token_id: agent_token_id,
|
|
status: status,
|
|
response_time_ms: check.response_time_ms,
|
|
checked_at: checked_at
|
|
}
|
|
|
|
case Monitoring.create_check(attrs) do
|
|
{:ok, _check} ->
|
|
update_device_status_from_check(device, status)
|
|
: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 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}
|
|
)
|
|
end
|
|
|
|
defp handle_agent_status_change(device, _old_status, :down) do
|
|
if !Alerts.has_active_alert?(device.id, :device_down) do
|
|
now = DateTime.truncate(DateTime.utc_now(), :second)
|
|
|
|
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:new",
|
|
{:new_alert, device.id, :device_down}
|
|
)
|
|
end
|
|
end
|
|
|
|
defp handle_agent_status_change(device, _old_status, :up) do
|
|
now = DateTime.truncate(DateTime.utc_now(), :second)
|
|
|
|
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,
|
|
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:resolved",
|
|
{:alert_resolved, device.id, :device_down}
|
|
)
|
|
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)
|
|
|
|
snmp_device = device.snmp_device
|
|
oid_values = normalize_oid_keys(result.oid_values)
|
|
# Use server time to avoid clock drift issues between agents
|
|
timestamp = DateTime.truncate(DateTime.utc_now(), :second)
|
|
|
|
Logger.info(
|
|
"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)
|
|
|
|
# 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
|
|
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
|
|
alias Towerops.Snmp.ArpDiscovery
|
|
alias Towerops.Snmp.MacDiscovery
|
|
alias Towerops.Snmp.NeighborDiscovery
|
|
alias Towerops.Snmp.Profiles.Base
|
|
|
|
# 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
|
|
alias Towerops.Snmp.NeighborDiscovery
|
|
|
|
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
|
|
alias Towerops.Snmp.ArpDiscovery
|
|
|
|
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
|
|
alias Towerops.Snmp.MacDiscovery
|
|
|
|
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
|
|
alias Towerops.Snmp.Profiles.Base
|
|
|
|
case Base.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
|
|
alias Towerops.Snmp.Profiles.Base
|
|
|
|
case Base.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
|
|
alias Towerops.Snmp.Profiles.Base
|
|
|
|
case Base.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 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
|
|
Float.round(parsed / sensor.sensor_divisor, 1)
|
|
else
|
|
_ -> nil
|
|
end
|
|
end
|
|
|
|
defp process_interface_stats_batch(interfaces, oid_values, timestamp) do
|
|
entries =
|
|
Enum.map(interfaces, fn iface ->
|
|
idx = iface.if_index
|
|
|
|
%{
|
|
interface_id: iface.id,
|
|
if_in_octets: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.31.1.1.1.6.#{idx}")),
|
|
if_out_octets: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.31.1.1.1.10.#{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)
|
|
|
|
Snmp.create_interface_stats_batch(entries)
|
|
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
|
|
|
|
# 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) 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
|
|
|
|
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, _backup} ->
|
|
BackupRequests.update_request_status(request.id, "success", nil)
|
|
Logger.info("Backup created", device_id: result.device_id, source: request.source)
|
|
|
|
{:error, reason} ->
|
|
mark_backup_failed(request, result, "Failed to create backup: #{inspect(reason)}")
|
|
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
|
|
|
|
## 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
|