feat: add process_device/2 topology pipeline

Orchestrates evidence collection from LLDP, CDP, MAC, and ARP
data. Groups evidence by remote device, resolves to managed
devices, upserts links with merged confidence. Broadcasts
topology changes via PubSub. Skips self-links.
This commit is contained in:
Graham McIntire 2026-02-12 13:02:14 -06:00
parent 9327fd3f13
commit 8e355fce79
No known key found for this signature in database
2 changed files with 234 additions and 0 deletions

View file

@ -15,6 +15,8 @@ defmodule Towerops.Topology do
alias Towerops.Topology.DeviceLink
alias Towerops.Topology.DeviceLinkEvidence
require Logger
@doc """
Builds lookup maps for matching evidence to managed devices in an organization.
@ -274,4 +276,111 @@ defmodule Towerops.Topology do
{:ok, Repo.reload!(link)}
end
end
@doc """
Main topology inference pipeline for a device. Called after each poll cycle.
Collects evidence from all sources, groups by remote identifier, upserts links,
broadcasts changes. Returns {:ok, :changed} or {:ok, :unchanged}.
"""
def process_device(device, organization_id) do
lookup = build_device_lookup(organization_id)
now = DateTime.utc_now()
# Collect all evidence
all_evidence =
collect_lldp_evidence(device.id) ++
collect_cdp_evidence(device.id) ++
collect_mac_evidence(device.id, lookup) ++
collect_arp_evidence(device.id, lookup)
if Enum.empty?(all_evidence) do
{:ok, :unchanged}
else
grouped = group_evidence_by_remote(all_evidence)
results =
grouped
|> Enum.map(&upsert_grouped_evidence(&1, device.id, lookup, now))
|> Enum.reject(&(&1 == :skip))
has_changes = Enum.any?(results, &match?({:ok, _}, &1))
if has_changes do
Phoenix.PubSub.broadcast(
Towerops.PubSub,
"topology:#{organization_id}",
{:topology_updated, organization_id}
)
{:ok, :changed}
else
{:ok, :unchanged}
end
end
end
@doc "List all device links where this device is the source, ordered by confidence."
def list_device_links(device_id) do
Repo.all(
from(l in DeviceLink,
where: l.source_device_id == ^device_id,
order_by: [desc: l.confidence]
)
)
end
defp upsert_grouped_evidence({_key, evidence_list}, device_id, lookup, now) do
best = Enum.max_by(evidence_list, & &1.confidence)
matched_id =
resolve_device(lookup, %{
mac: best.remote_mac,
ip: best.remote_ip,
name: best.remote_name
})
if matched_id == device_id do
:skip
else
upsert_link(%{
source_device_id: device_id,
target_device_id: matched_id,
source_interface_id: best.source_interface_id,
link_type: evidence_type_to_link_type(best.evidence_type),
confidence: combine_evidence_confidence(evidence_list),
discovered_remote_mac: best.remote_mac,
discovered_remote_ip: best.remote_ip,
discovered_remote_name: best.remote_name,
metadata: %{},
last_confirmed_at: now,
evidence:
Enum.map(evidence_list, fn ev ->
%{
evidence_type: ev.evidence_type,
evidence_data: ev.evidence_data,
observed_at: now
}
end)
})
end
end
defp group_evidence_by_remote(evidence_list) do
Enum.group_by(evidence_list, fn ev ->
ev.remote_mac || ev.remote_ip || ev.remote_name || "unknown"
end)
end
defp combine_evidence_confidence(evidence_list) do
evidence_list
|> Enum.map(& &1.confidence)
|> Enum.reduce(0.0, &merge_confidence/2)
end
defp evidence_type_to_link_type("lldp_neighbor"), do: "lldp"
defp evidence_type_to_link_type("cdp_neighbor"), do: "cdp"
defp evidence_type_to_link_type("mac_on_interface"), do: "mac_match"
defp evidence_type_to_link_type("arp_entry"), do: "arp_inference"
defp evidence_type_to_link_type("wireless_registration"), do: "wireless_association"
defp evidence_type_to_link_type(_), do: "mac_match"
end

View file

@ -414,6 +414,131 @@ defmodule Towerops.TopologyTest do
end
end
describe "process_device/2" do
test "creates links from LLDP neighbors and resolves target device",
%{router: router, switch: switch, router_iface: iface, organization: org} do
# Create SNMP device for switch so its MAC is resolvable
switch_snmp =
%Device{}
|> Device.changeset(%{device_id: switch.id, sys_name: "main-switch"})
|> Repo.insert!()
%Interface{}
|> Interface.changeset(%{
snmp_device_id: switch_snmp.id,
if_index: 1,
if_name: "ge-0/0/1",
if_phys_address: "bb:cc:dd:00:11:22"
})
|> Repo.insert!()
# Create LLDP neighbor on router pointing to switch's MAC
{:ok, _} =
Towerops.Snmp.upsert_neighbor(%{
device_id: router.id,
interface_id: iface.id,
protocol: "lldp",
remote_chassis_id: "bb:cc:dd:00:11:22",
remote_system_name: "Main-Switch",
remote_address: "10.0.0.2",
remote_port_id: "ge-0/0/1",
remote_capabilities: ["switch", "bridge"],
last_discovered_at: DateTime.utc_now()
})
assert {:ok, :changed} = Topology.process_device(router, org.id)
links = Topology.list_device_links(router.id)
assert length(links) == 1
[link] = links
assert link.source_device_id == router.id
assert link.target_device_id == switch.id
assert link.link_type == "lldp"
assert_in_delta link.confidence, 0.95, 0.01
end
test "creates discovered link when remote device is not managed",
%{router: router, router_iface: iface, organization: org} do
{:ok, _} =
Towerops.Snmp.upsert_neighbor(%{
device_id: router.id,
interface_id: iface.id,
protocol: "lldp",
remote_chassis_id: "ff:ff:ff:00:00:01",
remote_system_name: "unknown-device",
remote_address: "10.99.99.99",
remote_port_id: "eth0",
remote_capabilities: [],
last_discovered_at: DateTime.utc_now()
})
assert {:ok, :changed} = Topology.process_device(router, org.id)
links = Topology.list_device_links(router.id)
assert length(links) == 1
[link] = links
assert link.target_device_id == nil
assert link.discovered_remote_mac == "ff:ff:ff:00:00:01"
assert link.discovered_remote_name == "unknown-device"
assert link.discovered_remote_ip == "10.99.99.99"
end
test "returns :unchanged when no evidence exists",
%{router: router, organization: org} do
assert {:ok, :unchanged} = Topology.process_device(router, org.id)
end
test "does not create self-links",
%{router: router, router_iface: iface, organization: org} do
# LLDP neighbor pointing back to the same device's MAC
{:ok, _} =
Towerops.Snmp.upsert_neighbor(%{
device_id: router.id,
interface_id: iface.id,
protocol: "lldp",
remote_chassis_id: "aa:bb:cc:00:11:22",
remote_system_name: "Core-Router",
remote_address: "10.0.0.1",
remote_port_id: "eth1",
remote_capabilities: ["router"],
last_discovered_at: DateTime.utc_now()
})
assert {:ok, :unchanged} = Topology.process_device(router, org.id)
assert Topology.list_device_links(router.id) == []
end
end
describe "list_device_links/1" do
test "returns links ordered by confidence desc", %{router: router, router_iface: iface} do
now = DateTime.utc_now()
{:ok, _} =
Topology.upsert_link(%{
source_device_id: router.id,
source_interface_id: iface.id,
link_type: "lldp",
confidence: 0.95,
discovered_remote_mac: "aa:00:00:00:00:01",
last_confirmed_at: now
})
{:ok, _} =
Topology.upsert_link(%{
source_device_id: router.id,
source_interface_id: iface.id,
link_type: "arp_inference",
confidence: 0.6,
discovered_remote_mac: "aa:00:00:00:00:02",
last_confirmed_at: now
})
links = Topology.list_device_links(router.id)
assert length(links) == 2
assert hd(links).confidence > List.last(links).confidence
end
end
describe "collect_arp_evidence/2" do
test "returns evidence for ARP entries matching known devices", %{
organization: org,