feat: read what listens on which address, and which service it is, once for every driver #7

Merged
christianmanivong merged 1 commits from feat/listening-sockets into main 2026-10-06 16:19:23 +00:00
6 changed files with 543 additions and 1 deletions
+8
View File
@@ -86,6 +86,14 @@ dnf-automatic). `package_updates` holds the shared apt and dnf parsers: apt's su
an update's `origin`, a `-security` suite makes it a security update, and dnf's security
advisories do the same.
`ListeningSocketsMixin` (`get_listening_sockets`) is mixed in the same way: every listening
TCP and bound UDP socket from `ss -lntup`, with the systemd service or container behind it
from `/proc/<pid>/cgroup`, in one round trip. A driver supplies
`_run_listening_sockets_command(command, privileged=)`; the command arrives as one `sh -c`
argument, so a `sudo -n` prefix covers all of it. Without root `ss` names only the login
user's processes, and the reading says so (`attributed: false`) instead of failing. A host
without `ss` raises `ListeningSocketsUnavailable`.
**Update readers raise when they cannot read.** `get_available_updates` returns an empty
list only when nothing is pending; netOrk keeps "pending since" per package, and an empty
list for "don't know" would reset it.
+11
View File
@@ -42,6 +42,7 @@ instead of being restated on every role that happens to need it:
* :class:`~napalm_device_types.host_reboot.HostRebootMixin`
* :class:`~napalm_device_types.interface_filter.InterfaceFilterMixin`
* :class:`~napalm_device_types.kernel.KernelFactsMixin`
* :class:`~napalm_device_types.listening.ListeningSocketsMixin`
* :class:`~napalm_device_types.mac_acl.MacAclMixin`
* :class:`~napalm_device_types.nat_vpn.NatVpnMixin`
* :class:`~napalm_device_types.packages.PackageManagementMixin`
@@ -68,6 +69,12 @@ from napalm_device_types.host_reboot import HostRebootMixin
from napalm_device_types.interface_filter import InterfaceFilterMixin
from napalm_device_types.kernel import KERNEL_FACTS_COMMAND, KernelFactsMixin, parse_kernel_facts
from napalm_device_types.lag import add_lag_interfaces
from napalm_device_types.listening import (
LISTENING_SOCKETS_COMMAND,
ListeningSocketsMixin,
ListeningSocketsUnavailable,
parse_listening_sockets,
)
from napalm_device_types.mac_acl import MacAclMixin
from napalm_device_types.media import MediaDriver
from napalm_device_types.nat_vpn import NatVpnMixin
@@ -123,6 +130,10 @@ __all__ = [
"KernelFactsMixin",
"KERNEL_FACTS_COMMAND",
"parse_kernel_facts",
"LISTENING_SOCKETS_COMMAND",
"ListeningSocketsMixin",
"ListeningSocketsUnavailable",
"parse_listening_sockets",
"PhoneDriver",
"PingSweepMixin",
"PortSpec",
+218
View File
@@ -0,0 +1,218 @@
# -*- coding: utf-8 -*-
"""Listening sockets: what listens on which address, and which service it is.
Whether a service is reachable from outside its host is decided by what it
listens on -- ``0.0.0.0:5432`` is, ``127.0.0.1:5432`` is not -- and that is read
the same way on every Linux host. So the command and its parse live here once,
and a driver only carries the command across.
**One round trip.** ``ss -lntup`` lists every listening TCP and bound UDP
socket with the processes holding it; for each of those processes,
``/proc/<pid>/cgroup`` says which systemd service or container it runs in. The
report is framed, and a report whose end is missing raises: a list cut short
must never read as sockets that closed.
**Root, and without it.** Only root sees every process behind a socket.
The whole script therefore goes to the host as one ``sh -c`` argument -- a
driver that prefixes ``sudo -n`` would otherwise run only its first command as
root. When that call brings no report back (no sudo, a wrong password), the
command runs again without privilege: the sockets are still worth having, and
the reading says it is not ``attributed``.
**No ``-H``.** iproute2 before 4.10 has no option to leave out the header and
fails on it, which would read as nothing listening. The parse skips the header
instead. A host without ``ss`` at all (busybox, QNAP) raises
:class:`ListeningSocketsUnavailable`.
"""
from __future__ import annotations
import re
from shlex import quote
from typing import Dict, List, Optional, Tuple, TYPE_CHECKING
from napalm_device_types.models import ListeningSocketDict, ListeningSocketsDict
from napalm_device_types.terminal import strip_terminal_codes
_BEGIN = "SOCK_BEGIN"
_END = "SOCK_END"
_NO_SS = "no-ss"
_RC_RE = re.compile(r"^__SS_RC=(\d+)$")
#: One line, POSIX ``sh``, read-only. The frame markers are printed in two
#: halves so that a transport which echoes the command does not show them early.
#: ``ss`` is in ``/usr/sbin`` on some systems, outside a login user's ``PATH``.
LISTENING_SOCKETS_COMMAND = (
"PATH=$PATH:/usr/sbin:/sbin; "
"printf '%s%s\\n' SOCK_ BEGIN; "
"if command -v ss >/dev/null 2>&1; then "
"s=$(ss -lntup 2>&1); r=$?; echo '[ss]'; printf '%s\\n' \"$s\"; echo \"__SS_RC=$r\"; "
"echo '[cgroups]'; "
"for p in $(printf '%s\\n' \"$s\" | grep -o 'pid=[0-9]*' | cut -d= -f2 | sort -u); do "
"sed \"s|^|$p |\" /proc/$p/cgroup 2>/dev/null; done; "
"else echo '[no-ss]'; fi; "
"printf '%s%s\\n' SOCK_ END"
)
_PROTOCOLS = frozenset({"tcp", "udp"})
#: ``users:(("nginx",pid=901,fd=6),("nginx",pid=900,fd=6))``
_USER_RE = re.compile(r'\("((?:[^"\\]|\\.)*)",pid=(\d+),fd=\d+\)')
#: The service a cgroup path runs in: its deepest ``*.service`` component.
_SERVICE_RE = re.compile(r"/([^/]+)\.service(?=/|$)")
#: A container's cgroup: ``docker-<id>.scope`` (systemd driver), ``/docker/<id>`` (cgroupfs).
_CONTAINER_RE = re.compile(
r"(?:docker|libpod)-([0-9a-f]{64})\.scope|/(?:docker|libpod)/([0-9a-f]{64})(?=/|$)"
)
class ListeningSocketsUnavailable(NotImplementedError):
"""The host has no ``ss``; there is nothing to read and nothing to retry."""
def _frame(output: str) -> List[str]:
lines = [line.strip() for line in strip_terminal_codes(output).splitlines()]
try:
start = lines.index(_BEGIN)
end = lines.index(_END, start)
except ValueError:
raise ValueError("no intact listening socket report in the output") from None
return lines[start + 1 : end]
def _sections(lines: List[str]) -> Dict[str, List[str]]:
sections: Dict[str, List[str]] = {}
current: List[str] = []
for line in lines:
if line.startswith("[") and line.endswith("]") and " " not in line:
current = sections.setdefault(line[1:-1], [])
else:
current.append(line)
return sections
def _split_local(local: str) -> Optional[Tuple[str, Optional[str], int]]:
"""``[fe80::1%eth0]:546`` -> ``("fe80::1", "eth0", 546)``; None if no port."""
host, sep, port = local.rpartition(":")
if not sep or not port.isdigit():
return None
if host.startswith("[") and host.endswith("]"):
host = host[1:-1]
address, _, zone = host.partition("%")
return address or "*", zone or None, int(port)
def _cgroup_paths(lines: List[str]) -> Dict[int, str]:
"""Each process's cgroup path: the unified hierarchy, or systemd's under v1."""
paths: Dict[int, str] = {}
for line in lines:
pid, _, entry = line.partition(" ")
parts = entry.split(":", 2)
if not pid.isdigit() or len(parts) != 3:
continue
hierarchy, controllers, path = parts
if (hierarchy == "0" and controllers == "") or controllers == "name=systemd":
paths[int(pid)] = path
return paths
def _unit(path: Optional[str]) -> Optional[str]:
services = _SERVICE_RE.findall(path or "")
return services[-1] if services else None
def _container(path: Optional[str]) -> Optional[str]:
match = _CONTAINER_RE.search(path or "")
return (match.group(1) or match.group(2)) if match else None
def _socket(line: str, paths: Dict[int, str]) -> Optional[ListeningSocketDict]:
parts = line.split()
if len(parts) < 5 or parts[0] not in _PROTOCOLS:
return None
local = _split_local(parts[4])
if local is None:
return None
address, interface, port = local
users = _USER_RE.findall(line)
process, pid = (users[0][0], int(users[0][1])) if users else (None, None)
path = paths.get(pid) if pid is not None else None
return {
"proto": parts[0],
"address": address,
"port": port,
"interface": interface,
"process": process,
"pid": pid,
"unit": _unit(path),
"container_id": _container(path),
}
def parse_listening_sockets(output: str) -> List[ListeningSocketDict]:
"""Parse what :data:`LISTENING_SOCKETS_COMMAND` printed, sorted by protocol,
port and address.
:raises ListeningSocketsUnavailable: when the host has no ``ss``.
:raises ValueError: when the output carries no intact report, or ``ss`` failed.
"""
sections = _sections(_frame(output))
if _NO_SS in sections:
raise ListeningSocketsUnavailable("the host has no ss")
ss_lines = sections.get("ss", [])
statuses = [m.group(1) for m in map(_RC_RE.match, ss_lines) if m]
if not statuses or statuses[-1] != "0":
detail = " ".join(line for line in ss_lines if not _RC_RE.match(line))[:200]
raise ValueError(f"ss did not list the sockets: {detail or 'no exit status'}")
paths = _cgroup_paths(sections.get("cgroups", []))
sockets = [s for s in (_socket(line, paths) for line in ss_lines) if s is not None]
return sorted(sockets, key=lambda s: (s["proto"], s["port"], s["address"], s["interface"] or ""))
class ListeningSocketsMixin:
"""Adds :meth:`get_listening_sockets` to a driver that can run a command on a Linux host.
The template form (README, "Function classes"): the command and its parse
are the same everywhere, so they are concrete here, and a driver supplies
only :meth:`_run_listening_sockets_command` -- how a command reaches its
host, and how it gains root there. Mixed in by the drivers that can, not by
:class:`~napalm_device_types.os.OSDriver`: a Windows host is an OS driver
too, and ``hasattr(driver, "get_listening_sockets")`` has to stay truthful.
"""
if TYPE_CHECKING: # pragma: no cover - declared for type checkers only
def _run_listening_sockets_command(self, command: str, *, privileged: bool) -> str:
"""Run *command* on the host and return what it printed; as root
when *privileged*. The command is a single ``sh -c`` invocation."""
...
def get_listening_sockets(self) -> ListeningSocketsDict:
"""
Returns every listening TCP and bound UDP socket, with the process,
systemd service and container behind it.
* attributed (bool) - read as root, so every process is named
* sockets (list) - see :class:`~napalm_device_types.models.ListeningSocketDict`
Example::
{
"attributed": True,
"sockets": [
{"proto": "tcp", "address": "0.0.0.0", "port": 5432,
"interface": None, "process": "postgres", "pid": 812,
"unit": "postgresql@16-main", "container_id": None},
],
}
:raises ListeningSocketsUnavailable: if the host has no ``ss``.
:raises ValueError: if neither reading carried an intact report.
"""
command = f"sh -c {quote(LISTENING_SOCKETS_COMMAND)}"
try:
output = self._run_listening_sockets_command(command, privileged=True)
return {"attributed": True, "sockets": parse_listening_sockets(output)}
except ValueError:
pass
output = self._run_listening_sockets_command(command, privileged=False)
return {"attributed": False, "sockets": parse_listening_sockets(output)}
+31
View File
@@ -335,6 +335,37 @@ class KernelFactsDict(TypedDict):
config: Optional[Dict[str, str]] # build configuration, set options only; quotes stripped
class ListeningSocketDict(TypedDict):
"""A TCP socket that listens, or a UDP socket that is bound, on the host.
One entry per socket as ``ss`` lists it. ``address`` is printed the way ss
prints it, without brackets or zone: ``0.0.0.0`` and ``::`` are every
address of their family, ``*`` every address of both. Whether that is
reachable from outside the host is the consumer's call.
"""
proto: str # "tcp" or "udp"
address: str # "0.0.0.0", "::", "*", "127.0.0.1", "::ffff:127.0.0.1", ...
port: int
interface: Optional[str] # the %zone a socket is bound to ("lo", "eth0"), if any
process: Optional[str] # the first process holding it; None without one or without root
pid: Optional[int]
unit: Optional[str] # the process's systemd service, without ".service"
container_id: Optional[str] # the full container ID, for a container on the host network
class ListeningSocketsDict(TypedDict):
"""What ``ListeningSocketsMixin.get_listening_sockets`` read.
``attributed`` is False when the reading ran without root: ``ss`` then names
only the login user's own processes, so a socket without a process means
"not told", not "the kernel's".
"""
attributed: bool
sockets: List[ListeningSocketDict]
class PortForwardDict(TypedDict):
"""A port the WAN side can reach, forwarded to a host inside.
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "napalm-device-types"
version = "2.3.0"
version = "2.4.0"
description = "Abstract device-type base classes for NAPALM drivers"
readme = "README.md"
requires-python = ">=3.10"
+274
View File
@@ -0,0 +1,274 @@
"""get_listening_sockets: what listens on which address, and which service it is.
Whether a service is reachable from outside its host is decided by what it
listens on -- ``0.0.0.0:5432`` is, ``127.0.0.1:5432`` is not. Reading that is
the same on every Linux host: ``ss`` for the sockets, ``/proc/<pid>/cgroup``
for the systemd unit or container a process belongs to. So both live here once,
and a driver only carries the command across (#658 in netOrk).
"""
from __future__ import annotations
import shutil
import subprocess
import pytest
from napalm_device_types import OSDriver
from napalm_device_types.listening import (
LISTENING_SOCKETS_COMMAND,
ListeningSocketsMixin,
ListeningSocketsUnavailable,
parse_listening_sockets,
)
CID = "4f1c" + "0" * 60
SS = """Netid State Recv-Q Send-Q Local Address:Port Peer Address:PortProcess
udp UNCONN 0 0 127.0.0.53%lo:53 0.0.0.0:* users:(("systemd-resolve",pid=612,fd=13))
udp UNCONN 0 0 0.0.0.0%ens18:68 0.0.0.0:* users:(("systemd-network",pid=590,fd=22))
tcp LISTEN 0 4096 0.0.0.0:5432 0.0.0.0:* users:(("postgres",pid=812,fd=6))
tcp LISTEN 0 128 [::]:22 [::]:* users:(("sshd",pid=700,fd=4))
tcp LISTEN 0 511 *:80 *:* users:(("nginx",pid=901,fd=6),("nginx",pid=900,fd=6))
tcp LISTEN 0 4096 [::ffff:127.0.0.1]:8125 *:* users:(("statsd",pid=950,fd=3))
tcp LISTEN 0 4096 0.0.0.0:8080 0.0.0.0:* users:(("docker-proxy",pid=1200,fd=4))
tcp LISTEN 0 4096 0.0.0.0:9100 0.0.0.0:* users:(("node_exporter",pid=1300,fd=3))
tcp LISTEN 0 64 0.0.0.0:2049 0.0.0.0:*
"""
CGROUPS = f"""612 0::/system.slice/systemd-resolved.service
590 0::/system.slice/systemd-networkd.service
812 0::/system.slice/system-postgresql.slice/postgresql@16-main.service
700 0::/system.slice/ssh.service
900 0::/system.slice/nginx.service
901 0::/system.slice/nginx.service
950 0::/user.slice/user-1000.slice/session-3.scope
1200 0::/system.slice/docker.service
1300 0::/system.slice/docker-{CID}.scope
"""
def _wire(ss: str = SS, cgroups: str = CGROUPS, *, rc: int = 0, noise: str = "") -> str:
"""The report as the command prints it, framed."""
return f"{noise}SOCK_BEGIN\n[ss]\n{ss}__SS_RC={rc}\n[cgroups]\n{cgroups}SOCK_END\n"
def _by_port(sockets: list) -> dict:
return {(s["proto"], s["port"], s["address"]): s for s in sockets}
class TestParsing:
def test_a_socket_comes_with_its_process_and_unit(self):
sockets = _by_port(parse_listening_sockets(_wire()))
assert sockets[("tcp", 5432, "0.0.0.0")] == {
"proto": "tcp",
"address": "0.0.0.0",
"port": 5432,
"interface": None,
"process": "postgres",
"pid": 812,
"unit": "postgresql@16-main",
"container_id": None,
}
@pytest.mark.parametrize(
"local, address, interface",
[
("[::]:22", "::", None),
("*:80", "*", None),
("127.0.0.53%lo:53", "127.0.0.53", "lo"),
("0.0.0.0%ens18:68", "0.0.0.0", "ens18"),
("[::ffff:127.0.0.1]:8125", "::ffff:127.0.0.1", None),
("[fe80::1%eth0]:546", "fe80::1", "eth0"),
(":::22", "::", None), # iproute2 4.9 prints IPv6 without brackets
],
)
def test_every_address_form(self, local: str, address: str, interface):
line = f"tcp LISTEN 0 128 {local} *:* users:((\"x\",pid=1,fd=3))\n"
[socket] = parse_listening_sockets(_wire(line, "1 0::/system.slice/x.service\n"))
assert (socket["address"], socket["interface"]) == (address, interface)
def test_the_first_of_several_processes_names_the_socket(self):
nginx = _by_port(parse_listening_sockets(_wire()))[("tcp", 80, "*")]
assert (nginx["process"], nginx["pid"], nginx["unit"]) == ("nginx", 901, "nginx")
def test_a_socket_without_a_process_is_kept(self):
"""The kernel's nfsd, or every socket of another user without root."""
nfs = _by_port(parse_listening_sockets(_wire()))[("tcp", 2049, "0.0.0.0")]
assert (nfs["process"], nfs["pid"], nfs["unit"], nfs["container_id"]) == (
None,
None,
None,
None,
)
def test_the_header_is_no_socket(self):
sockets = parse_listening_sockets(_wire())
assert len(sockets) == 9
assert all(s["proto"] in ("tcp", "udp") for s in sockets)
def test_sorted_by_protocol_port_and_address(self):
keys = [(s["proto"], s["port"], s["address"]) for s in parse_listening_sockets(_wire())]
assert keys == sorted(keys)
def test_nothing_listening_is_an_empty_list(self):
assert parse_listening_sockets(_wire(SS.splitlines(True)[0], "")) == []
class TestCgroups:
def _unit_and_container(self, cgroup_lines: str) -> tuple:
line = 'tcp LISTEN 0 128 0.0.0.0:1 0.0.0.0:* users:(("p",pid=5,fd=3))\n'
[socket] = parse_listening_sockets(_wire(line, cgroup_lines))
return socket["unit"], socket["container_id"]
def test_a_container_on_the_host_network(self):
assert self._unit_and_container(f"5 0::/system.slice/docker-{CID}.scope\n") == (
None,
CID,
)
def test_a_container_under_the_cgroupfs_driver(self):
assert self._unit_and_container(f"5 0::/docker/{CID}\n") == (None, CID)
def test_cgroup_v1_reads_the_systemd_hierarchy(self):
lines = "5 12:pids:/user.slice\n5 1:name=systemd:/system.slice/ssh.service\n"
assert self._unit_and_container(lines) == ("ssh", None)
def test_an_escaped_template_instance(self):
lines = "5 0::/system.slice/system-wg\\x2dquick.slice/wg-quick@wg0.service\n"
assert self._unit_and_container(lines) == ("wg-quick@wg0", None)
def test_a_login_session_is_no_unit(self):
assert self._unit_and_container("5 0::/user.slice/user-1000.slice/session-3.scope\n") == (
None,
None,
)
def test_a_process_gone_before_its_cgroup_was_read(self):
assert self._unit_and_container("") == (None, None)
def test_the_docker_daemons_own_proxy_stays_raw(self):
"""docker-proxy runs in docker.service. Telling it apart from a real
listener is the consumer's business; the reading reports what is there."""
proxy = _by_port(parse_listening_sockets(_wire()))[("tcp", 8080, "0.0.0.0")]
assert (proxy["process"], proxy["unit"]) == ("docker-proxy", "docker")
class TestFailures:
def test_whatever_surrounds_the_frame_is_ignored(self):
noisy = _wire(noise="$ sh -c '...'\nWelcome to Ubuntu\n")
assert len(parse_listening_sockets(noisy)) == 9
def test_output_without_the_frame_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets("sudo: a password is required\n")
def test_a_report_cut_short_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets(_wire()[:-len("SOCK_END\n")])
def test_a_host_without_ss_is_unavailable(self):
with pytest.raises(ListeningSocketsUnavailable):
parse_listening_sockets("SOCK_BEGIN\n[no-ss]\nSOCK_END\n")
def test_ss_failing_raises(self):
with pytest.raises(ValueError):
parse_listening_sockets(_wire("ss: invalid option -- 'p'\n", "", rc=1))
def test_unavailable_is_not_a_value_error(self):
"""A caller retries a broken report without privilege; a host without
ss has nothing to retry."""
assert not issubclass(ListeningSocketsUnavailable, ValueError)
class TestTheCommand:
def test_the_frame_is_not_in_the_command_itself(self):
"""An echoing transport prints the command back; the markers must only
appear once the command has run."""
assert "SOCK_BEGIN" not in LISTENING_SOCKETS_COMMAND
assert "SOCK_END" not in LISTENING_SOCKETS_COMMAND
def test_it_writes_and_changes_nothing(self):
for verb in ("sudo", " > ", ">>", "kill", "rm ", "systemctl "):
assert verb not in LISTENING_SOCKETS_COMMAND
def test_it_is_posix_sh(self):
result = subprocess.run(
["sh", "-n", "-c", LISTENING_SOCKETS_COMMAND], capture_output=True, text=True
)
assert result.returncode == 0, result.stderr
@pytest.mark.skipif(shutil.which("ss") is None, reason="needs ss on this host")
def test_it_runs_and_parses_on_this_host(self):
out = subprocess.run(
["sh", "-c", LISTENING_SOCKETS_COMMAND], capture_output=True, text=True, timeout=60
).stdout
sockets = parse_listening_sockets(out)
assert all(s["proto"] in ("tcp", "udp") and 0 < s["port"] < 65536 for s in sockets)
class _Driver(ListeningSocketsMixin):
def __init__(self, privileged: str, unprivileged: str = "") -> None:
self.answers = {True: privileged, False: unprivileged}
self.calls: list = []
def _run_listening_sockets_command(self, command: str, *, privileged: bool) -> str:
self.calls.append((command, privileged))
return self.answers[privileged]
class TestTheTemplate:
def test_a_driver_supplies_only_the_transport(self):
driver = _Driver(_wire())
reading = driver.get_listening_sockets()
assert reading["attributed"] is True
assert len(reading["sockets"]) == 9
[(command, privileged)] = driver.calls
assert privileged is True
def test_the_command_reaches_root_as_one_shell(self):
"""``sudo -n a; b`` runs only ``a`` as root: the whole script goes as
one ``sh -c`` argument."""
driver = _Driver(_wire())
driver.get_listening_sockets()
[(command, _privileged)] = driver.calls
assert command.startswith("sh -c '")
assert subprocess.run(["sh", "-n", "-c", command]).returncode == 0
def test_without_root_it_reads_what_it_can(self):
"""Without sudo, ss names only the user's own processes: the sockets
are still worth having, marked as not attributed."""
driver = _Driver("sudo: a password is required\n", _wire(cgroups=""))
reading = driver.get_listening_sockets()
assert reading["attributed"] is False
assert len(reading["sockets"]) == 9
assert [p for _c, p in driver.calls] == [True, False]
def test_a_host_without_ss_is_not_asked_twice(self):
driver = _Driver("SOCK_BEGIN\n[no-ss]\nSOCK_END\n")
with pytest.raises(ListeningSocketsUnavailable):
driver.get_listening_sockets()
assert len(driver.calls) == 1
def test_not_every_os_driver_has_it(self):
"""A Windows host is an OSDriver too and has no ss: ``hasattr`` has to
stay a truthful answer, so the drivers that can mix this in themselves."""
assert not issubclass(OSDriver, ListeningSocketsMixin)
assert not hasattr(OSDriver, "get_listening_sockets")