towerops/lib/towerops_web/channels/agent_channel.ex
Graham McIntire dc32e16c68
fix: allow cloud pollers to access devices from any organization
Cloud pollers (is_cloud_poller = true, organization_id = nil) were
being blocked from polling devices because verify_device_organization
was checking if the device's organization matched the agent's
organization.

Since cloud pollers don't belong to any organization (organization_id
is nil), they should be allowed to poll devices from any organization.

Changes:
- verify_device_organization now allows access if organization_id is nil
- This fixes "Device not in agent's organization" errors for cloud pollers

Error before:
  Device b7e8dd99-0b19-496d-97ae-9f2601d16464 not in agent's organization
2026-01-25 11:15:19 -06:00

481 lines
14 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.MonitoringCheck
alias Towerops.Agent.SnmpDevice
alias Towerops.Agent.SnmpQuery
alias Towerops.Agent.SnmpResult
alias Towerops.Agents
alias Towerops.Devices
alias Towerops.Monitoring
alias Towerops.Snmp
alias Towerops.Snmp.Discovery
alias Towerops.Workers.DiscoveryWorker
require Logger
@impl true
@spec join(String.t(), map(), Phoenix.Socket.t()) ::
{:ok, Phoenix.Socket.t()} | {:error, map()}
def join("agent:" <> _agent_id, %{"token" => token}, socket) do
# Verify agent token from join payload
case Agents.verify_agent_token(token) do
{:ok, agent_token} ->
socket =
socket
|> assign(:agent_token_id, agent_token.id)
|> assign(:organization_id, agent_token.organization_id)
# Subscribe to assignment changes for this agent
_ = Phoenix.PubSub.subscribe(Towerops.PubSub, "agent:#{agent_token.id}:assignments")
# Update last_seen_at and IP on join
remote_ip = get_remote_ip(socket)
_ = Agents.update_agent_token_heartbeat(agent_token.id, remote_ip, %{})
# Send initial job list
send(self(), :send_jobs)
{:ok, socket}
{:error, _} ->
{:error, %{reason: "unauthorized"}}
end
end
def join("agent:" <> _agent_id, _payload, _socket) do
{:error, %{reason: "missing token"}}
end
@impl true
@spec handle_info(:send_jobs, Phoenix.Socket.t()) :: {:noreply, Phoenix.Socket.t()}
def handle_info(:send_jobs, socket) do
jobs = build_jobs_for_agent(socket.assigns.agent_token_id)
job_list = %AgentJobList{jobs: jobs}
binary = AgentJobList.encode(job_list)
push(socket, "jobs", %{binary: Base.encode64(binary)})
{:noreply, socket}
end
# Handle PubSub broadcast when device assignments change
def handle_info({:assignments_changed, _event}, socket) do
Logger.info("Agent assignments changed, sending updated job list",
agent_token_id: socket.assigns.agent_token_id
)
send(self(), :send_jobs)
{:noreply, socket}
end
@impl true
@spec handle_in(String.t(), map(), Phoenix.Socket.t()) :: {:noreply, Phoenix.Socket.t()}
def handle_in("result", %{"binary" => binary_b64}, socket) do
binary = Base.decode64!(binary_b64)
result = SnmpResult.decode(binary)
_ = process_snmp_result(socket.assigns.organization_id, result)
{:noreply, socket}
end
def handle_in("heartbeat", %{"binary" => binary_b64}, socket) do
binary = Base.decode64!(binary_b64)
heartbeat = AgentHeartbeat.decode(binary)
metadata = %{
"version" => heartbeat.version,
"hostname" => heartbeat.hostname,
"uptime_seconds" => heartbeat.uptime_seconds
}
_ =
Agents.update_agent_token_heartbeat(
socket.assigns.agent_token_id,
heartbeat.ip_address,
metadata
)
{:noreply, socket}
end
@impl true
def handle_in("error", %{"binary" => binary_b64}, socket) do
binary = Base.decode64!(binary_b64)
error = AgentError.decode(binary)
Logger.error("Agent job error",
agent_token_id: socket.assigns.agent_token_id,
device_id: error.device_id,
error: error.message
)
{:noreply, socket}
end
def handle_in("monitoring_check", %{"binary" => binary_b64}, socket) do
binary = Base.decode64!(binary_b64)
check = MonitoringCheck.decode(binary)
_ =
process_monitoring_check(
socket.assigns.organization_id,
socket.assigns.agent_token_id,
check
)
{:noreply, socket}
end
# Private helpers
@spec build_jobs_for_agent(Ecto.UUID.t()) :: [AgentJob.t()]
defp build_jobs_for_agent(agent_token_id) do
agent_token_id
|> Agents.list_agent_polling_targets()
|> Enum.map(&build_job_for_device/1)
end
defp build_job_for_device(device) do
if needs_discovery?(device) do
build_discovery_job(device)
else
build_polling_job(device)
end
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 build_discovery_job(device) do
%AgentJob{
job_id: "discover:#{device.id}",
job_type: :DISCOVER,
device_id: device.id,
snmp_device: %SnmpDevice{
ip: device.ip_address,
community: device.snmp_community,
version: device.snmp_version,
port: device.snmp_port || 161,
monitoring_enabled: device.monitoring_enabled || false,
check_interval_seconds: device.check_interval_seconds || 60
},
queries: build_discovery_queries()
}
end
defp build_polling_job(device) do
_snmp_device = device.snmp_device
%AgentJob{
job_id: "poll:#{device.id}",
job_type: :POLL,
device_id: device.id,
snmp_device: %SnmpDevice{
ip: device.ip_address,
community: device.snmp_community,
version: device.snmp_version,
port: device.snmp_port || 161,
monitoring_enabled: device.monitoring_enabled || false,
check_interval_seconds: device.check_interval_seconds || 60
},
queries: build_polling_queries(device)
}
end
defp build_discovery_queries do
[
# System info (GET)
%SnmpQuery{
query_type: :GET,
oids: [
# sysDescr
"1.3.6.1.2.1.1.1.0",
# sysObjectID
"1.3.6.1.2.1.1.2.0",
# sysUpTime
"1.3.6.1.2.1.1.3.0",
# sysContact
"1.3.6.1.2.1.1.4.0",
# sysName
"1.3.6.1.2.1.1.5.0",
# sysLocation
"1.3.6.1.2.1.1.6.0"
]
},
# Interface table (WALK)
%SnmpQuery{
query_type: :WALK,
# ifTable
oids: ["1.3.6.1.2.1.2.2.1"]
},
# Interface extensions (WALK)
%SnmpQuery{
query_type: :WALK,
# ifXTable
oids: ["1.3.6.1.2.1.31.1.1.1"]
},
# Vendor-specific MIBs (WALK)
%SnmpQuery{
query_type: :WALK,
oids: [
# Cisco
"1.3.6.1.4.1.9",
# MikroTik
"1.3.6.1.4.1.14988",
# Ubiquiti
"1.3.6.1.4.1.41112"
]
},
# Neighbor discovery (WALK)
%SnmpQuery{
query_type: :WALK,
oids: [
# LLDP
"1.0.8802.1.1.2.1.4.1.1",
# Cisco CDP
"1.3.6.1.4.1.9.9.23"
]
}
]
end
defp build_polling_queries(device) do
sensor_oids = Enum.map(device.snmp_device.sensors, & &1.sensor_oid)
interface_oids =
Enum.flat_map(device.snmp_device.interfaces, fn iface ->
idx = iface.if_index
[
# ifInOctets
"1.3.6.1.2.1.2.2.1.10.#{idx}",
# ifOutOctets
"1.3.6.1.2.1.2.2.1.16.#{idx}",
# ifInErrors
"1.3.6.1.2.1.2.2.1.14.#{idx}",
# ifOutErrors
"1.3.6.1.2.1.2.2.1.20.#{idx}",
# ifInDiscards
"1.3.6.1.2.1.2.2.1.13.#{idx}",
# ifOutDiscards
"1.3.6.1.2.1.2.2.1.19.#{idx}"
]
end)
# Neighbor discovery OIDs
neighbor_oids = [
# LLDP
"1.0.8802.1.1.2.1.4.1.1",
# Cisco CDP
"1.3.6.1.4.1.9.9.23"
]
[
%SnmpQuery{
query_type: :GET,
oids: sensor_oids ++ interface_oids
},
%SnmpQuery{
query_type: :WALK,
oids: neighbor_oids
}
]
end
defp process_snmp_result(organization_id, result) do
with {:ok, device} <- fetch_device(result.device_id),
:ok <- verify_device_organization(device, organization_id) do
process_job_result(device, result)
else
{:error, :device_not_found} ->
Logger.error("Device not found: #{result.device_id}")
{:error, :wrong_organization} ->
Logger.error("Device #{result.device_id} not in agent's organization")
end
end
defp fetch_device(device_id) do
case Devices.get_device_with_details(device_id) do
nil -> {:error, :device_not_found}
device -> {:ok, device}
end
end
defp verify_device_organization(device, organization_id) do
# Cloud pollers (organization_id = nil) can access any device
if is_nil(organization_id) or device.site.organization_id == organization_id do
:ok
else
{:error, :wrong_organization}
end
end
defp process_job_result(device, %{job_type: :DISCOVER} = result) do
process_discovery_result(device, result)
end
defp process_job_result(device, %{job_type: :POLL} = result) do
process_polling_result(device, result)
end
defp process_discovery_result(device, _result) do
# Agent confirmed device is reachable, trigger full discovery from server
# This hybrid approach simplifies the initial implementation:
# - Agent acts as remote executor
# - Server does SNMP queries and parsing using existing logic
Logger.info("Discovery results received for #{device.name}, triggering full discovery")
enqueue_discovery(device.id)
end
defp process_polling_result(device, result) do
snmp_device = device.snmp_device
oid_values = Map.new(result.oid_values)
timestamp = DateTime.from_unix!(result.timestamp, :second)
# Process sensor readings
Enum.each(snmp_device.sensors, fn sensor ->
if value = Map.get(oid_values, sensor.sensor_oid) do
Snmp.create_sensor_reading(%{
sensor_id: sensor.id,
value: parse_float(value) / sensor.sensor_divisor,
status: "ok",
checked_at: timestamp
})
end
end)
# Process interface stats
Enum.each(device.snmp_device.interfaces, fn iface ->
idx = iface.if_index
Snmp.create_interface_stat(%{
interface_id: iface.id,
if_in_octets: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.10.#{idx}")),
if_out_octets: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.16.#{idx}")),
if_in_errors: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.14.#{idx}")),
if_out_errors: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.20.#{idx}")),
if_in_discards: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.13.#{idx}")),
if_out_discards: parse_integer(Map.get(oid_values, "1.3.6.1.2.1.2.2.1.19.#{idx}")),
checked_at: timestamp
})
end)
# Note: Neighbor discovery is handled by separate polling worker
# since it requires complex LLDP/CDP parsing
end
defp process_monitoring_check(organization_id, agent_token_id, check) do
with {:ok, device} <- fetch_device(check.device_id),
:ok <- verify_device_organization(device, organization_id) do
# Parse status from protobuf string ("success" or "failure")
status =
case check.status do
"success" -> :success
"failure" -> :failure
_ -> :failure
end
timestamp = DateTime.from_unix!(check.timestamp, :second)
# Create monitoring check record with agent tracking
Monitoring.create_check(%{
device_id: device.id,
agent_token_id: agent_token_id,
status: status,
response_time_ms: check.response_time_ms,
checked_at: timestamp
})
else
{:error, :device_not_found} ->
Logger.error("Device not found for monitoring check: #{check.device_id}")
{:error, :wrong_organization} ->
Logger.error("Device #{check.device_id} not in agent's organization")
end
end
# Enqueue discovery job - safe to call in test environment
defp enqueue_discovery(device_id) do
if Application.get_env(:towerops, :env) == :test do
# In test, run synchronously
Task.start(fn -> Discovery.discover_device(Devices.get_device!(device_id)) end)
else
# In dev/prod, enqueue to Oban
DiscoveryWorker.enqueue(device_id)
end
end
defp parse_integer(nil), do: nil
defp parse_integer(value) when is_integer(value), do: value
defp parse_integer(value) when is_binary(value) do
case Integer.parse(value) do
{int, _} -> int
:error -> nil
end
end
defp parse_integer(_), do: nil
defp parse_float(value) when is_float(value), do: value
defp parse_float(value) when is_integer(value), do: value / 1.0
defp parse_float(value) when is_binary(value) do
case Float.parse(value) do
{float, _} -> float
:error -> nil
end
end
defp parse_float(_), do: nil
defp get_remote_ip(socket) do
case socket.transport_pid do
nil ->
nil
pid when is_pid(pid) ->
case :ranch.get_addr(pid) do
{ip, _port} when is_tuple(ip) -> format_ip(ip)
_ -> nil
end
end
rescue
_ -> nil
end
defp format_ip({a, b, c, d}), do: "#{a}.#{b}.#{c}.#{d}"
defp format_ip({a, b, c, d, e, f, g, h}),
do:
"#{Integer.to_string(a, 16)}:#{Integer.to_string(b, 16)}:#{Integer.to_string(c, 16)}:#{Integer.to_string(d, 16)}:#{Integer.to_string(e, 16)}:#{Integer.to_string(f, 16)}:#{Integer.to_string(g, 16)}:#{Integer.to_string(h, 16)}"
end