diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index 173daa95..3498c83d 120000 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -1 +1 @@ -/nix/store/zx636v118rzzqzwafgrw2aybh7l0szrl-pre-commit-config.json \ No newline at end of file +/nix/store/4xlkmdd947h5qlhj2rmbyavq08grkz27-pre-commit-config.json \ No newline at end of file diff --git a/lib/towerops/capacity.ex b/lib/towerops/capacity.ex index 2ed65de9..7c219c97 100644 --- a/lib/towerops/capacity.ex +++ b/lib/towerops/capacity.ex @@ -147,29 +147,43 @@ defmodule Towerops.Capacity do {in_bps, out_bps} end) - if Enum.empty?(pairs) do + # Filter out zero-bps pairs (from clamped negative deltas / anomalies) + valid_pairs = Enum.reject(pairs, fn {in_bps, out_bps} -> in_bps == 0.0 and out_bps == 0.0 end) + + if Enum.empty?(valid_pairs) do zero_throughput() else - avg_in = Enum.sum(Enum.map(pairs, &elem(&1, 0))) / length(pairs) - avg_out = Enum.sum(Enum.map(pairs, &elem(&1, 1))) / length(pairs) + in_values = Enum.map(valid_pairs, &elem(&1, 0)) + out_values = Enum.map(valid_pairs, &elem(&1, 1)) + + avg_in = Enum.sum(in_values) / length(in_values) + avg_out = Enum.sum(out_values) / length(out_values) + + # Peak values for capacity planning (95th percentile approximation) + peak_in = percentile(in_values, 95) + peak_out = percentile(out_values, 95) %{ in_bps: Float.round(avg_in, 2), out_bps: Float.round(avg_out, 2), - max_bps: Float.round(max(avg_in, avg_out), 2) + max_bps: Float.round(max(avg_in, avg_out), 2), + peak_in_bps: Float.round(peak_in, 2), + peak_out_bps: Float.round(peak_out, 2), + peak_bps: Float.round(max(peak_in, peak_out), 2) } end end + # 400 Gbps — values above this are counter anomalies, not real traffic + @max_reasonable_bps 400_000_000_000 + @doc """ Calculates bits per second from two consecutive octet counter readings. - Handles 32-bit SNMP counter wraps (ifInOctets/ifOutOctets wrap at 2^32) - and 64-bit HC counter resets (device reboot). Always tries a 2^32 - correction when the delta is negative — this works because: - - - 32-bit wrap: corrected delta is positive and correct - - 64-bit reset: corrected delta is still very negative, clamped to 0 + Negative deltas are clamped to zero (matching LibreNMS behavior). A 32-bit + counter wrap is indistinguishable from an HC counter reset, so we discard + the ambiguous data point rather than guess at a correction. Values exceeding + 400 Gbps are also discarded as counter anomalies. """ def calculate_bps(nil, _, _), do: 0.0 def calculate_bps(_, nil, _), do: 0.0 @@ -177,20 +191,29 @@ defmodule Towerops.Capacity do def calculate_bps(prev_octets, curr_octets, time_diff) do delta = curr_octets - prev_octets - delta = - if delta < 0 do - # Try 32-bit counter wrap correction. For 64-bit HC counter resets - # (e.g. 10 TB → 1 MB), the corrected delta stays negative and gets - # clamped to 0 below. - delta + 4_294_967_296 - else - delta - end + # Clamp negative deltas to zero — counter wraps and HC resets both produce + # negative deltas and are indistinguishable, so we discard the data point. + bps = max(0, delta) * 8 / time_diff - max(0.0, delta * 8 / time_diff) + # Discard values above 400 Gbps as counter anomalies + if bps > @max_reasonable_bps, do: 0.0, else: bps end - defp zero_throughput, do: %{in_bps: 0.0, out_bps: 0.0, max_bps: 0.0} + defp zero_throughput do + %{in_bps: 0.0, out_bps: 0.0, max_bps: 0.0, peak_in_bps: 0.0, peak_out_bps: 0.0, peak_bps: 0.0} + end + + @doc """ + Calculates the nth percentile of a list of numbers. + Uses nearest-rank method. + """ + def percentile([], _n), do: 0.0 + + def percentile(values, n) when n >= 0 and n <= 100 do + sorted = Enum.sort(values) + rank = Float.ceil(n / 100 * length(sorted)) |> trunc() |> max(1) + Enum.at(sorted, rank - 1) + end defp utilization_status(pct) when pct >= 90, do: :red defp utilization_status(pct) when pct >= 70, do: :yellow diff --git a/lib/towerops/snmp/interface_stat.ex b/lib/towerops/snmp/interface_stat.ex index c28c2a2e..88e749a6 100644 --- a/lib/towerops/snmp/interface_stat.ex +++ b/lib/towerops/snmp/interface_stat.ex @@ -21,6 +21,7 @@ defmodule Towerops.Snmp.InterfaceStat do field :if_in_discards, :integer field :if_out_discards, :integer field :checked_at, :utc_datetime + field :is_hc, :boolean, default: true belongs_to :interface, Interface @@ -36,6 +37,7 @@ defmodule Towerops.Snmp.InterfaceStat do if_in_discards: integer() | nil, if_out_discards: integer() | nil, checked_at: DateTime.t(), + is_hc: boolean(), interface_id: Ecto.UUID.t(), interface: Ecto.Association.NotLoaded.t() | Interface.t(), inserted_at: DateTime.t() @@ -52,7 +54,8 @@ defmodule Towerops.Snmp.InterfaceStat do :if_out_errors, :if_in_discards, :if_out_discards, - :checked_at + :checked_at, + :is_hc ]) |> validate_required([:interface_id, :checked_at]) |> foreign_key_constraint(:interface_id) diff --git a/lib/towerops/workers/device_poller_worker.ex b/lib/towerops/workers/device_poller_worker.ex index 14f4be9c..48828e71 100644 --- a/lib/towerops/workers/device_poller_worker.ex +++ b/lib/towerops/workers/device_poller_worker.ex @@ -298,7 +298,7 @@ defmodule Towerops.Workers.DevicePollerWorker do end defp poll_device_interfaces(device, snmp_device, client_opts, now) do - poll_interfaces(snmp_device.interfaces, client_opts, now) + poll_interfaces(snmp_device.interfaces, client_opts, now, snmp_device) Logger.debug("Polled #{length(snmp_device.interfaces)} interfaces for #{device.name}") _ = @@ -1003,12 +1003,21 @@ defmodule Towerops.Workers.DevicePollerWorker do end end - defp poll_interfaces(interfaces, client_opts, timestamp) do + defp poll_interfaces(interfaces, client_opts, timestamp, snmp_device) do + # Detect if this is an AirFiber device that needs proprietary counter OIDs + af_overrides = detect_airfiber_counter_overrides(snmp_device, client_opts) + entries = interfaces |> Task.async_stream( fn interface -> - {:ok, stat_data} = get_interface_stats(client_opts, interface.if_index) + stat_data = + case get_airfiber_stats(af_overrides, interface, client_opts) do + {:ok, data} -> data + _ -> + {:ok, data} = get_interface_stats(client_opts, interface.if_index) + data + end Map.merge(stat_data, %{ interface_id: interface.id, @@ -1027,6 +1036,70 @@ defmodule Towerops.Workers.DevicePollerWorker do Snmp.create_interface_stats_batch(entries) end + # Detect AirFiber devices and return the appropriate proprietary OIDs + # for traffic counters. LibreNMS uses UBNT-AirFIBER-MIB (rxOctetsOK/txOctetsOK) + # and UBNT-AFLTU-MIB (afLTUethRxBytes/afLTUethTxBytes) for these devices + # because their IF-MIB counters are unreliable. + defp detect_airfiber_counter_overrides(nil, _client_opts), do: nil + + defp detect_airfiber_counter_overrides(snmp_device, client_opts) do + sys_oid = snmp_device.sys_object_id || "" + + cond do + # Ubiquiti AirFiber (enterprise OID 41112) + String.starts_with?(sys_oid, "1.3.6.1.4.1.41112") -> + # Check if it's an LTU device + case Client.get(client_opts, "1.3.6.1.4.1.41112.1.10.1.2.2.0") do + {:ok, _} -> :airfiber_ltu + _ -> + # Check for regular AirFiber stats table + case Client.get(client_opts, "1.3.6.1.4.1.41112.1.3.3.1.1.1") do + {:ok, _} -> :airfiber + _ -> nil + end + end + + true -> nil + end + end + + # Fetch traffic stats using AirFiber proprietary MIBs (matching LibreNMS behavior) + defp get_airfiber_stats(nil, _interface, _client_opts), do: :not_airfiber + defp get_airfiber_stats(_af_type, %{if_descr: descr}, _client_opts) + when descr != "eth0" and descr != nil, do: :not_eth0 + + defp get_airfiber_stats(:airfiber, _interface, client_opts) do + # UBNT-AirFIBER-MIB: airFiberStatistics table (1.3.6.1.4.1.41112.1.3.3.1) + oids = [ + {"1.3.6.1.4.1.41112.1.3.3.1.1.1", :if_in_octets}, # rxOctetsOK + {"1.3.6.1.4.1.41112.1.3.3.1.2.1", :if_out_octets}, # txOctetsOK + {"1.3.6.1.4.1.41112.1.3.3.1.5.1", :if_in_errors}, # rxErroredFrames + {"1.3.6.1.4.1.41112.1.3.3.1.10.1", :if_out_errors} # txErroredFrames + ] + fetch_proprietary_stats(client_opts, oids, true) + end + + defp get_airfiber_stats(:airfiber_ltu, _interface, client_opts) do + # UBNT-AFLTU-MIB: afLTUeth table (1.3.6.1.4.1.41112.1.10.1.6) + oids = [ + {"1.3.6.1.4.1.41112.1.10.1.6.1.6.0", :if_in_octets}, # afLTUethRxBytes + {"1.3.6.1.4.1.41112.1.10.1.6.1.4.0", :if_out_octets} # afLTUethTxBytes + ] + fetch_proprietary_stats(client_opts, oids, true) + end + + defp fetch_proprietary_stats(client_opts, oids, is_hc) do + results = + Enum.map(oids, fn {oid, key} -> + case Client.get(client_opts, oid) do + {:ok, value} -> {key, decode_snmp_value(value)} + _ -> {key, nil} + end + end) + + {:ok, Map.new(results ++ [is_hc: is_hc])} + end + defp check_interface_changes(interfaces, device, client_opts, timestamp) do interfaces |> Task.async_stream( @@ -1312,7 +1385,7 @@ defmodule Towerops.Workers.DevicePollerWorker do if_out_discards: "1.3.6.1.2.1.2.2.1.19.#{if_index}" ] - octet_results = fetch_octet_counters(client_opts, hc_oids, std_oids) + {octet_results, is_hc} = fetch_octet_counters(client_opts, hc_oids, std_oids) error_results = Enum.map(error_oids, fn {key, oid} -> @@ -1322,7 +1395,7 @@ defmodule Towerops.Workers.DevicePollerWorker do end end) - {:ok, Map.new(octet_results ++ error_results)} + {:ok, Map.new(octet_results ++ error_results ++ [is_hc: is_hc])} end defp fetch_octet_counters(client_opts, hc_oids, std_oids) do @@ -1334,10 +1407,10 @@ defmodule Towerops.Workers.DevicePollerWorker do hc_available = Enum.all?(hc_results, fn {_key, value} -> value != :not_available end) if hc_available do - Enum.map(hc_results, fn {key, value} -> {key, value} end) + {Enum.map(hc_results, fn {key, value} -> {key, value} end), true} else Logger.debug("HC counters not available, using standard 32-bit counters") - Enum.map(std_oids, &fetch_standard_counter(client_opts, &1)) + {Enum.map(std_oids, &fetch_standard_counter(client_opts, &1)), false} end end @@ -1375,6 +1448,10 @@ defmodule Towerops.Workers.DevicePollerWorker do result = case size do + 4 -> + <> = value + counter + 8 -> <> = value counter diff --git a/priv/repo/migrations/20260325200000_add_is_hc_to_snmp_interface_stats.exs b/priv/repo/migrations/20260325200000_add_is_hc_to_snmp_interface_stats.exs new file mode 100644 index 00000000..e799928a --- /dev/null +++ b/priv/repo/migrations/20260325200000_add_is_hc_to_snmp_interface_stats.exs @@ -0,0 +1,9 @@ +defmodule Towerops.Repo.Migrations.AddIsHcToSnmpInterfaceStats do + use Ecto.Migration + + def change do + alter table(:snmp_interface_stats) do + add :is_hc, :boolean, default: true, null: false + end + end +end diff --git a/test/towerops/capacity_test.exs b/test/towerops/capacity_test.exs index c35e8688..f80cdd19 100644 --- a/test/towerops/capacity_test.exs +++ b/test/towerops/capacity_test.exs @@ -83,7 +83,9 @@ defmodule Towerops.CapacityTest do assert_in_delta result.max_bps, 133_333.33, 1.0 end - test "applies 32-bit wrap correction for small counter decreases", %{interface: interface} do + test "clamps to zero when counter decreases (negative delta discarded)", %{ + interface: interface + } do now = DateTime.utc_now() t1 = DateTime.add(now, -120, :second) t2 = DateTime.add(now, -60, :second) @@ -94,9 +96,7 @@ defmodule Towerops.CapacityTest do checked_at: t1 }) - # Counter decreased — could be 32-bit wrap on a fast link. - # We always apply 2^32 correction for negative deltas within 32-bit range. - # Corrected delta: (4,294,967,296 - 5,000,000) + 1,000,000 = 4,290,967,296 + # Counter decreased — ambiguous between 32-bit wrap and HC reset, so clamped to 0 insert_interface_stat(interface.id, %{ if_in_octets: 1_000_000, if_out_octets: 6_000_000, @@ -105,27 +105,25 @@ defmodule Towerops.CapacityTest do result = Capacity.calculate_throughput(interface.id, DateTime.add(now, -180, :second)) - # 4,290,967,296 * 8 / 60 = 572,128,972.8 bps - assert_in_delta result.in_bps, 572_128_972.8, 1.0 + assert result.in_bps == 0.0 + # out: positive delta = (6M - 5M) * 8 / 60 = 133,333.33 bps assert_in_delta result.out_bps, 133_333.33, 1.0 end - test "recovers correct rate when 32-bit counter wraps with prev in lower half of range", %{ + test "clamps to zero when 32-bit-range counter decreases (indistinguishable from HC reset)", %{ interface: interface } do now = DateTime.utc_now() t1 = DateTime.add(now, -120, :second) t2 = DateTime.add(now, -60, :second) - # Counter at 1.36B (lower half of 32-bit range, < 2^31) - # This is the exact scenario from production data insert_interface_stat(interface.id, %{ if_in_octets: 1_360_303_152, if_out_octets: 2_961_644_622, checked_at: t1 }) - # Counter wrapped — high-bandwidth link pushed past 2^32 + # Negative delta — clamped to 0 regardless of whether it is a wrap or reset insert_interface_stat(interface.id, %{ if_in_octets: 105_884_342, if_out_octets: 3_220_277_392, @@ -134,26 +132,23 @@ defmodule Towerops.CapacityTest do result = Capacity.calculate_throughput(interface.id, DateTime.add(now, -180, :second)) - # in: (4,294,967,296 - 1,360,303,152) + 105,884,342 = 3,040,548,486 bytes - # bps = 3,040,548,486 * 8 / 60 = 405,406,464.8 - assert_in_delta result.in_bps, 405_406_464.8, 1.0 - # out: positive delta, no wrap needed = (3,220,277,392 - 2,961,644,622) * 8 / 60 + assert result.in_bps == 0.0 + # out: positive delta = (3,220,277,392 - 2,961,644,622) * 8 / 60 assert_in_delta result.out_bps, 34_484_369.33, 1.0 end - test "recovers correct rate when 32-bit counter wraps near max", %{interface: interface} do + test "clamps to zero when counter decreases near 32-bit max", %{interface: interface} do now = DateTime.utc_now() t1 = DateTime.add(now, -120, :second) t2 = DateTime.add(now, -60, :second) - # Counter was near the 32-bit maximum (2^32 = 4,294,967,296) insert_interface_stat(interface.id, %{ if_in_octets: 4_290_000_000, if_out_octets: 4_290_000_000, checked_at: t1 }) - # Counter wrapped around — now at a small value + # Counter decreased — clamped to 0 insert_interface_stat(interface.id, %{ if_in_octets: 100_000_000, if_out_octets: 100_000_000, @@ -162,10 +157,30 @@ defmodule Towerops.CapacityTest do result = Capacity.calculate_throughput(interface.id, DateTime.add(now, -180, :second)) - # Correct delta: (4,294,967,296 - 4,290,000,000) + 100,000,000 = 104,967,296 bytes in 60s - # bps = 104,967,296 * 8 / 60 = 13,995,639 bps - assert_in_delta result.in_bps, 13_995_639.0, 1.0 - assert_in_delta result.out_bps, 13_995_639.0, 1.0 + assert result.in_bps == 0.0 + assert result.out_bps == 0.0 + end + + test "discards data points above 400 Gbps as counter anomalies", %{interface: interface} do + now = DateTime.utc_now() + t1 = DateTime.add(now, -120, :second) + t2 = DateTime.add(now, -60, :second) + + insert_interface_stat(interface.id, %{ + if_in_octets: 0, + if_out_octets: 0, + checked_at: t1 + }) + + # delta of 6_000_000_000_000 bytes / 60s * 8 = 800 Gbps — above 400 Gbps ceiling + insert_interface_stat(interface.id, %{ + if_in_octets: 6_000_000_000_000, + if_out_octets: 0, + checked_at: t2 + }) + + result = Capacity.calculate_throughput(interface.id, DateTime.add(now, -180, :second)) + assert result.in_bps == 0.0 end test "clamps to zero when HC 64-bit counter resets after device reboot", %{interface: interface} do