towerops/test/towerops_web/channels/agent_channel_processing_test.exs
Graham McIntire d85b42c3ee fix: replace remaining Process.sleep calls with receive/timeout patterns
Production code (lib/): all Process.sleep calls replaced with receive
after timeout - retry backoff, exponential backoff, batch delays,
poll intervals, and settle waits.

Test code: poll_until helpers across 4 agent channel test files now
use receive after instead of Process.sleep. Various standalone
Process.sleep(N) calls replaced with :sys.get_state sync barriers
or receive after. deferred_discovery test simulation sleeps reduced
500ms -> 30ms.
2026-06-02 15:27:23 -05:00

364 lines
11 KiB
Elixir

defmodule ToweropsWeb.AgentChannelProcessingTest do
@moduledoc """
Additional coverage for ToweropsWeb.AgentChannel — focused on processing
paths (poll/discovery/check/lldp/mikrotik) and lifecycle handle_info events
not exercised in agent_channel_test.exs.
"""
use Towerops.DataCase, async: false
import Phoenix.ChannelTest
alias Towerops.AccountsFixtures
alias Towerops.Agent.AgentHeartbeat
alias Towerops.Agent.AgentJobList
alias Towerops.Agent.MikrotikResult
alias Towerops.Agent.MikrotikSentence
alias Towerops.Agent.SnmpResult
alias Towerops.AgentsFixtures
alias Towerops.DevicesFixtures
alias Towerops.OrganizationsFixtures
alias Towerops.Proto.Encode
alias Towerops.Proto.Types.LldpNeighbor
alias Towerops.Proto.Types.LldpTopologyResult
alias Towerops.Snmp.AgentDiscovery
alias ToweropsWeb.AgentSocket
@endpoint ToweropsWeb.Endpoint
setup do
user = AccountsFixtures.user_fixture()
organization = OrganizationsFixtures.organization_fixture(user.id)
device =
DevicesFixtures.device_fixture(%{
organization_id: organization.id,
snmp_version: "2c",
snmp_community: "public"
})
{:ok, agent_token, token_string} =
AgentsFixtures.agent_token_fixture(organization.id)
{:ok, _assignment} =
AgentsFixtures.agent_assignment_fixture(agent_token.id, device.id)
{:ok, socket} = connect(AgentSocket, %{})
{:ok, _, socket} =
subscribe_and_join(socket, "agent:#{agent_token.id}", %{"token" => token_string})
assert_push "jobs", _initial_jobs
%{
socket: socket,
device: device,
agent_token: agent_token,
token_string: token_string,
organization: organization
}
end
defp encode_payload(struct) do
binary = struct.__struct__.encode(struct)
%{"binary" => Base.encode64(binary)}
end
defp poll_until(fun, opts \\ []) do
max_attempts = Keyword.get(opts, :max_attempts, 100)
delay_ms = Keyword.get(opts, :delay_ms, 20)
Enum.reduce_while(1..max_attempts, nil, fn _, _ ->
case fun.() do
nil ->
receive do
after
delay_ms -> :ok
end
{:cont, nil}
false ->
receive do
after
delay_ms -> :ok
end
{:cont, nil}
result ->
{:halt, result}
end
end) || fun.()
end
defp build_heartbeat do
%AgentHeartbeat{
version: "1.0.0",
hostname: "test-agent",
uptime_seconds: 3600,
ip_address: "",
arch: "x86_64"
}
end
# ── Discovery error path ──────────────────────────────────────────
describe "discovery error path" do
test "discovery with empty oid_values does not update last_discovery_at", %{
socket: socket,
device: device
} do
assert is_nil(Towerops.Devices.get_device(device.id).last_discovery_at)
result = %SnmpResult{
device_id: device.id,
job_type: :DISCOVER,
job_id: "discover:#{device.id}",
timestamp: DateTime.to_unix(DateTime.utc_now()),
oid_values: %{}
}
push(socket, "result", encode_payload(result))
ref = push(socket, "heartbeat", encode_payload(build_heartbeat()))
refute_reply ref, :error
assert is_nil(Towerops.Devices.get_device(device.id).last_discovery_at)
end
end
# ── Polling status transitions (process_polling_result) ─────────
describe "process_polling_result status transitions" do
setup %{device: device} do
oid_values = %{
"1.3.6.1.2.1.1.1.0" => "Test Device",
"1.3.6.1.2.1.1.2.0" => "1.3.6.1.4.1.9.1.1",
"1.3.6.1.2.1.1.3.0" => "123456",
"1.3.6.1.2.1.1.4.0" => "admin@test.com",
"1.3.6.1.2.1.1.5.0" => "test-device",
"1.3.6.1.2.1.1.6.0" => "Test Location"
}
{:ok, _discovered} = AgentDiscovery.process_agent_discovery(device, oid_values)
device_with_details = Towerops.Devices.get_device_with_details(device.id)
%{device: device_with_details}
end
test "successful poll transitions device from down to up and broadcasts", %{
socket: socket,
device: device
} do
{:ok, device} = Towerops.Devices.update_device_status(device, :down)
Phoenix.PubSub.subscribe(Towerops.PubSub, "device:#{device.id}")
result = %SnmpResult{
device_id: device.id,
job_type: :POLL,
job_id: "poll:#{device.id}",
timestamp: DateTime.to_unix(DateTime.utc_now()),
oid_values: %{}
}
push(socket, "result", encode_payload(result))
assert_receive {:device_status_changed, _, :up, _}, 1_000
assert poll_until(fn ->
d = Towerops.Devices.get_device(device.id)
if d && d.status == :up, do: d
end)
end
test "poll for already-up device does not re-broadcast status change", %{
socket: socket,
device: device
} do
{:ok, _device} = Towerops.Devices.update_device_status(device, :up)
Phoenix.PubSub.subscribe(Towerops.PubSub, "device:#{device.id}")
result = %SnmpResult{
device_id: device.id,
job_type: :POLL,
job_id: "poll:#{device.id}",
timestamp: DateTime.to_unix(DateTime.utc_now()),
oid_values: %{}
}
push(socket, "result", encode_payload(result))
refute_receive {:device_status_changed, _, _, _}, 100
assert_receive {:state_sensors_updated, _}, 1_000
end
end
# ── handle_in lldp_topology_result ───────────────────────────────
describe "handle_in lldp_topology_result" do
test "stores LLDP neighbors from topology result", %{socket: socket, device: device} do
result = %LldpTopologyResult{
device_id: device.id,
job_id: Ecto.UUID.generate(),
local_system_name: "test-device",
neighbors: [
%LldpNeighbor{
neighbor_name: "switch-A",
local_port: "Eth1/0/1",
remote_port: "Gi0/1",
remote_port_id: "Gi0/1",
management_addresses: ["10.0.0.1"]
},
%LldpNeighbor{
neighbor_name: "switch-B",
local_port: "Eth1/0/2",
remote_port: "Gi0/2",
remote_port_id: "Gi0/2",
management_addresses: ["10.0.0.2"]
}
],
timestamp: DateTime.to_unix(DateTime.utc_now())
}
binary = Encode.encode_lldp_topology_result(result)
push(socket, "lldp_topology_result", %{"binary" => Base.encode64(binary)})
neighbors =
poll_until(
fn ->
n = Towerops.Topology.list_lldp_neighbors(device.id)
if length(n) == 2, do: n
end,
max_attempts: 100,
delay_ms: 20
)
assert length(neighbors) == 2
assert Enum.any?(neighbors, &(&1.neighbor_name == "switch-A"))
assert Enum.any?(neighbors, &(&1.neighbor_name == "switch-B"))
end
test "empty neighbor list does not crash", %{socket: socket, device: device} do
result = %LldpTopologyResult{
device_id: device.id,
job_id: Ecto.UUID.generate(),
local_system_name: "test-device",
neighbors: [],
timestamp: DateTime.to_unix(DateTime.utc_now())
}
binary = Encode.encode_lldp_topology_result(result)
push(socket, "lldp_topology_result", %{"binary" => Base.encode64(binary)})
ref = push(socket, "heartbeat", encode_payload(build_heartbeat()))
refute_reply ref, :error
end
test "invalid base64 lldp_topology_result is handled gracefully", %{socket: socket} do
push(socket, "lldp_topology_result", %{"binary" => "not-valid-base64!!!"})
ref = push(socket, "heartbeat", encode_payload(build_heartbeat()))
refute_reply ref, :error
end
test "invalid protobuf lldp_topology_result is handled gracefully", %{socket: socket} do
push(socket, "lldp_topology_result", %{"binary" => Base.encode64("not-a-protobuf")})
ref = push(socket, "heartbeat", encode_payload(build_heartbeat()))
refute_reply ref, :error
end
end
# ── mikrotik_result non-backup path ─────────────────────────────
describe "mikrotik_result non-backup job_id" do
test "non-backup job_id is processed without crash", %{socket: socket, device: device} do
result = %MikrotikResult{
device_id: device.id,
job_id: Ecto.UUID.generate(),
sentences: [%MikrotikSentence{attributes: %{"name" => "router1"}}],
error: "",
timestamp: DateTime.to_unix(DateTime.utc_now())
}
push(socket, "mikrotik_result", encode_payload(result))
ref = push(socket, "heartbeat", encode_payload(build_heartbeat()))
refute_reply ref, :error
end
end
# ── backup_requested via PubSub ─────────────────────────────────
describe "backup_requested via PubSub" do
test "PubSub broadcast routes backup job to channel", %{
agent_token: agent_token,
device: device
} do
job_id = "backup:#{device.id}"
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"agent:#{agent_token.id}:backup",
{:backup_requested, device.id, job_id}
)
assert_push "backup_job", %{binary: jobs_binary}, 200
{:ok, decoded} = Base.decode64(jobs_binary)
{:ok, job_list} = AgentJobList.decode(decoded)
assert hd(job_list.jobs).job_id == job_id
end
end
# ── latency_probe_jobs handle_info ──────────────────────────────
describe "latency_probe_jobs handle_info" do
test "sends 5 PING jobs per device for cloud poller", %{
socket: socket,
device: device
} do
send(socket.channel_pid, {:latency_probe_jobs, [device]})
assert_push "jobs", %{binary: jobs_binary}, 500
{:ok, decoded} = Base.decode64(jobs_binary)
{:ok, job_list} = AgentJobList.decode(decoded)
ping_jobs = Enum.filter(job_list.jobs, &(&1.job_type == :PING))
assert length(ping_jobs) == 5
assert Enum.all?(ping_jobs, &(&1.device_id == device.id))
end
test "empty device list does not push", %{socket: socket} do
send(socket.channel_pid, {:latency_probe_jobs, []})
refute_push "jobs", %{}, 100
end
end
# ── checks_changed handle_info ──────────────────────────────────
describe "checks_changed handle_info" do
test "rebuilds and pushes check jobs", %{
socket: socket,
organization: organization,
agent_token: agent_token
} do
{:ok, _check} =
Towerops.Monitoring.create_check(%{
organization_id: organization.id,
agent_token_id: agent_token.id,
check_type: "tcp",
name: "TCP Check",
interval_seconds: 60,
timeout_ms: 5000,
enabled: true,
config: %{"host" => "example.com", "port" => 443}
})
send(socket.channel_pid, :checks_changed)
assert_push "check_jobs", %{binary: _binary}, 500
end
end
end