diff --git a/lib/towerops/snmp.ex b/lib/towerops/snmp.ex index 724ab678..c7055412 100644 --- a/lib/towerops/snmp.ex +++ b/lib/towerops/snmp.ex @@ -22,6 +22,7 @@ defmodule Towerops.Snmp do alias Towerops.Snmp.SensorReading alias Towerops.Snmp.StateSensor alias Towerops.Snmp.Storage + alias Towerops.Snmp.StorageReading alias Towerops.Snmp.Vlan @doc """ @@ -662,4 +663,83 @@ defmodule Towerops.Snmp do |> Storage.changeset(attrs) |> Repo.update() end + + # Storage reading queries + + @doc """ + Gets recent storage readings for a storage entry. + + ## Options + + - `:limit` - Maximum number of readings to return (default: 100) + - `:since` - Only return readings after this datetime + """ + def get_storage_readings(storage_id, opts \\ []) do + limit = Keyword.get(opts, :limit, 100) + since = Keyword.get(opts, :since) + + query = + StorageReading + |> where([r], r.storage_id == ^storage_id) + |> order_by([r], desc: r.checked_at) + |> limit(^limit) + + query = + if since do + where(query, [r], r.checked_at >= ^since) + else + query + end + + Repo.all(query) + end + + @doc """ + Gets the latest storage reading for a storage entry. + """ + def get_latest_storage_reading(storage_id) do + StorageReading + |> where([r], r.storage_id == ^storage_id) + |> order_by([r], desc: r.checked_at) + |> limit(1) + |> Repo.one() + end + + @doc """ + Gets the latest storage readings for multiple storage entries in a single query. + + Returns a map of %{storage_id => reading} for efficient batch loading. + Storage entries without readings will not be present in the map. + + ## Examples + + iex> get_latest_storage_readings_batch([storage_id1, storage_id2]) + %{storage_id1 => %StorageReading{}, storage_id2 => %StorageReading{}} + """ + def get_latest_storage_readings_batch(storage_ids) when is_list(storage_ids) do + if Enum.empty?(storage_ids) do + %{} + else + # Use DISTINCT ON to get only the latest reading per storage + query = + from(r in StorageReading, + where: r.storage_id in ^storage_ids, + distinct: r.storage_id, + order_by: [asc: r.storage_id, desc: r.checked_at] + ) + + query + |> Repo.all() + |> Map.new(fn reading -> {reading.storage_id, reading} end) + end + end + + @doc """ + Records a new storage reading. + """ + def create_storage_reading(attrs) do + %StorageReading{} + |> StorageReading.changeset(attrs) + |> Repo.insert() + end end diff --git a/lib/towerops/snmp/poller_worker.ex b/lib/towerops/snmp/poller_worker.ex index c7f9f145..dfb11675 100644 --- a/lib/towerops/snmp/poller_worker.ex +++ b/lib/towerops/snmp/poller_worker.ex @@ -169,6 +169,9 @@ defmodule Towerops.Snmp.PollerWorker do end), Task.Supervisor.async_nolink(PollerTaskSupervisor, fn -> poll_device_processors(device, snmp_device, client_opts, now) + end), + Task.Supervisor.async_nolink(PollerTaskSupervisor, fn -> + poll_device_storage(device, snmp_device, client_opts, now) end) ] @@ -285,6 +288,86 @@ defmodule Towerops.Snmp.PollerWorker do Logger.error("Error polling processors for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}") end + defp poll_device_storage(device, snmp_device, client_opts, now) do + storage_entries = snmp_device.storage || [] + + if storage_entries != [] do + poll_storage(storage_entries, client_opts, now) + Logger.debug("Polled #{length(storage_entries)} storage entries for #{device.name}") + + # Broadcast storage update event to device-specific topic + Phoenix.PubSub.broadcast( + Towerops.PubSub, + "device:#{device.id}", + {:storage_updated, device.id} + ) + end + rescue + error -> + Logger.error("Error polling storage for #{device.name}: #{inspect(error)}\n#{Exception.format_stacktrace()}") + end + + defp poll_storage(storage_entries, client_opts, timestamp) do + # Poll storage entries in parallel with max_concurrency + storage_entries + |> Task.async_stream( + fn storage -> + result = poll_storage_value(storage, client_opts) + handle_storage_poll_result(storage, result, timestamp) + end, + max_concurrency: 2, + timeout: 40_000, + on_timeout: :kill_task + ) + |> Stream.run() + end + + defp poll_storage_value(storage, client_opts) do + # HOST-RESOURCES-MIB storage OIDs + used_oid = "1.3.6.1.2.1.25.2.3.1.6.#{storage.storage_index}" + size_oid = "1.3.6.1.2.1.25.2.3.1.5.#{storage.storage_index}" + alloc_oid = "1.3.6.1.2.1.25.2.3.1.4.#{storage.storage_index}" + + with {:ok, used_units} <- Client.get(client_opts, used_oid), + {:ok, size_units} <- Client.get(client_opts, size_oid), + {:ok, alloc_units} <- Client.get(client_opts, alloc_oid), + true <- is_integer(used_units) and is_integer(size_units) and is_integer(alloc_units), + true <- size_units > 0 do + # Calculate bytes from allocation units + used_bytes = used_units * alloc_units + total_bytes = size_units * alloc_units + usage_percent = used_bytes / total_bytes * 100 + + {:ok, %{used_bytes: used_bytes, total_bytes: total_bytes, usage_percent: usage_percent}} + else + {:error, reason} -> {:error, reason} + false -> {:error, :invalid_values} + _ -> {:error, :unknown} + end + end + + defp handle_storage_poll_result(storage, {:ok, values}, timestamp) do + # Create a reading + Snmp.create_storage_reading(%{ + storage_id: storage.id, + used_bytes: values.used_bytes, + total_bytes: values.total_bytes, + usage_percent: values.usage_percent, + checked_at: timestamp + }) + + # Update storage entry with latest values + Snmp.update_storage(storage, %{ + used_bytes: values.used_bytes, + total_bytes: values.total_bytes, + last_checked_at: timestamp + }) + end + + defp handle_storage_poll_result(storage, {:error, reason}, _timestamp) do + Logger.debug("Failed to poll storage #{storage.description}: #{inspect(reason)}") + end + defp poll_processors(processors, client_opts, timestamp) do # Poll processors in parallel with max_concurrency processors diff --git a/lib/towerops/snmp/profiles/vendors/routeros.ex b/lib/towerops/snmp/profiles/vendors/routeros.ex index fc9c18f3..7d168393 100644 --- a/lib/towerops/snmp/profiles/vendors/routeros.ex +++ b/lib/towerops/snmp/profiles/vendors/routeros.ex @@ -585,7 +585,9 @@ defmodule Towerops.Snmp.Profiles.Vendors.Routeros do unit_type = get_int_value(values, 4) if name && value do - {sensor_type, sensor_unit, divisor} = gauge_unit_to_sensor_type(unit_type) + # Get base type from unit, then check if name suggests different type + base_type = gauge_unit_to_sensor_type(unit_type) + {sensor_type, sensor_unit, divisor} = maybe_override_sensor_type(base_type, name) [ %{ @@ -612,6 +614,42 @@ defmodule Towerops.Snmp.Profiles.Vendors.Routeros do defp gauge_unit_to_sensor_type(@gauge_unit_status), do: {"state", "", 1} defp gauge_unit_to_sensor_type(_), do: {"unknown", "", 1} + # Override sensor type based on name when the device reports incorrect unit type + # This handles cases where MikroTik devices report wrong mtxrGaugeUnit values + defp maybe_override_sensor_type({sensor_type, _unit, _divisor} = type_info, name) when is_binary(name) do + name_lower = String.downcase(name) + inferred_type = infer_type_from_name(name_lower) + + case {sensor_type, inferred_type} do + # Type matches or no override needed + {type, type} -> type_info + {_, nil} -> type_info + # Override with inferred type + {_, inferred} -> sensor_type_config(inferred) + end + end + + defp maybe_override_sensor_type(type_info, _name), do: type_info + + # Infer sensor type from name keywords + defp infer_type_from_name(name_lower) do + cond do + String.contains?(name_lower, "voltage") -> "voltage" + String.contains?(name_lower, ["temperature", "temp"]) -> "temperature" + String.contains?(name_lower, ["current", "ampere"]) -> "current" + String.contains?(name_lower, ["power", "watt"]) -> "power" + String.contains?(name_lower, ["fan", "rpm"]) -> "fanspeed" + true -> nil + end + end + + # Sensor type configurations for overrides + defp sensor_type_config("voltage"), do: {"voltage", "V", 10} + defp sensor_type_config("temperature"), do: {"temperature", "°C", 10} + defp sensor_type_config("current"), do: {"current", "A", 1000} + defp sensor_type_config("power"), do: {"power", "W", 10} + defp sensor_type_config("fanspeed"), do: {"fanspeed", "RPM", 1} + # Discover sensors from mtxrOpticalTable (SFP/optical transceiver sensors) defp discover_optical_sensors(client_opts) do case Client.walk(client_opts, @optical_table) do diff --git a/lib/towerops/snmp/storage_reading.ex b/lib/towerops/snmp/storage_reading.ex new file mode 100644 index 00000000..e63fe8b4 --- /dev/null +++ b/lib/towerops/snmp/storage_reading.ex @@ -0,0 +1,45 @@ +defmodule Towerops.Snmp.StorageReading do + @moduledoc """ + Time-series storage reading schema. + + Stores individual storage readings (used/total bytes, usage percent, timestamp) + collected during SNMP polling cycles. + """ + use Ecto.Schema + + import Ecto.Changeset + + alias Towerops.Snmp.Storage + + @primary_key {:id, :binary_id, autogenerate: true} + @foreign_key_type :binary_id + schema "snmp_storage_readings" do + field :used_bytes, :integer + field :total_bytes, :integer + field :usage_percent, :float + field :checked_at, :utc_datetime + + belongs_to :storage, Storage + + timestamps(type: :utc_datetime, updated_at: false) + end + + @type t :: %__MODULE__{ + id: Ecto.UUID.t(), + used_bytes: integer() | nil, + total_bytes: integer() | nil, + usage_percent: float() | nil, + checked_at: DateTime.t(), + storage_id: Ecto.UUID.t(), + storage: Ecto.Association.NotLoaded.t() | Storage.t(), + inserted_at: DateTime.t() + } + + @doc false + def changeset(reading, attrs) do + reading + |> cast(attrs, [:storage_id, :used_bytes, :total_bytes, :usage_percent, :checked_at]) + |> validate_required([:storage_id, :checked_at]) + |> foreign_key_constraint(:storage_id) + end +end diff --git a/lib/towerops_web/live/device_live/show.ex b/lib/towerops_web/live/device_live/show.ex index ad0e7a64..49c58136 100644 --- a/lib/towerops_web/live/device_live/show.ex +++ b/lib/towerops_web/live/device_live/show.ex @@ -118,6 +118,11 @@ defmodule ToweropsWeb.DeviceLive.Show do {:noreply, load_equipment_data(socket, socket.assigns.device.id)} end + @impl true + def handle_info({:storage_updated, _device_id}, socket) do + {:noreply, load_equipment_data(socket, socket.assigns.device.id)} + end + @impl true def handle_info({:monitoring_check_updated, _device_id}, socket) do {:noreply, load_equipment_data(socket, socket.assigns.device.id)} diff --git a/lib/towerops_web/live/device_live/show.html.heex b/lib/towerops_web/live/device_live/show.html.heex index 192864c0..e18c390f 100644 --- a/lib/towerops_web/live/device_live/show.html.heex +++ b/lib/towerops_web/live/device_live/show.html.heex @@ -479,9 +479,14 @@ else 0.0 end %> -
+ <.link + navigate={ + ~p"/devices/#{@device.id}/graph/storage_volume?storage_id=#{storage.id}" + } + class="block hover:bg-gray-50 dark:hover:bg-gray-700/50 -mx-2 px-2 py-1 rounded transition-colors" + >
- + {storage.description || "Storage #{storage.storage_index}"} @@ -509,7 +514,7 @@ {String.replace(storage.storage_type || "other", "_", " ") |> String.capitalize()}
-
+ <% end %> diff --git a/lib/towerops_web/live/graph_live/show.ex b/lib/towerops_web/live/graph_live/show.ex index 2714e482..440806fb 100644 --- a/lib/towerops_web/live/graph_live/show.ex +++ b/lib/towerops_web/live/graph_live/show.ex @@ -21,6 +21,7 @@ defmodule ToweropsWeb.GraphLive.Show do range = Map.get(params, "range", "24h") interface_id = Map.get(params, "interface_id") sensor_id = Map.get(params, "sensor_id") + storage_id = Map.get(params, "storage_id") socket = socket @@ -29,6 +30,7 @@ defmodule ToweropsWeb.GraphLive.Show do |> assign(:range, range) |> assign(:interface_id, interface_id) |> assign(:sensor_id, sensor_id) + |> assign(:storage_id, storage_id) |> load_graph_data() {:noreply, socket} @@ -54,8 +56,15 @@ defmodule ToweropsWeb.GraphLive.Show do base_params end - if assigns.sensor_id do - Map.put(base_params, "sensor_id", assigns.sensor_id) + base_params = + if assigns.sensor_id do + Map.put(base_params, "sensor_id", assigns.sensor_id) + else + base_params + end + + if assigns.storage_id do + Map.put(base_params, "storage_id", assigns.storage_id) else base_params end @@ -69,6 +78,7 @@ defmodule ToweropsWeb.GraphLive.Show do range = socket.assigns.range interface_id = socket.assigns.interface_id sensor_id = socket.assigns.sensor_id + storage_id = socket.assigns.storage_id device = Devices.get_device!(device_id) @@ -83,6 +93,9 @@ defmodule ToweropsWeb.GraphLive.Show do sensor_type == "traffic" -> {load_traffic_chart_data(device_id, range), nil} + sensor_type == "storage_volume" && storage_id -> + {load_storage_chart_data(storage_id, range), get_storage_name(storage_id)} + sensor_id -> # Load only the specific sensor when sensor_id is provided {load_single_sensor_chart_data(sensor_id, range), get_sensor_name(sensor_id)} @@ -119,6 +132,7 @@ defmodule ToweropsWeb.GraphLive.Show do defp get_chart_config("processors"), do: {"Processor Usage", "%", false} defp get_chart_config("memory"), do: {"Memory Usage", "%", false} defp get_chart_config("storage"), do: {"Storage Usage", "%", false} + defp get_chart_config("storage_volume"), do: {"Storage Usage", "%", false} defp get_chart_config("temperature"), do: {"Temperature", "°C", true} defp get_chart_config("voltage"), do: {"Voltage", "V", true} defp get_chart_config("traffic"), do: {"Overall Traffic", "bps", true} @@ -165,6 +179,45 @@ defmodule ToweropsWeb.GraphLive.Show do end end + defp load_storage_chart_data(storage_id, range) do + case Snmp.get_storage(storage_id) do + nil -> + nil + + storage -> + since = get_datetime_from_range(range) + limit = get_limit_for_range(range) + dataset = storage_to_dataset(storage, since, limit) + Jason.encode!(%{datasets: [dataset]}) + end + end + + defp get_storage_name(storage_id) do + case Snmp.get_storage(storage_id) do + nil -> nil + storage -> storage.description || "Storage #{storage.storage_index}" + end + end + + defp storage_to_dataset(storage, since, limit) do + readings = + storage.id + |> Snmp.get_storage_readings(since: since, limit: limit) + |> Enum.reverse() + + %{ + label: storage.description || "Storage #{storage.storage_index}", + data: Enum.map(readings, &storage_reading_to_chart_point/1) + } + end + + defp storage_reading_to_chart_point(reading) do + %{ + x: DateTime.to_unix(reading.checked_at, :millisecond), + y: if(reading.usage_percent, do: Float.round(reading.usage_percent, 1)) + } + end + defp build_sensor_chart_json(device, sensor_types, range) do sensors = device.sensors diff --git a/priv/repo/migrations/20260122234739_add_snmp_storage_readings.exs b/priv/repo/migrations/20260122234739_add_snmp_storage_readings.exs new file mode 100644 index 00000000..5c1f8de6 --- /dev/null +++ b/priv/repo/migrations/20260122234739_add_snmp_storage_readings.exs @@ -0,0 +1,28 @@ +defmodule Towerops.Repo.Migrations.AddSnmpStorageReadings do + use Ecto.Migration + + def up do + create table(:snmp_storage_readings, primary_key: false) do + add :id, :binary_id, primary_key: true + + add :storage_id, references(:snmp_storage, type: :binary_id, on_delete: :delete_all), + null: false + + add :used_bytes, :bigint + add :total_bytes, :bigint + add :usage_percent, :float + add :checked_at, :utc_datetime, null: false + + timestamps(type: :utc_datetime, updated_at: false) + end + + # Create indexes for query performance + create index(:snmp_storage_readings, [:storage_id]) + create index(:snmp_storage_readings, [:storage_id, :checked_at]) + create index(:snmp_storage_readings, [:checked_at]) + end + + def down do + drop table(:snmp_storage_readings) + end +end diff --git a/test/towerops/snmp_test.exs b/test/towerops/snmp_test.exs index 1499cd67..ea4f2b11 100644 --- a/test/towerops/snmp_test.exs +++ b/test/towerops/snmp_test.exs @@ -11,6 +11,8 @@ defmodule Towerops.SnmpTest do alias Towerops.Snmp.Sensor alias Towerops.Snmp.SensorReading alias Towerops.Snmp.SnmpMock + alias Towerops.Snmp.Storage + alias Towerops.Snmp.StorageReading setup do user = user_fixture() @@ -1309,4 +1311,332 @@ defmodule Towerops.SnmpTest do assert Snmp.get_neighbor(neighbor2.id) end end + + describe "get_storage_readings/2" do + test "returns recent readings for a storage entry", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + for i <- 1..5 do + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 1000 * i, + total_bytes: 10_000, + usage_percent: 10.0 * i, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + end + + readings = Snmp.get_storage_readings(storage.id) + assert length(readings) == 5 + end + + test "respects limit option", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + for i <- 1..10 do + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 1000 * i, + total_bytes: 10_000, + usage_percent: 10.0 * i, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + end + + readings = Snmp.get_storage_readings(storage.id, limit: 3) + assert length(readings) == 3 + end + + test "respects since option", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + past = DateTime.add(DateTime.utc_now(), -3600, :second) + recent = DateTime.add(DateTime.utc_now(), -60, :second) + + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 1000, + total_bytes: 10_000, + usage_percent: 10.0, + checked_at: past + }) + |> Repo.insert!() + + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 5000, + total_bytes: 10_000, + usage_percent: 50.0, + checked_at: recent + }) + |> Repo.insert!() + + since_time = DateTime.add(DateTime.utc_now(), -120, :second) + readings = Snmp.get_storage_readings(storage.id, since: since_time) + assert length(readings) == 1 + assert hd(readings).usage_percent == 50.0 + end + + test "returns empty list for storage with no readings", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + assert Snmp.get_storage_readings(storage.id) == [] + end + end + + describe "get_latest_storage_reading/1" do + test "returns most recent reading", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + old_reading = + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 1000, + total_bytes: 10_000, + usage_percent: 10.0, + checked_at: DateTime.add(DateTime.utc_now(), -3600, :second) + }) + |> Repo.insert!() + + new_reading = + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 5000, + total_bytes: 10_000, + usage_percent: 50.0, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + + latest = Snmp.get_latest_storage_reading(storage.id) + assert latest.id == new_reading.id + refute latest.id == old_reading.id + end + + test "returns nil for storage with no readings", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + assert Snmp.get_latest_storage_reading(storage.id) == nil + end + end + + describe "get_latest_storage_readings_batch/1" do + test "returns empty map for empty storage list" do + assert Snmp.get_latest_storage_readings_batch([]) == %{} + end + + test "returns readings for multiple storage entries", %{snmp_device: snmp_device} do + storage1 = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + storage2 = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 2, + storage_type: "ram", + description: "Physical Memory" + }) + |> Repo.insert!() + + reading1 = + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage1.id, + used_bytes: 5000, + total_bytes: 10_000, + usage_percent: 50.0, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + + reading2 = + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage2.id, + used_bytes: 2000, + total_bytes: 8000, + usage_percent: 25.0, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + + batch = Snmp.get_latest_storage_readings_batch([storage1.id, storage2.id]) + assert map_size(batch) == 2 + assert batch[storage1.id].id == reading1.id + assert batch[storage2.id].id == reading2.id + end + + test "returns only latest reading per storage", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + _old_reading = + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 1000, + total_bytes: 10_000, + usage_percent: 10.0, + checked_at: DateTime.add(DateTime.utc_now(), -3600, :second) + }) + |> Repo.insert!() + + new_reading = + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage.id, + used_bytes: 5000, + total_bytes: 10_000, + usage_percent: 50.0, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + + batch = Snmp.get_latest_storage_readings_batch([storage.id]) + assert map_size(batch) == 1 + assert batch[storage.id].id == new_reading.id + assert batch[storage.id].usage_percent == 50.0 + end + + test "excludes storage without readings", %{snmp_device: snmp_device} do + storage_with_reading = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + storage_without_reading = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 2, + storage_type: "ram", + description: "Physical Memory" + }) + |> Repo.insert!() + + %StorageReading{} + |> StorageReading.changeset(%{ + storage_id: storage_with_reading.id, + used_bytes: 5000, + total_bytes: 10_000, + usage_percent: 50.0, + checked_at: DateTime.utc_now() + }) + |> Repo.insert!() + + batch = Snmp.get_latest_storage_readings_batch([storage_with_reading.id, storage_without_reading.id]) + assert map_size(batch) == 1 + assert Map.has_key?(batch, storage_with_reading.id) + refute Map.has_key?(batch, storage_without_reading.id) + end + end + + describe "create_storage_reading/1" do + test "creates a storage reading with valid attributes", %{snmp_device: snmp_device} do + storage = + %Storage{} + |> Storage.changeset(%{ + snmp_device_id: snmp_device.id, + storage_index: 1, + storage_type: "fixed_disk", + description: "/dev/sda1" + }) + |> Repo.insert!() + + {:ok, reading} = + Snmp.create_storage_reading(%{ + storage_id: storage.id, + used_bytes: 5000, + total_bytes: 10_000, + usage_percent: 50.0, + checked_at: DateTime.utc_now() + }) + + assert reading.storage_id == storage.id + assert reading.used_bytes == 5000 + assert reading.total_bytes == 10_000 + assert reading.usage_percent == 50.0 + end + + test "returns error with invalid attributes" do + {:error, changeset} = Snmp.create_storage_reading(%{used_bytes: 1000}) + assert changeset.errors[:storage_id] + end + end end