From ecec8ba845bafb04a503bc5fe7f088a37c999b41 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Mon, 11 May 2026 18:02:21 -0500 Subject: [PATCH] feat(mikrotik): webhook receiver for MikroTik RouterOS events Receives webhooks from MikroTik RouterOS when devices come online, extracts MAC/IP data, and reconciles with Gaiia inventory. --- lib/towerops/integrations/integration.ex | 2 +- lib/towerops/mikrotik.ex | 147 ++++++++++ lib/towerops/mikrotik/event.ex | 63 ++++ lib/towerops/mikrotik/gaiia_resolver.ex | 268 ++++++++++++++++++ lib/towerops/mikrotik/ip_assignment.ex | 63 ++++ .../workers/mikrotik_webhook_worker.ex | 143 ++++++++++ .../api/v1/mikrotik_webhook_controller.ex | 113 ++++++++ lib/towerops_web/router.ex | 2 + ...1175355_create_mikrotik_webhook_tables.exs | 63 ++++ test/towerops/mikrotik_test.exs | 196 +++++++++++++ .../workers/mikrotik_webhook_worker_test.exs | 120 ++++++++ .../v1/mikrotik_webhook_controller_test.exs | 159 +++++++++++ 12 files changed, 1338 insertions(+), 1 deletion(-) create mode 100644 lib/towerops/mikrotik.ex create mode 100644 lib/towerops/mikrotik/event.ex create mode 100644 lib/towerops/mikrotik/gaiia_resolver.ex create mode 100644 lib/towerops/mikrotik/ip_assignment.ex create mode 100644 lib/towerops/workers/mikrotik_webhook_worker.ex create mode 100644 lib/towerops_web/controllers/api/v1/mikrotik_webhook_controller.ex create mode 100644 priv/repo/migrations/20260511175355_create_mikrotik_webhook_tables.exs create mode 100644 test/towerops/mikrotik_test.exs create mode 100644 test/towerops/workers/mikrotik_webhook_worker_test.exs create mode 100644 test/towerops_web/controllers/api/v1/mikrotik_webhook_controller_test.exs diff --git a/lib/towerops/integrations/integration.ex b/lib/towerops/integrations/integration.ex index fedbd950..90109ae8 100644 --- a/lib/towerops/integrations/integration.ex +++ b/lib/towerops/integrations/integration.ex @@ -14,7 +14,7 @@ defmodule Towerops.Integrations.Integration do @primary_key {:id, :binary_id, autogenerate: true} @foreign_key_type :binary_id - @valid_providers ~w(preseem gaiia pagerduty netbox sonar splynx visp uisp cn_maestro) + @valid_providers ~w(preseem gaiia pagerduty netbox sonar splynx visp uisp cn_maestro mikrotik) @valid_sync_statuses ~w(never success partial failed) schema "integrations" do diff --git a/lib/towerops/mikrotik.ex b/lib/towerops/mikrotik.ex new file mode 100644 index 00000000..0f494ce9 --- /dev/null +++ b/lib/towerops/mikrotik.ex @@ -0,0 +1,147 @@ +defmodule Towerops.Mikrotik do + @moduledoc """ + Context for MikroTik router-sourced DHCP / PPPoE webhook events. + + Routers POST lease-script and PPP on-up/on-down events to + `POST /api/v1/webhooks/mikrotik/:organization_id/{dhcp,pppoe}`. + + The controller authenticates with a shared secret, enqueues an Oban job, + and ACKs immediately. The worker (`Towerops.Workers.MikrotikWebhookWorker`) + resolves the MAC / username to a Gaiia account, writes one row to + `mikrotik_events`, and upserts the current state in `mikrotik_ip_assignments`. + + Authoritative resolution path: + 1. caller-MAC (PPPoE) / lease-MAC (DHCP) -> Gaiia InventoryItem.mac + custom field -> Account assignment. + 2. Fallback: PPPoE username / DHCP host-name fuzzy match against + Account.displayName. + 3. If neither resolves, the event is still recorded with + `resolved_account_id=NULL` so later backfill can pick it up. + """ + + import Ecto.Query + + alias Towerops.Mikrotik.Event + alias Towerops.Mikrotik.IpAssignment + alias Towerops.Repo + + @doc """ + Inserts a single MikroTik event row. Returns `{:ok, event}` or `{:error, changeset}`. + """ + def record_event(attrs) do + %Event{} + |> Event.changeset(attrs) + |> Repo.insert() + end + + @doc """ + Upserts the current-state row for a given (organization_id, ip). + + - On `assign` events, sets/refreshes mac, account, last_seen_at, is_active=true. + - On `release` events, flips is_active=false but preserves the most-recent + resolution so we keep the audit trail. + + `attrs` must include at minimum: organization_id, ip, source, event_type. + Optional: mac, router, username, host, gaiia_account_id, + gaiia_inventory_item_id, account_name, resolved_via. + """ + def upsert_assignment(attrs) do + now = DateTime.utc_now() + + event_type = Map.fetch!(attrs, :event_type) + org_id = Map.fetch!(attrs, :organization_id) + ip = Map.fetch!(attrs, :ip) + + base_attrs = + attrs + |> Map.delete(:event_type) + |> Map.put(:is_active, event_type == "assign") + |> Map.put(:last_seen_at, now) + |> Map.put(:first_seen_at, now) + + case Repo.get_by(IpAssignment, organization_id: org_id, ip: ip) do + nil -> + %IpAssignment{} + |> IpAssignment.changeset(base_attrs) + |> Repo.insert() + + existing -> + update_attrs = + base_attrs + |> Map.put(:first_seen_at, existing.first_seen_at || now) + |> maybe_keep_resolution(existing, event_type) + + existing + |> IpAssignment.changeset(update_attrs) + |> Repo.update() + end + end + + # On release events, only flip is_active; keep prior resolution fields so the + # audit trail isn't lost when a customer disconnects. + defp maybe_keep_resolution(attrs, existing, "release") do + Map.merge( + %{ + mac: existing.mac, + username: existing.username, + host: existing.host, + router: existing.router, + source: existing.source, + gaiia_account_id: existing.gaiia_account_id, + gaiia_inventory_item_id: existing.gaiia_inventory_item_id, + account_name: existing.account_name, + resolved_via: existing.resolved_via + }, + attrs + ) + end + + defp maybe_keep_resolution(attrs, _existing, _event_type), do: attrs + + @doc "List recent events for an organization (most recent first)." + def list_recent_events(organization_id, limit \\ 100) do + Event + |> where(organization_id: ^organization_id) + |> order_by([e], desc: e.received_at) + |> limit(^limit) + |> Repo.all() + end + + @doc "List currently-active IP assignments for an organization." + def list_active_assignments(organization_id) do + IpAssignment + |> where(organization_id: ^organization_id, is_active: true) + |> order_by(:ip) + |> Repo.all() + end + + @doc "Get the current assignment for a specific IP (if any)." + def get_assignment(organization_id, ip) do + Repo.get_by(IpAssignment, organization_id: organization_id, ip: ip) + end + + @doc "Normalize a MAC string to lowercase 12-char hex (no separators)." + def normalize_mac(nil), do: nil + def normalize_mac(""), do: nil + + def normalize_mac(mac) when is_binary(mac) do + cleaned = + mac + |> String.downcase() + |> String.replace(~r/[^0-9a-f]/, "") + + case cleaned do + <<_::binary-size(12)>> -> cleaned + _ -> nil + end + end + + @doc "Format a normalized MAC back to `XX:XX:XX:XX:XX:XX` (uppercase)." + def format_mac(nil), do: nil + + def format_mac(<>) do + [a, b, c, d, e, f] |> Enum.join(":") |> String.upcase() + end + + def format_mac(other), do: other +end diff --git a/lib/towerops/mikrotik/event.ex b/lib/towerops/mikrotik/event.ex new file mode 100644 index 00000000..0070ea26 --- /dev/null +++ b/lib/towerops/mikrotik/event.ex @@ -0,0 +1,63 @@ +defmodule Towerops.Mikrotik.Event do + @moduledoc """ + Append-only log of MikroTik DHCP / PPPoE webhook events. + + Stored as-received from the router lease-script / PPP profile on-up/on-down + hooks. The Oban worker writes one row per inbound webhook and never updates + these rows after the fact — historical timeline of every IP assignment event. + """ + + use Ecto.Schema + + import Ecto.Changeset + + @primary_key {:id, :binary_id, autogenerate: true} + @foreign_key_type :binary_id + + @sources ~w(dhcp pppoe) + @event_types ~w(assign release) + + schema "mikrotik_events" do + field :source, :string + field :event_type, :string + field :router, :string + field :mac, :string + field :ip, :string + field :username, :string + field :host, :string + field :server, :string + field :remote_id, :string + field :raw_payload, :map + field :resolved_account_id, :string + field :resolved_via, :string + field :received_at, :utc_datetime_usec + + belongs_to :organization, Towerops.Organizations.Organization + + timestamps(type: :utc_datetime, updated_at: false) + end + + def changeset(event, attrs) do + event + |> cast(attrs, [ + :organization_id, + :source, + :event_type, + :router, + :mac, + :ip, + :username, + :host, + :server, + :remote_id, + :raw_payload, + :resolved_account_id, + :resolved_via, + :received_at + ]) + |> validate_required([:organization_id, :source, :event_type, :raw_payload, :received_at]) + |> validate_inclusion(:source, @sources) + |> validate_inclusion(:event_type, @event_types) + |> foreign_key_constraint(:organization_id) + end +end diff --git a/lib/towerops/mikrotik/gaiia_resolver.ex b/lib/towerops/mikrotik/gaiia_resolver.ex new file mode 100644 index 00000000..a90eb1fd --- /dev/null +++ b/lib/towerops/mikrotik/gaiia_resolver.ex @@ -0,0 +1,268 @@ +defmodule Towerops.Mikrotik.GaiiaResolver do + @moduledoc """ + Resolves a MikroTik webhook payload to a Gaiia Account by walking the + inventory item assignment chain. + + Two strategies, in order: + + 1. **MAC match** — find the InventoryItem whose `mac` custom-field equals + the caller-id (PPPoE) or lease MAC (DHCP), then read its Account + assignment. Authoritative when the modem has been entered in Gaiia + inventory and assigned to the right account. + + 2. **Fuzzy name match** — for PPPoE, smush the username (`judydevine`) + and compare against the smushed `Account.displayName` (`Judy Devine`). + For DHCP, smush the lease host-name. Brittle but covered most of our + fleet in the Python port of this resolver. + + Gaiia's public schema doesn't expose PPPoE secrets / RADIUS, so until that + changes we depend on these two paths. See `scripts/pppoe_freeloaders.py` + in the network repo for the reference implementation we ported. + """ + + alias Towerops.Gaiia.Client + alias Towerops.Integrations + alias Towerops.Mikrotik + + require Logger + + @typedoc """ + The resolved identity for a MikroTik event, or `nil` if no match was found. + + - `:account_id` — Gaiia GlobalID string for the Account + - `:inventory_item_id` — Gaiia GlobalID string for the InventoryItem (may be nil if only name matched) + - `:account_name` — Account.displayName at time of resolution + - `:via` — `"mac"`, `"username"`, `"host"`, or `"username-ambiguous"` etc. + """ + @type resolution :: %{ + required(:account_id) => String.t(), + required(:via) => String.t(), + optional(:inventory_item_id) => String.t() | nil, + optional(:account_name) => String.t() | nil + } + + @doc """ + Resolve a webhook event for the given organization. + + Returns `{:ok, resolution}` on a hit, `{:ok, nil}` when neither strategy + matched, or `{:error, reason}` when Gaiia is unreachable / mis-configured. + + `event_attrs` must include `:source` (`"dhcp"` | `"pppoe"`), plus the + relevant fields: `:mac`, `:username`, `:host`. + """ + def resolve(organization_id, event_attrs) do + with {:ok, integration} <- get_gaiia_integration(organization_id), + api_key when is_binary(api_key) <- get_in(integration.credentials, ["api_key"]) do + mac = Mikrotik.normalize_mac(event_attrs[:mac]) + lookup_name = event_attrs[:username] || event_attrs[:host] + + case try_mac(api_key, mac) do + {:ok, %{} = res} -> + {:ok, res} + + {:ok, nil} -> + try_name(api_key, lookup_name) + + {:error, _} = err -> + err + end + else + {:error, :not_found} -> + {:error, :gaiia_integration_not_configured} + + nil -> + {:error, :gaiia_api_key_missing} + end + end + + # --- MAC match --- + + defp try_mac(_api_key, nil), do: {:ok, nil} + + defp try_mac(api_key, normalized_mac) do + formatted = Mikrotik.format_mac(normalized_mac) + + query = """ + query InvByMac($mac: String!) { + inventoryItems(first: 5, filter: { fieldValue: $mac }) { + edges { + node { + id + assignation { + assigneeType + assignee { ... on Account { id displayName } } + } + fields(first: 10) { + edges { + node { data modelField { name } } + } + } + } + } + } + } + """ + + case Client.query(api_key, query, %{"mac" => formatted}) do + {:ok, %{"inventoryItems" => %{"edges" => edges}}} -> + {:ok, pick_mac_match(edges, normalized_mac)} + + {:error, {:graphql_errors, _}} -> + # Server-side filter not supported; degrade to no-MAC-match so the + # caller can try the name-based fallback. The dedicated freeloaders + # syncer covers the bulk-index path; per-event we accept this miss. + Logger.debug("Gaiia inventoryItems mac filter unsupported; skipping MAC strategy") + {:ok, nil} + + {:error, reason} -> + {:error, reason} + end + end + + defp pick_mac_match([], _), do: nil + + defp pick_mac_match(edges, normalized_mac) do + edges + |> Enum.map(& &1["node"]) + |> Enum.find(fn node -> node_has_mac?(node, normalized_mac) end) + |> case do + nil -> + nil + + node -> + case get_in(node, ["assignation", "assignee"]) do + %{"id" => acct_id, "displayName" => name} -> + %{ + account_id: acct_id, + inventory_item_id: node["id"], + account_name: name, + via: "mac" + } + + _ -> + nil + end + end + end + + defp node_has_mac?(node, normalized_mac) do + node + |> get_in(["fields", "edges"]) + |> List.wrap() + |> Enum.any?(fn edge -> + data = get_in(edge, ["node", "data"]) + field_name = get_in(edge, ["node", "modelField", "name"]) + + is_binary(field_name) and field_name |> String.downcase() |> String.starts_with?("mac") and + Mikrotik.normalize_mac(data) == normalized_mac + end) + end + + # --- Name match --- + + defp try_name(_api_key, nil), do: {:ok, nil} + defp try_name(_api_key, ""), do: {:ok, nil} + + defp try_name(api_key, raw_name) do + smushed = smush(raw_name) + + if smushed == "" do + {:ok, nil} + else + query = """ + query AcctsByName($search: String!) { + accounts(first: 25, filter: { displayName: $search }) { + edges { + node { + id + displayName + status { name } + billingSubscriptions(first: 10) { edges { node { status } } } + } + } + } + } + """ + + case Client.query(api_key, query, %{"search" => raw_name}) do + {:ok, %{"accounts" => %{"edges" => edges}}} -> + {:ok, pick_name_match(edges, smushed)} + + {:error, {:graphql_errors, _}} -> + # If the displayName filter shape isn't accepted, we don't have + # an efficient fallback per-event. Bulk reconciliation handles + # the rest. + Logger.debug("Gaiia accounts displayName filter unsupported; skipping name strategy") + {:ok, nil} + + {:error, reason} -> + {:error, reason} + end + end + end + + defp pick_name_match(edges, smushed) do + candidates = + edges + |> Enum.map(& &1["node"]) + |> Enum.filter(fn n -> smush(n["displayName"] || "") == smushed end) + + case candidates do + [] -> + nil + + [only] -> + %{ + account_id: only["id"], + inventory_item_id: nil, + account_name: only["displayName"], + via: "username" + } + + many -> + paying = Enum.filter(many, &paying?/1) + + case paying do + [one] -> + %{ + account_id: one["id"], + inventory_item_id: nil, + account_name: one["displayName"], + via: "username-paying" + } + + _ -> + first = hd(many) + + %{ + account_id: first["id"], + inventory_item_id: nil, + account_name: first["displayName"], + via: "username-ambiguous" + } + end + end + end + + defp paying?(account) do + account + |> get_in(["billingSubscriptions", "edges"]) + |> List.wrap() + |> Enum.any?(fn edge -> + status = get_in(edge, ["node", "status"]) + status in ["ACTIVE", "TRIAL"] + end) + end + + defp smush(nil), do: "" + + defp smush(s) when is_binary(s) do + s + |> String.downcase() + |> String.replace(~r/[^a-z0-9]/, "") + end + + defp get_gaiia_integration(organization_id) do + Integrations.get_integration(organization_id, "gaiia") + end +end diff --git a/lib/towerops/mikrotik/ip_assignment.ex b/lib/towerops/mikrotik/ip_assignment.ex new file mode 100644 index 00000000..a2e5601a --- /dev/null +++ b/lib/towerops/mikrotik/ip_assignment.ex @@ -0,0 +1,63 @@ +defmodule Towerops.Mikrotik.IpAssignment do + @moduledoc """ + Current-state IP <-> customer mapping derived from MikroTik events. + + One row per (organization_id, ip). Upserted on every `assign` event; + `is_active` is flipped to false on a matching `release` event. Joining + this against `Towerops.Gaiia.Account` data is how downstream code answers + "which customer is using this IP right now?". + """ + + use Ecto.Schema + + import Ecto.Changeset + + @primary_key {:id, :binary_id, autogenerate: true} + @foreign_key_type :binary_id + + @sources ~w(dhcp pppoe) + + schema "mikrotik_ip_assignments" do + field :ip, :string + field :mac, :string + field :source, :string + field :router, :string + field :username, :string + field :host, :string + field :gaiia_account_id, :string + field :gaiia_inventory_item_id, :string + field :account_name, :string + field :resolved_via, :string + field :first_seen_at, :utc_datetime_usec + field :last_seen_at, :utc_datetime_usec + field :is_active, :boolean, default: true + + belongs_to :organization, Towerops.Organizations.Organization + + timestamps(type: :utc_datetime) + end + + def changeset(assignment, attrs) do + assignment + |> cast(attrs, [ + :organization_id, + :ip, + :mac, + :source, + :router, + :username, + :host, + :gaiia_account_id, + :gaiia_inventory_item_id, + :account_name, + :resolved_via, + :first_seen_at, + :last_seen_at, + :is_active + ]) + |> validate_required([:organization_id, :ip, :source, :first_seen_at, :last_seen_at]) + |> validate_inclusion(:source, @sources) + |> unique_constraint([:organization_id, :ip]) + |> foreign_key_constraint(:organization_id) + end +end diff --git a/lib/towerops/workers/mikrotik_webhook_worker.ex b/lib/towerops/workers/mikrotik_webhook_worker.ex new file mode 100644 index 00000000..0e2b391f --- /dev/null +++ b/lib/towerops/workers/mikrotik_webhook_worker.ex @@ -0,0 +1,143 @@ +defmodule Towerops.Workers.MikrotikWebhookWorker do + @moduledoc """ + Processes a MikroTik DHCP / PPPoE webhook event asynchronously. + + Pipeline: + 1. Resolve the modem MAC (or fallback username/host) to a Gaiia Account + via `Towerops.Mikrotik.GaiiaResolver`. + 2. Append a row to `mikrotik_events` capturing the raw payload + the + resolved account id (or nil). + 3. Upsert the current-state row in `mikrotik_ip_assignments`. + + The controller has already ACK'd the router, so transient Gaiia errors + are retried by Oban (max_attempts 5). Permanent resolution misses (no + matching MAC/name) are NOT failures: the event is still recorded with + `resolved_account_id=nil`, and bulk reconciliation jobs can backfill + later. + """ + + use Oban.Worker, queue: :default, max_attempts: 5 + + alias Towerops.Mikrotik + alias Towerops.Mikrotik.GaiiaResolver + + require Logger + + @impl Oban.Worker + def perform(%Oban.Job{args: args}) do + organization_id = Map.fetch!(args, "organization_id") + source = Map.fetch!(args, "source") + event_type = Map.fetch!(args, "event_type") + payload = Map.fetch!(args, "payload") + received_at = parse_received_at(args["received_at"]) + + extracted = extract_fields(source, payload) + + resolution_attrs = Map.put(extracted, :source, source) + + {resolved, resolve_error} = + case GaiiaResolver.resolve(organization_id, resolution_attrs) do + {:ok, %{} = r} -> + {r, nil} + + {:ok, nil} -> + {nil, nil} + + {:error, reason} -> + Logger.warning("Mikrotik webhook: Gaiia resolver error: #{inspect(reason)}") + {nil, reason} + end + + event_attrs = build_event_attrs(organization_id, source, event_type, extracted, payload, resolved, received_at) + assignment_attrs = build_assignment_attrs(organization_id, source, event_type, extracted, resolved, received_at) + + with {:ok, _event} <- Mikrotik.record_event(event_attrs), + {:ok, _assn} <- Mikrotik.upsert_assignment(assignment_attrs) do + case resolve_error do + nil -> :ok + # Returning an error tuple makes Oban retry the job; transient + # Gaiia outages will resolve on retry and we'll re-record the + # event (events table dedupes nothing today, that's intentional — + # the audit log is allowed to have duplicates from retries). + :rate_limited -> {:error, :rate_limited} + {:rate_limited, _} -> {:error, :rate_limited} + _ -> :ok + end + end + end + + defp parse_received_at(nil), do: DateTime.utc_now() + + defp parse_received_at(iso) when is_binary(iso) do + case DateTime.from_iso8601(iso) do + {:ok, dt, _} -> dt + _ -> DateTime.utc_now() + end + end + + # MikroTik DHCP lease-script payload: + # { bound: "1"|"0", mac, ip, server, remoteId } + defp extract_fields("dhcp", payload) do + %{ + mac: payload["mac"], + ip: payload["ip"], + server: payload["server"], + remote_id: payload["remoteId"], + username: nil, + host: nil, + router: payload["router"] + } + end + + # MikroTik PPPoE on-up/on-down payload: + # { bound: "1"|"0", username, callingstationid, framedip } + defp extract_fields("pppoe", payload) do + %{ + mac: payload["callingstationid"], + ip: payload["framedip"], + server: nil, + remote_id: nil, + username: payload["username"], + host: nil, + router: payload["router"] + } + end + + defp build_event_attrs(org_id, source, event_type, extracted, payload, resolved, received_at) do + %{ + organization_id: org_id, + source: source, + event_type: event_type, + router: extracted.router, + mac: extracted.mac, + ip: extracted.ip, + username: extracted.username, + host: extracted.host, + server: extracted.server, + remote_id: extracted.remote_id, + raw_payload: payload, + resolved_account_id: resolved && resolved.account_id, + resolved_via: resolved && resolved.via, + received_at: received_at + } + end + + defp build_assignment_attrs(org_id, source, event_type, extracted, resolved, received_at) do + %{ + organization_id: org_id, + ip: extracted.ip, + mac: extracted.mac, + source: source, + event_type: event_type, + router: extracted.router, + username: extracted.username, + host: extracted.host, + gaiia_account_id: resolved && resolved.account_id, + gaiia_inventory_item_id: resolved && resolved[:inventory_item_id], + account_name: resolved && resolved[:account_name], + resolved_via: resolved && resolved.via, + first_seen_at: received_at, + last_seen_at: received_at + } + end +end diff --git a/lib/towerops_web/controllers/api/v1/mikrotik_webhook_controller.ex b/lib/towerops_web/controllers/api/v1/mikrotik_webhook_controller.ex new file mode 100644 index 00000000..f7652d88 --- /dev/null +++ b/lib/towerops_web/controllers/api/v1/mikrotik_webhook_controller.ex @@ -0,0 +1,113 @@ +defmodule ToweropsWeb.Api.V1.MikrotikWebhookController do + @moduledoc """ + Webhook endpoints for MikroTik routers. + + Routers POST DHCP lease-script and PPPoE on-up/on-down events here. Auth is + a per-org shared secret stored on a `mikrotik` integration (header + `X-Towerops-Token`). The handler enqueues an Oban job and ACKs immediately + so the router's `/tool fetch` doesn't block the lease/PPP state machine. + + Routes: + POST /api/v1/webhooks/mikrotik/:organization_id/dhcp + POST /api/v1/webhooks/mikrotik/:organization_id/pppoe + """ + + use ToweropsWeb, :controller + + alias Towerops.Integrations + alias Towerops.Workers.MikrotikWebhookWorker + + require Logger + + @token_header "x-towerops-token" + + def dhcp(conn, params), do: handle(conn, params, "dhcp") + + def pppoe(conn, params), do: handle(conn, params, "pppoe") + + defp handle(conn, %{"organization_id" => organization_id} = params, source) do + payload = Map.delete(params, "organization_id") + + with {:ok, integration} <- get_integration(organization_id), + :ok <- verify_token(conn, integration), + {:ok, event_type} <- parse_event_type(payload) do + args = %{ + "organization_id" => organization_id, + "source" => source, + "event_type" => event_type, + "payload" => payload, + "received_at" => DateTime.to_iso8601(DateTime.utc_now()) + } + + case Oban.insert(MikrotikWebhookWorker.new(args)) do + {:ok, _job} -> + json(conn, %{status: "ok"}) + + {:error, changeset} -> + Logger.error("MikroTik webhook enqueue failed: #{inspect(changeset.errors)}") + + conn + |> put_status(:internal_server_error) + |> json(%{error: "Failed to enqueue"}) + end + else + {:error, :not_found} -> + conn + |> put_status(:not_found) + |> json(%{error: "MikroTik integration not found"}) + + {:error, :no_token_configured} -> + conn + |> put_status(:forbidden) + |> json(%{error: "Webhook token not configured"}) + + {:error, :invalid_token} -> + conn + |> put_status(:unauthorized) + |> json(%{error: "Invalid webhook token"}) + + {:error, :missing_bound} -> + conn + |> put_status(:bad_request) + |> json(%{error: "Missing or invalid 'bound' field"}) + end + end + + defp get_integration(organization_id) do + case Integrations.get_integration(organization_id, "mikrotik") do + {:ok, _integration} = ok -> ok + {:error, :not_found} -> {:error, :not_found} + end + end + + defp verify_token(conn, integration) do + expected = get_in(integration.credentials, ["webhook_token"]) + presented = conn |> Plug.Conn.get_req_header(@token_header) |> List.first() + + cond do + is_nil(expected) or expected == "" -> + Logger.warning("Mikrotik webhook: rejected - no webhook_token configured for org") + {:error, :no_token_configured} + + is_nil(presented) -> + {:error, :invalid_token} + + Plug.Crypto.secure_compare(expected, presented) -> + :ok + + true -> + {:error, :invalid_token} + end + end + + # MikroTik scripts send bound as a string ("1" / "0"). Map to assign/release. + defp parse_event_type(%{"bound" => bound}) do + case to_string(bound) do + "1" -> {:ok, "assign"} + "0" -> {:ok, "release"} + _ -> {:error, :missing_bound} + end + end + + defp parse_event_type(_), do: {:error, :missing_bound} +end diff --git a/lib/towerops_web/router.ex b/lib/towerops_web/router.ex index 314449bb..38164a0f 100644 --- a/lib/towerops_web/router.ex +++ b/lib/towerops_web/router.ex @@ -242,6 +242,8 @@ defmodule ToweropsWeb.Router do post "/gaiia/:organization_id", GaiiaWebhookController, :create post "/pagerduty/:organization_id", PagerdutyWebhookController, :create + post "/mikrotik/:organization_id/dhcp", MikrotikWebhookController, :dhcp + post "/mikrotik/:organization_id/pppoe", MikrotikWebhookController, :pppoe end # Admin API routes (requires superuser API token) diff --git a/priv/repo/migrations/20260511175355_create_mikrotik_webhook_tables.exs b/priv/repo/migrations/20260511175355_create_mikrotik_webhook_tables.exs new file mode 100644 index 00000000..42c680bd --- /dev/null +++ b/priv/repo/migrations/20260511175355_create_mikrotik_webhook_tables.exs @@ -0,0 +1,63 @@ +defmodule Towerops.Repo.Migrations.CreateMikrotikWebhookTables do + use Ecto.Migration + + def change do + create table(:mikrotik_events, primary_key: false) do + add :id, :binary_id, primary_key: true + + add :organization_id, + references(:organizations, type: :binary_id, on_delete: :delete_all), + null: false + + add :source, :string, null: false + add :event_type, :string, null: false + add :router, :string + add :mac, :string + add :ip, :string + add :username, :string + add :host, :string + add :server, :string + add :remote_id, :string + add :raw_payload, :map, null: false + add :resolved_account_id, :string + add :resolved_via, :string + add :received_at, :utc_datetime_usec, null: false + + timestamps(type: :utc_datetime, updated_at: false) + end + + create index(:mikrotik_events, [:organization_id, :received_at]) + create index(:mikrotik_events, [:organization_id, :ip]) + create index(:mikrotik_events, [:organization_id, :mac]) + create index(:mikrotik_events, [:organization_id, :username]) + + create table(:mikrotik_ip_assignments, primary_key: false) do + add :id, :binary_id, primary_key: true + + add :organization_id, + references(:organizations, type: :binary_id, on_delete: :delete_all), + null: false + + add :ip, :string, null: false + add :mac, :string + add :source, :string, null: false + add :router, :string + add :username, :string + add :host, :string + add :gaiia_account_id, :string + add :gaiia_inventory_item_id, :string + add :account_name, :string + add :resolved_via, :string + add :first_seen_at, :utc_datetime_usec, null: false + add :last_seen_at, :utc_datetime_usec, null: false + add :is_active, :boolean, default: true, null: false + + timestamps(type: :utc_datetime) + end + + create unique_index(:mikrotik_ip_assignments, [:organization_id, :ip]) + create index(:mikrotik_ip_assignments, [:organization_id, :mac]) + create index(:mikrotik_ip_assignments, [:organization_id, :gaiia_account_id]) + create index(:mikrotik_ip_assignments, [:organization_id, :is_active]) + end +end diff --git a/test/towerops/mikrotik_test.exs b/test/towerops/mikrotik_test.exs new file mode 100644 index 00000000..8c7dbe25 --- /dev/null +++ b/test/towerops/mikrotik_test.exs @@ -0,0 +1,196 @@ +defmodule Towerops.MikrotikTest do + use Towerops.DataCase, async: true + + import Towerops.AccountsFixtures + import Towerops.OrganizationsFixtures + + alias Towerops.Mikrotik + + setup do + user = user_fixture() + org = organization_fixture(user.id) + %{org: org} + end + + describe "normalize_mac/1 and format_mac/1" do + test "normalize strips separators and lowercases" do + assert Mikrotik.normalize_mac("AA:BB:CC:11:22:33") == "aabbcc112233" + assert Mikrotik.normalize_mac("aa-bb-cc-11-22-33") == "aabbcc112233" + assert Mikrotik.normalize_mac("AABBCC112233") == "aabbcc112233" + end + + test "normalize returns nil for blank or wrong-length input" do + assert Mikrotik.normalize_mac(nil) == nil + assert Mikrotik.normalize_mac("") == nil + assert Mikrotik.normalize_mac("not-a-mac") == nil + assert Mikrotik.normalize_mac("aabbcc") == nil + end + + test "format_mac round-trips with normalize_mac" do + normalized = Mikrotik.normalize_mac("aa:bb:cc:11:22:33") + assert Mikrotik.format_mac(normalized) == "AA:BB:CC:11:22:33" + end + end + + describe "record_event/1" do + test "inserts a valid assign event", %{org: org} do + {:ok, event} = + Mikrotik.record_event(%{ + organization_id: org.id, + source: "dhcp", + event_type: "assign", + mac: "aabbcc112233", + ip: "10.0.0.5", + raw_payload: %{"bound" => "1"}, + received_at: DateTime.utc_now() + }) + + assert event.source == "dhcp" + assert event.event_type == "assign" + assert event.ip == "10.0.0.5" + end + + test "rejects invalid source", %{org: org} do + {:error, changeset} = + Mikrotik.record_event(%{ + organization_id: org.id, + source: "wifi", + event_type: "assign", + raw_payload: %{}, + received_at: DateTime.utc_now() + }) + + refute changeset.valid? + assert {"is invalid", _} = changeset.errors[:source] + end + end + + describe "upsert_assignment/1" do + test "inserts a fresh row on first assign", %{org: org} do + now = DateTime.utc_now() + + {:ok, assn} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "100.64.0.10", + mac: "aabbcc112233", + source: "pppoe", + event_type: "assign", + username: "alice", + gaiia_account_id: "account_abc", + account_name: "Alice Q", + resolved_via: "mac", + first_seen_at: now, + last_seen_at: now + }) + + assert assn.is_active + assert assn.gaiia_account_id == "account_abc" + end + + test "updates the existing row on a subsequent assign for same IP", %{org: org} do + now = DateTime.utc_now() + + {:ok, first} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "100.64.0.10", + mac: "aa", + source: "pppoe", + event_type: "assign", + first_seen_at: now, + last_seen_at: now + }) + + Process.sleep(2) + + {:ok, second} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "100.64.0.10", + mac: "bb", + source: "pppoe", + event_type: "assign", + first_seen_at: now, + last_seen_at: DateTime.utc_now() + }) + + assert second.id == first.id + assert second.mac == "bb" + assert second.is_active + end + + test "release flips is_active and preserves prior resolution", %{org: org} do + now = DateTime.utc_now() + + {:ok, _} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "100.64.0.10", + mac: "aabbcc112233", + source: "pppoe", + event_type: "assign", + gaiia_account_id: "account_abc", + account_name: "Alice Q", + resolved_via: "mac", + first_seen_at: now, + last_seen_at: now + }) + + {:ok, after_release} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "100.64.0.10", + source: "pppoe", + event_type: "release", + first_seen_at: now, + last_seen_at: DateTime.utc_now() + }) + + refute after_release.is_active + # Prior resolution preserved through the release. + assert after_release.gaiia_account_id == "account_abc" + assert after_release.account_name == "Alice Q" + assert after_release.mac == "aabbcc112233" + end + end + + describe "list_active_assignments/1" do + test "returns only active rows for the org", %{org: org} do + now = DateTime.utc_now() + + {:ok, _} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "1.1.1.1", + source: "dhcp", + event_type: "assign", + first_seen_at: now, + last_seen_at: now + }) + + {:ok, _} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "2.2.2.2", + source: "dhcp", + event_type: "assign", + first_seen_at: now, + last_seen_at: now + }) + + {:ok, _} = + Mikrotik.upsert_assignment(%{ + organization_id: org.id, + ip: "2.2.2.2", + source: "dhcp", + event_type: "release", + first_seen_at: now, + last_seen_at: now + }) + + ips = org.id |> Mikrotik.list_active_assignments() |> Enum.map(& &1.ip) + assert ips == ["1.1.1.1"] + end + end +end diff --git a/test/towerops/workers/mikrotik_webhook_worker_test.exs b/test/towerops/workers/mikrotik_webhook_worker_test.exs new file mode 100644 index 00000000..c8c01de5 --- /dev/null +++ b/test/towerops/workers/mikrotik_webhook_worker_test.exs @@ -0,0 +1,120 @@ +defmodule Towerops.Workers.MikrotikWebhookWorkerTest do + use Towerops.DataCase, async: true + use Oban.Testing, repo: Towerops.Repo + + import Towerops.AccountsFixtures + import Towerops.OrganizationsFixtures + + alias Towerops.Mikrotik + alias Towerops.Workers.MikrotikWebhookWorker + + setup do + user = user_fixture() + org = organization_fixture(user.id) + %{org: org} + end + + # The worker records the event and upserts the current state even when the + # resolver can't reach Gaiia (no integration configured). resolved_account_id + # stays nil; later reconciliation can fill it in. + test "DHCP assign with no Gaiia integration records event and creates inactive resolution", %{ + org: org + } do + payload = %{ + "bound" => "1", + "mac" => "AA:BB:CC:11:22:33", + "ip" => "100.64.0.50", + "server" => "verona cgnat", + "remoteId" => "" + } + + args = %{ + "organization_id" => org.id, + "source" => "dhcp", + "event_type" => "assign", + "payload" => payload, + "received_at" => DateTime.to_iso8601(DateTime.utc_now()) + } + + assert :ok = perform_job(MikrotikWebhookWorker, args) + + [event] = Mikrotik.list_recent_events(org.id) + assert event.source == "dhcp" + assert event.event_type == "assign" + assert event.mac == "AA:BB:CC:11:22:33" + assert event.ip == "100.64.0.50" + assert event.resolved_account_id == nil + + assn = Mikrotik.get_assignment(org.id, "100.64.0.50") + assert assn.mac == "AA:BB:CC:11:22:33" + assert assn.source == "dhcp" + assert assn.is_active + assert assn.gaiia_account_id == nil + end + + test "PPPoE assign records username + caller MAC", %{org: org} do + payload = %{ + "bound" => "1", + "username" => "judydevine", + "callingstationid" => "08:62:66:3B:9F:D0", + "framedip" => "100.64.0.11" + } + + args = %{ + "organization_id" => org.id, + "source" => "pppoe", + "event_type" => "assign", + "payload" => payload, + "received_at" => DateTime.to_iso8601(DateTime.utc_now()) + } + + assert :ok = perform_job(MikrotikWebhookWorker, args) + + assn = Mikrotik.get_assignment(org.id, "100.64.0.11") + assert assn.username == "judydevine" + assert assn.mac == "08:62:66:3B:9F:D0" + assert assn.source == "pppoe" + assert assn.is_active + end + + test "release event flips is_active false and preserves prior identity", %{org: org} do + iso_now = fn -> DateTime.to_iso8601(DateTime.utc_now()) end + + assign_args = %{ + "organization_id" => org.id, + "source" => "pppoe", + "event_type" => "assign", + "payload" => %{ + "bound" => "1", + "username" => "alice", + "callingstationid" => "AA:BB:CC:11:22:33", + "framedip" => "100.64.0.99" + }, + "received_at" => iso_now.() + } + + :ok = perform_job(MikrotikWebhookWorker, assign_args) + + release_args = %{ + "organization_id" => org.id, + "source" => "pppoe", + "event_type" => "release", + "payload" => %{ + "bound" => "0", + "username" => "alice", + "callingstationid" => "AA:BB:CC:11:22:33", + "framedip" => "100.64.0.99" + }, + "received_at" => iso_now.() + } + + :ok = perform_job(MikrotikWebhookWorker, release_args) + + assn = Mikrotik.get_assignment(org.id, "100.64.0.99") + refute assn.is_active + assert assn.username == "alice" + assert assn.mac == "AA:BB:CC:11:22:33" + + assert length(Mikrotik.list_recent_events(org.id)) == 2 + end +end diff --git a/test/towerops_web/controllers/api/v1/mikrotik_webhook_controller_test.exs b/test/towerops_web/controllers/api/v1/mikrotik_webhook_controller_test.exs new file mode 100644 index 00000000..4577dd2c --- /dev/null +++ b/test/towerops_web/controllers/api/v1/mikrotik_webhook_controller_test.exs @@ -0,0 +1,159 @@ +defmodule ToweropsWeb.Api.V1.MikrotikWebhookControllerTest do + use ToweropsWeb.ConnCase, async: true + use Oban.Testing, repo: Towerops.Repo + + import Towerops.AccountsFixtures + import Towerops.OrganizationsFixtures + + alias Towerops.Integrations + alias Towerops.Workers.MikrotikWebhookWorker + + @webhook_token "test-mikrotik-token-xyz123" + + setup do + user = user_fixture() + org = organization_fixture(user.id) + + {:ok, integration} = + Integrations.create_integration(org.id, %{ + provider: "mikrotik", + enabled: true, + credentials: %{"webhook_token" => @webhook_token} + }) + + %{org: org, integration: integration} + end + + defp post_dhcp(conn, org_id, payload, token \\ @webhook_token) do + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-towerops-token", token) + |> post(~p"/api/v1/webhooks/mikrotik/#{org_id}/dhcp", Jason.encode!(payload)) + end + + defp post_pppoe(conn, org_id, payload, token \\ @webhook_token) do + conn + |> put_req_header("content-type", "application/json") + |> put_req_header("x-towerops-token", token) + |> post(~p"/api/v1/webhooks/mikrotik/#{org_id}/pppoe", Jason.encode!(payload)) + end + + describe "authentication" do + test "returns 401 without token header", %{conn: conn, org: org} do + conn = + conn + |> put_req_header("content-type", "application/json") + |> post( + ~p"/api/v1/webhooks/mikrotik/#{org.id}/dhcp", + Jason.encode!(%{"bound" => "1", "mac" => "aa:bb:cc:dd:ee:ff", "ip" => "10.0.0.1"}) + ) + + assert json_response(conn, 401)["error"] =~ "Invalid" + end + + test "returns 401 with wrong token", %{conn: conn, org: org} do + conn = post_dhcp(conn, org.id, %{"bound" => "1", "mac" => "aa", "ip" => "1.1.1.1"}, "wrong-token") + + assert json_response(conn, 401)["error"] =~ "Invalid" + end + + test "returns 404 when org has no mikrotik integration", %{conn: conn} do + other_user = user_fixture() + other_org = organization_fixture(other_user.id) + + conn = post_dhcp(conn, other_org.id, %{"bound" => "1", "mac" => "aa", "ip" => "1.1.1.1"}) + + assert json_response(conn, 404)["error"] =~ "not found" + end + end + + describe "DHCP webhook" do + test "200 on valid assign event, enqueues worker", %{conn: conn, org: org} do + payload = %{ + "bound" => "1", + "mac" => "AA:BB:CC:11:22:33", + "ip" => "100.64.0.50", + "server" => "verona cgnat", + "remoteId" => "0102030405" + } + + conn = post_dhcp(conn, org.id, payload) + + assert json_response(conn, 200) == %{"status" => "ok"} + + assert_enqueued( + worker: MikrotikWebhookWorker, + args: %{ + "organization_id" => org.id, + "source" => "dhcp", + "event_type" => "assign", + "payload" => %{ + "bound" => "1", + "mac" => "AA:BB:CC:11:22:33", + "ip" => "100.64.0.50", + "server" => "verona cgnat", + "remoteId" => "0102030405" + } + } + ) + end + + test "200 on valid release event (bound=0)", %{conn: conn, org: org} do + conn = post_dhcp(conn, org.id, %{"bound" => "0", "mac" => "aa", "ip" => "1.1.1.1"}) + + assert json_response(conn, 200) == %{"status" => "ok"} + + assert_enqueued( + worker: MikrotikWebhookWorker, + args: %{"event_type" => "release", "source" => "dhcp"} + ) + end + + test "400 when bound is missing", %{conn: conn, org: org} do + conn = post_dhcp(conn, org.id, %{"mac" => "aa", "ip" => "1.1.1.1"}) + assert json_response(conn, 400)["error"] =~ "bound" + end + end + + describe "PPPoE webhook" do + test "200 on valid assign event with PPPoE-shape payload", %{conn: conn, org: org} do + payload = %{ + "bound" => "1", + "username" => "judydevine", + "callingstationid" => "08:62:66:3B:9F:D0", + "framedip" => "100.64.0.11" + } + + conn = post_pppoe(conn, org.id, payload) + + assert json_response(conn, 200) == %{"status" => "ok"} + + assert_enqueued( + worker: MikrotikWebhookWorker, + args: %{ + "source" => "pppoe", + "event_type" => "assign", + "payload" => payload + } + ) + end + + test "passes router field through when included", %{conn: conn, org: org} do + payload = %{ + "bound" => "1", + "username" => "alice", + "callingstationid" => "aa", + "framedip" => "1.1.1.1", + "router" => "verona" + } + + conn = post_pppoe(conn, org.id, payload) + assert json_response(conn, 200) == %{"status" => "ok"} + + assert_enqueued( + worker: MikrotikWebhookWorker, + args: %{"payload" => %{"router" => "verona"}} + ) + end + end +end