# -*- coding: utf-8 -*- # Licensed under the Apache License, Version 2.0 (the "License"); # you may not use this file except in compliance with the License. # You may obtain a copy of the License at # # http://www.apache.org/licenses/LICENSE-2.0 # # Unless required by applicable law or agreed to in writing, software # distributed under the License is distributed on an "AS IS" BASIS, # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. """Ping support for OPNsense — NAPALM ``ping()`` plus a batched ``ping_sweep()``. OPNsense has no synchronous "ping once and tell me the result" endpoint. ``/api/diagnostics/ping`` is a *job* API: a job is created (``set``), started (``start``), then runs in the background writing statistics that are read back via ``search_jobs``, and finally has to be stopped and removed again. A single ping therefore costs five API calls and roughly one second of waiting. That makes the generic sequential sweep in ``napalm-device-types`` far too slow here, but it also hands us something better: jobs are independent and run on the firewall in parallel, and ``search_jobs`` reports *all* of them in one response. :meth:`OPNsensePingMixin.ping_sweep` therefore works in batches — create and start a whole batch, wait once, read every result with a single request, then clean the batch up. Cost per batch is constant in waiting time instead of linear in hosts. ``PING_SWEEP_BATCH_SIZE`` bounds how many ping processes the firewall is asked to run at once; ``PING_SWEEP_MAX_TARGETS`` bounds the sweep as a whole. Both are deliberately conservative — this runs on production firewalls. Why the results are *polled* rather than read once: ``search_jobs`` signals the running ping with ``SIGINFO`` and then parses whatever the process has written to its log so far. The statistics line therefore lands in the log slightly after the request that triggered it, so the first read of a healthy host can still show zero probes. Every job is asked repeatedly until it reports the requested probe count or the time budget runs out. Not supported by the API and therefore ignored: ``ttl``, ``vrf``. """ from __future__ import annotations import logging import time from ipaddress import ip_address from typing import Any, Callable, Iterable, Optional logger = logging.getLogger(__name__) _PING_API = "/api/diagnostics/ping" #: Where the ping fields live when the model cannot be read. Matches what #: OPNsense 24/25 ship: {"ping": {"settings": {...}}}. _DEFAULT_PING_MODEL_PATH = ("ping", "settings") #: Seconds between two ``search_jobs`` polls while waiting for probe results. _POLL_INTERVAL = 0.5 #: The API's ``interval`` field is in whole seconds and has a minimum of 1, #: so one probe takes at least this long. _PROBE_INTERVAL_SECONDS = 1 #: Extra grace on top of the expected probe duration before giving up on a job. _POLL_GRACE_SECONDS = 1.0 class OPNsensePingMixin: """NAPALM ``ping()`` / ``ping_sweep()`` on the OPNsense diagnostics ping API.""" #: Ping jobs started on the firewall at the same time. PING_SWEEP_BATCH_SIZE: int = 32 #: Hosts probed in a single sweep. A /24 fits; a /16 must be split by #: the caller — 65k background jobs is not something to point at a firewall. PING_SWEEP_MAX_TARGETS: int = 512 # Provided by the driver this mixin is mixed into. _ping_model_path_cache: Optional[tuple[str, ...]] = None # ── NAPALM ─────────────────────────────────────────────────────────────── def ping( self, destination: str, source: str = "", ttl: int = 255, timeout: int = 2, size: int = 100, count: int = 5, vrf: str = "", ) -> dict[str, Any]: """Ping *destination* from the firewall. Runs the full job lifecycle and returns NAPALM's standard ping shape: ``{"success": {...}}`` with probe counts and RTTs, or ``{"error": ...}`` if the job could not be created or its statistics could not be read. """ try: job_id = self._ping_job_create(destination, source=source, size=size) except Exception as exc: # noqa: BLE001 - reported to the caller as-is return {"error": str(exc)} try: self._post(f"{_PING_API}/start/{job_id}", {}) rows = self._await_ping_rows([job_id], count=count, timeout=timeout) except Exception as exc: # noqa: BLE001 return {"error": str(exc)} finally: self._ping_job_cleanup([job_id]) return self._row_to_napalm_reply(destination, rows.get(job_id), count=count) def ping_sweep( self, destinations: Iterable[str], *, count: int = 1, timeout: int = 1, max_targets: Optional[int] = None, on_progress: Optional[Callable[[int, int], None]] = None, should_stop: Optional[Callable[[], bool]] = None, ) -> dict[str, Any]: """Ping many destinations, one batch of parallel firewall jobs at a time. Same contract as :meth:`napalm_device_types.ping_sweep.PingSweepMixin.ping_sweep`; ``should_stop`` is polled between batches rather than between hosts, because a batch is the smallest unit of work here. """ limit = self.PING_SWEEP_MAX_TARGETS if max_targets is None else max_targets targets = list(destinations) truncated = len(targets) > limit if truncated: targets = targets[:limit] total = len(targets) entries: list[dict[str, Any]] = [] for start in range(0, total, self.PING_SWEEP_BATCH_SIZE): if should_stop is not None and should_stop(): break batch = targets[start : start + self.PING_SWEEP_BATCH_SIZE] entries.extend(self._sweep_batch(batch, count=count, timeout=timeout)) if on_progress is not None: on_progress(len(entries), total) return { "entries": entries, "scanned": len(entries), "alive_count": sum(1 for entry in entries if entry["alive"]), "truncated": truncated, } # ── Batch execution ────────────────────────────────────────────────────── def _sweep_batch( self, destinations: list[str], *, count: int, timeout: int ) -> list[dict[str, Any]]: """Create, start, harvest and remove one batch of ping jobs.""" job_ids: dict[str, str] = {} # destination -> job id errors: dict[str, str] = {} # destination -> creation error for destination in destinations: try: job_ids[destination] = self._ping_job_create(destination) except Exception as exc: # noqa: BLE001 - one bad host must not kill the batch errors[destination] = str(exc) rows: dict[str, dict[str, Any]] = {} batch_error: str | None = None try: for job_id in job_ids.values(): self._post(f"{_PING_API}/start/{job_id}", {}) rows = self._await_ping_rows(list(job_ids.values()), count=count, timeout=timeout) except Exception as exc: # noqa: BLE001 - the whole batch loses its results batch_error = str(exc) finally: self._ping_job_cleanup(list(job_ids.values())) entries: list[dict[str, Any]] = [] for destination in destinations: if destination in errors: entries.append(self._parse_ping_reply(destination, {"error": errors[destination]})) elif batch_error is not None: entries.append(self._parse_ping_reply(destination, {"error": batch_error})) else: reply = self._row_to_napalm_reply( destination, rows.get(job_ids[destination]), count=count ) entries.append(self._parse_ping_reply(destination, reply)) return entries # ── Job lifecycle ──────────────────────────────────────────────────────── def _ping_job_create(self, destination: str, *, source: str = "", size: int = 100) -> str: """Create a ping job for *destination* and return its job id.""" settings = { "hostname": destination, "fam": self._address_family(destination), "packetsize": str(size), "interval": str(_PROBE_INTERVAL_SECONDS), "description": "netork ping sweep", } if source: settings["source_address"] = source response = self._post(f"{_PING_API}/set", self._wrap_in_model(settings)) job_id = response.get("uuid") or response.get("id") if not job_id: raise RuntimeError(f"ping job creation failed for {destination}: {response}") return str(job_id) def _ping_job_cleanup(self, job_ids: list[str]) -> None: """Stop and remove *job_ids*, never raising — cleanup is best effort. Every job left behind keeps a ping process and a log file in ``/tmp/ping`` on the firewall, so this runs even when the sweep itself failed. """ for job_id in job_ids: for action in ("stop", "remove"): try: self._post(f"{_PING_API}/{action}/{job_id}", {}) except Exception as exc: # noqa: BLE001 logger.debug("ping job %s: %s failed: %s", job_id, action, exc) def _await_ping_rows( self, job_ids: list[str], *, count: int, timeout: int ) -> dict[str, dict[str, Any]]: """Poll ``search_jobs`` until every job sent *count* probes, or time is up. Returns the last rows seen, keyed by job id — a job that never reported is simply absent, which the caller reads as "no answer". """ if not job_ids: return {} wanted = set(job_ids) budget = max(timeout, 1) * max(count, 1) + _POLL_GRACE_SECONDS attempts = max(1, int(budget / _POLL_INTERVAL)) rows: dict[str, dict[str, Any]] = {} for attempt in range(attempts): if attempt: time.sleep(_POLL_INTERVAL) rows = { job_id: row for job_id, row in self._read_ping_jobs().items() if job_id in wanted } if wanted and all( _as_int(rows.get(job_id, {}).get("send")) >= count for job_id in wanted ): break return rows def _read_ping_jobs(self) -> dict[str, dict[str, Any]]: """All ping jobs currently known to the firewall, keyed by job id.""" response = self._get(f"{_PING_API}/search_jobs") rows = response.get("rows") or [] result: dict[str, dict[str, Any]] = {} for row in rows: job_id = row.get("uuid") or row.get("id") if job_id: result[str(job_id)] = row return result def _wrap_in_model(self, settings: dict[str, Any]) -> dict[str, Any]: """Nest *settings* under the keys ``set`` expects, outermost last.""" payload: dict[str, Any] = settings for key in reversed(self._ping_model_path()): payload = {key: payload} return payload def _ping_model_path(self) -> tuple[str, ...]: """Keys from the model root down to the fields, e.g. ``("ping", "settings")``. Read once from ``GET /api/diagnostics/ping/get`` rather than hardcoded, so a model rename in a future release does not silently break job creation. The response mirrors the model, so the path is the chain of single-dict wrappers around the fields:: {"ping": {"settings": {"hostname": "", "fam": {...}, ...}}} The descent stops at the first level that holds more than one key — that is the field level, where ``fam`` is a dict too and following it would land inside a form field. Getting this wrong is expensive and quiet: OPNsense answers a job whose hostname arrived at the wrong node with ``ping.settings.hostname: A value is required``, and a sweep turns 254 such refusals into a network that appears to hold nothing. """ if self._ping_model_path_cache is not None: return self._ping_model_path_cache path = _DEFAULT_PING_MODEL_PATH try: node = self._get(f"{_PING_API}/get") found: list[str] = [] while isinstance(node, dict) and len(node) == 1: key, value = next(iter(node.items())) if not isinstance(value, dict): break found.append(key) node = value if found: path = tuple(found) except Exception as exc: # noqa: BLE001 - the default is what firmware ships logger.debug("Could not read ping model path, assuming %r: %s", path, exc) self._ping_model_path_cache = path return path # ── Parsing ────────────────────────────────────────────────────────────── @staticmethod def _address_family(destination: str) -> str: """``"ip"`` for IPv4 / hostnames, ``"ip6"`` for IPv6 literals.""" try: return "ip6" if ip_address(destination).version == 6 else "ip" except ValueError: return "ip" @staticmethod def _row_to_napalm_reply( destination: str, row: Optional[dict[str, Any]], *, count: int ) -> dict[str, Any]: """Turn one ``search_jobs`` row into NAPALM's ping reply shape. A missing row, or one that never sent a probe, is reported as "all probes lost" rather than as an error: from the caller's point of view an unreachable host and a host the firewall never got around to pinging look the same, and the sweep only asks "did it answer?". """ if row is None: return _all_lost(count) sent = _as_int(row.get("send")) received = _as_int(row.get("received")) if sent == 0: return _all_lost(count) avg = _as_float(row.get("avg")) return { "success": { "probes_sent": sent, "packet_loss": max(sent - received, 0), "rtt_min": _as_float(row.get("min")), "rtt_avg": avg, "rtt_max": _as_float(row.get("max")), "rtt_stddev": _as_float(row.get("std-dev")), "results": [{"ip_address": destination, "rtt": avg}] if received else [], } } def _all_lost(count: int) -> dict[str, Any]: probes = max(count, 1) return { "success": { "probes_sent": probes, "packet_loss": probes, "rtt_min": 0.0, "rtt_avg": 0.0, "rtt_max": 0.0, "rtt_stddev": 0.0, "results": [], } } def _as_int(value: Any, default: int = 0) -> int: try: return int(value) except (TypeError, ValueError): return default def _as_float(value: Any, default: float = 0.0) -> float: try: return float(value) except (TypeError, ValueError): return default