From 8e355fce7997c00659245a3ed59e4ed7535d4335 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Thu, 12 Feb 2026 13:02:14 -0600 Subject: [PATCH] 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. --- lib/towerops/topology.ex | 109 ++++++++++++++++++++++++++++ test/towerops/topology_test.exs | 125 ++++++++++++++++++++++++++++++++ 2 files changed, 234 insertions(+) diff --git a/lib/towerops/topology.ex b/lib/towerops/topology.ex index 845d478c..893cab36 100644 --- a/lib/towerops/topology.ex +++ b/lib/towerops/topology.ex @@ -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 diff --git a/test/towerops/topology_test.exs b/test/towerops/topology_test.exs index 4e5dca5a..08bf6084 100644 --- a/test/towerops/topology_test.exs +++ b/test/towerops/topology_test.exs @@ -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,