#!/usr/bin/env python3 """ Prometheus exporter for linuxptp nodes. Designed for: * GNSS/PPS grandmasters (ptp4l + phc2sys + optional chrony) * Linux boundary clocks (including boundary_clock_jbod) * Ordinary-clock / downstream Linux clients The exporter reads ptp4l through its read-only Unix-domain management socket (/var/run/ptp4lro by default) using pmc. It does not modify PTP state. Optional data sources: * chronyc tracking / sources - useful on a GNSS-backed grandmaster * systemd service state * recent phc2sys journal output, to expose the PHC<->system/PHC servo offset """ from __future__ import annotations import argparse import math import re import socket import subprocess import time import os import threading from typing import Dict, Iterable, List, Optional, Tuple from prometheus_client import REGISTRY, start_http_server from prometheus_client.core import ( CounterMetricFamily, GaugeMetricFamily, InfoMetricFamily, ) PMC_DATASETS = [ "DEFAULT_DATA_SET", "CURRENT_DATA_SET", "PARENT_DATA_SET", "TIME_PROPERTIES_DATA_SET", "TIME_STATUS_NP", "PORT_PROPERTIES_NP", "PORT_STATS_NP", "PORT_SERVICE_STATS_NP", ] PMC_HEADER_RE = re.compile( r"^\s*(?P\S+)\s+seq\s+\d+\s+RESPONSE\s+MANAGEMENT\s+(?P[A-Z0-9_]+)\s*$" ) PMC_KV_RE = re.compile(r"^\s*(?P[A-Za-z0-9_.-]+)\s+(?P.+?)\s*$") PHC2SYS_RE = re.compile( r"(?P[A-Za-z0-9_.:/-]+)\s+phc\s+offset\s+" r"(?P[+-]?\d+)\s+\S+\s+freq\s+(?P[+-]?\d+)" r"(?:\s+delay\s+(?P\d+))?" ) CHRONY_SOURCE_RE = re.compile( r"^\s*(?P[\^=#])(?P[*+\-?x~])\s+(?P\S+)" ) def run_command(argv: List[str], timeout: float = 3.0) -> str: proc = subprocess.run( argv, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, timeout=timeout, check=False, ) if proc.returncode != 0: raise RuntimeError( f"command failed ({proc.returncode}): {' '.join(argv)}\n{proc.stdout.strip()}" ) return proc.stdout def parse_scalar(value: str): value = value.strip() low = value.lower() if low == "true": return 1.0 if low == "false": return 0.0 try: if low.startswith("0x"): return float(int(low, 16)) return float(value) except ValueError: return value def parse_pmc(text: str) -> List[Dict[str, object]]: """Parse pmc's multi-response textual output into response dictionaries.""" blocks: List[Dict[str, object]] = [] current: Optional[Dict[str, object]] = None for line in text.splitlines(): if line.startswith("sending:"): continue match = PMC_HEADER_RE.match(line) if match: current = { "_source": match.group("source"), "_dataset": match.group("dataset"), } blocks.append(current) continue if current is None: continue match = PMC_KV_RE.match(line) if not match: continue current[match.group("key")] = parse_scalar(match.group("value")) return blocks def first_block(blocks: Iterable[Dict[str, object]], dataset: str) -> Optional[Dict[str, object]]: for block in blocks: if block.get("_dataset") == dataset: return block return None def blocks_for(blocks: Iterable[Dict[str, object]], dataset: str) -> List[Dict[str, object]]: return [b for b in blocks if b.get("_dataset") == dataset] def as_float(value, default: float = math.nan) -> float: if isinstance(value, (int, float)): return float(value) if isinstance(value, str): try: if value.lower().startswith("0x"): return float(int(value, 16)) return float(value) except ValueError: pass return default def as_str(value, default: str = "") -> str: if value is None: return default return str(value) def normalize_port_state(state: str) -> str: # linuxptp releases may use SLAVE/MASTER or CLIENT/SERVER terminology. state = state.strip().upper() aliases = { "SLAVE": "client", "CLIENT": "client", "MASTER": "server", "SERVER": "server", "GRAND_MASTER": "server", "PASSIVE": "passive", "LISTENING": "listening", "UNCALIBRATED": "uncalibrated", "FAULTY": "faulty", "DISABLED": "disabled", "INITIALIZING": "initializing", "PRE_MASTER": "pre_master", } return aliases.get(state, state.lower()) def parse_chrony_tracking(text: str) -> Dict[str, object]: result: Dict[str, object] = {} for line in text.splitlines(): if ":" not in line: continue key, raw = [part.strip() for part in line.split(":", 1)] result[key] = raw def first_number(key: str) -> Optional[float]: raw = result.get(key) if raw is None: return None m = re.search(r"[+-]?(?:\d+(?:\.\d*)?|\.\d+)(?:[eE][+-]?\d+)?", str(raw)) return float(m.group(0)) if m else None parsed: Dict[str, object] = {} numeric_map = { "Stratum": "stratum", "Last offset": "last_offset_seconds", "RMS offset": "rms_offset_seconds", "Residual freq": "residual_frequency_ppm", "Skew": "skew_ppm", "Root delay": "root_delay_seconds", "Root dispersion": "root_dispersion_seconds", "Update interval": "update_interval_seconds", } for chrony_key, out_key in numeric_map.items(): value = first_number(chrony_key) if value is not None: parsed[out_key] = value # Frequency has a direction word: "12.3 ppm fast" or "slow". freq = first_number("Frequency") if freq is not None: raw = str(result.get("Frequency", "")).lower() if "slow" in raw: freq = -abs(freq) elif "fast" in raw: freq = abs(freq) parsed["frequency_ppm"] = freq # "System time : 0.000001 seconds fast/slow of NTP time" system_time = first_number("System time") if system_time is not None: raw = str(result.get("System time", "")).lower() # Positive = local system clock is ahead/fast, negative = behind/slow. if "slow" in raw: system_time = -abs(system_time) elif "fast" in raw: system_time = abs(system_time) parsed["system_time_offset_seconds"] = system_time if "Leap status" in result: parsed["leap_status"] = str(result["Leap status"]) if "Reference ID" in result: parsed["reference_id"] = str(result["Reference ID"]) return parsed def parse_chrony_sources(text: str) -> List[Dict[str, str]]: rows: List[Dict[str, str]] = [] for line in text.splitlines(): m = CHRONY_SOURCE_RE.match(line) if m: rows.append(m.groupdict()) return rows def systemd_active(service: str) -> float: try: out = run_command(["systemctl", "is-active", service], timeout=2.0).strip() return 1.0 if out == "active" else 0.0 except Exception: return 0.0 def read_phc2sys_journal(service: str) -> Dict[str, Dict[str, float]]: """ Return the latest phc2sys servo sample per destination clock found in the last journal lines. Values are nanoseconds / ppb as printed by phc2sys. """ text = run_command( [ "journalctl", "-u", service, "-n", "250", "--no-pager", "-o", "cat", ], timeout=3.0, ) latest: Dict[str, Dict[str, float]] = {} for line in text.splitlines(): m = PHC2SYS_RE.search(line) if not m: continue row = { "offset_ns": float(m.group("offset")), "frequency_ppb": float(m.group("freq")), } if m.group("delay") is not None: row["delay_ns"] = float(m.group("delay")) latest[m.group("clock")] = row return latest class LinuxPTPCollector: def __init__( self, role: str, uds: str, pmc_binary: str, command_timeout: float, chrony: bool, ptp4l_service: Optional[str], phc2sys_service: Optional[str], pmc_runtime_dir: str, ): self.role = role self.uds = uds self.pmc_binary = pmc_binary self.command_timeout = command_timeout self.chrony = chrony self.ptp4l_service = ptp4l_service self.phc2sys_service = phc2sys_service self.pmc_runtime_dir = pmc_runtime_dir self.hostname = socket.gethostname() def _pmc(self): """ Query each management dataset independently. linuxptp versions differ in which *_NP management IDs they expose, and the read-only UDS may intentionally not return some port-scoped data. One failed optional query must therefore not discard valid clock data. """ blocks: List[Dict[str, object]] = [] status: Dict[str, float] = {} errors: Dict[str, str] = {} for dataset in PMC_DATASETS: # `-s` is the ptp4l SERVER socket. `pmc` also needs its own # CLIENT-side Unix socket. Without `-i`, pmc defaults to # /var/run/pmc.$pid, which an unprivileged exporter cannot create. # Use our systemd-owned RuntimeDirectory instead. local_uds = os.path.join( self.pmc_runtime_dir, f"pmc.{os.getpid()}.{threading.get_ident()}.{time.time_ns()}", ) argv = [ self.pmc_binary, "-u", "-i", local_uds, "-s", self.uds, "-b", "0", f"GET {dataset}", ] try: output = run_command(argv, timeout=self.command_timeout) parsed = parse_pmc(output) if parsed: blocks.extend(parsed) status[dataset] = 1.0 else: status[dataset] = 0.0 errors[dataset] = "no_response" except Exception as exc: status[dataset] = 0.0 errors[dataset] = type(exc).__name__ return blocks, status, errors def collect(self): start = time.monotonic() scrape_success = 1.0 error_stage = "" try: blocks, dataset_status, dataset_errors = self._pmc() core_datasets = { "DEFAULT_DATA_SET", "CURRENT_DATA_SET", "PARENT_DATA_SET", "TIME_PROPERTIES_DATA_SET", "TIME_STATUS_NP", } core_ok = any( dataset_status.get(dataset, 0.0) == 1.0 for dataset in core_datasets ) if not core_ok: scrape_success = 0.0 error_stage = "pmc" except Exception: blocks = [] dataset_status = {} dataset_errors = {} scrape_success = 0.0 error_stage = "pmc" role_info = InfoMetricFamily( "ptp_linux_node", "Linux PTP exporter node information", labels=["hostname", "role"], ) role_info.add_metric([self.hostname, self.role], {}) yield role_info # -------- Core datasets -------- current = first_block(blocks, "CURRENT_DATA_SET") if current: metric = GaugeMetricFamily( "ptp_master_offset_nanoseconds", "Current PTP offset from the selected immediate master", ) metric.add_metric([], as_float(current.get("offsetFromMaster"))) yield metric metric = GaugeMetricFamily( "ptp_mean_path_delay_nanoseconds", "Current estimated mean path delay to the immediate master", ) metric.add_metric([], as_float(current.get("meanPathDelay"))) yield metric metric = GaugeMetricFamily( "ptp_steps_removed", "Number of boundary-clock steps from the grandmaster", ) metric.add_metric([], as_float(current.get("stepsRemoved"))) yield metric time_status = first_block(blocks, "TIME_STATUS_NP") if time_status: metric = GaugeMetricFamily( "ptp_gm_present", "Whether linuxptp reports a grandmaster as present", ) metric.add_metric([], as_float(time_status.get("gmPresent"), 0.0)) yield metric # TIME_STATUS_NP master_offset is useful too, and often an integer. metric = GaugeMetricFamily( "ptp_time_status_master_offset_nanoseconds", "Master offset from linuxptp TIME_STATUS_NP", ) metric.add_metric([], as_float(time_status.get("master_offset"))) yield metric gm = as_str(time_status.get("gmIdentity")) if gm: info = InfoMetricFamily( "ptp_grandmaster", "Currently selected PTP grandmaster identity", labels=["identity"], ) info.add_metric([gm], {}) yield info parent = first_block(blocks, "PARENT_DATA_SET") if parent: info = InfoMetricFamily( "ptp_parent", "Immediate parent and grandmaster identity", labels=["parent_port_identity", "grandmaster_identity"], ) info.add_metric( [ as_str(parent.get("parentPortIdentity")), as_str(parent.get("grandmasterIdentity")), ], {}, ) yield info gauges = { "ptp_grandmaster_priority1": ("grandmasterPriority1", "Grandmaster priority1"), "ptp_grandmaster_priority2": ("grandmasterPriority2", "Grandmaster priority2"), "ptp_grandmaster_clock_class": ("gm.ClockClass", "Grandmaster clockClass"), "ptp_grandmaster_clock_accuracy_code": ( "gm.ClockAccuracy", "Grandmaster IEEE 1588 clockAccuracy code", ), "ptp_grandmaster_offset_scaled_log_variance": ( "gm.OffsetScaledLogVariance", "Grandmaster offsetScaledLogVariance", ), } for metric_name, (field, help_text) in gauges.items(): if field in parent: g = GaugeMetricFamily(metric_name, help_text) g.add_metric([], as_float(parent.get(field))) yield g default = first_block(blocks, "DEFAULT_DATA_SET") if default: info = InfoMetricFamily( "ptp_local_clock", "Local linuxptp clock identity", labels=["identity"], ) info.add_metric([as_str(default.get("clockIdentity"))], {}) yield info for metric_name, field, help_text in [ ("ptp_domain_number", "domainNumber", "PTP domain number"), ("ptp_local_clock_class", "clockClass", "Local clockClass"), ("ptp_local_clock_accuracy_code", "clockAccuracy", "Local clockAccuracy code"), ("ptp_local_priority1", "priority1", "Local priority1"), ("ptp_local_priority2", "priority2", "Local priority2"), ("ptp_local_number_ports", "numberPorts", "Number of PTP ports"), ]: if field in default: g = GaugeMetricFamily(metric_name, help_text) g.add_metric([], as_float(default.get(field))) yield g time_props = first_block(blocks, "TIME_PROPERTIES_DATA_SET") if time_props: fields = [ ("ptp_current_utc_offset_seconds", "currentUtcOffset", "PTP TAI-UTC offset"), ("ptp_current_utc_offset_valid", "currentUtcOffsetValid", "UTC offset valid flag"), ("ptp_timescale", "ptpTimescale", "PTP timescale flag"), ("ptp_time_traceable", "timeTraceable", "PTP time traceable flag"), ("ptp_frequency_traceable", "frequencyTraceable", "PTP frequency traceable flag"), ("ptp_leap61", "leap61", "Positive leap-second flag"), ("ptp_leap59", "leap59", "Negative leap-second flag"), ("ptp_time_source_code", "timeSource", "IEEE 1588 timeSource code"), ] for metric_name, field, help_text in fields: if field in time_props: g = GaugeMetricFamily(metric_name, help_text) g.add_metric([], as_float(time_props.get(field))) yield g # -------- Port properties and counters -------- port_properties = blocks_for(blocks, "PORT_PROPERTIES_NP") port_to_interface: Dict[str, str] = {} port_state = GaugeMetricFamily( "ptp_port_state", "Current PTP port state; one sample with value 1 is emitted per port", labels=["interface", "port_identity", "state"], ) timestamping = GaugeMetricFamily( "ptp_port_timestamping", "Timestamping mode reported by linuxptp", labels=["interface", "port_identity", "mode"], ) for port in port_properties: identity = as_str(port.get("portIdentity")) interface = as_str(port.get("interface"), identity) port_to_interface[identity] = interface state = normalize_port_state(as_str(port.get("portState"))) port_state.add_metric([interface, identity, state], 1.0) timestamping.add_metric( [interface, identity, as_str(port.get("timestamping")).lower()], 1.0, ) if port_properties: yield port_state yield timestamping packet_counter = CounterMetricFamily( "ptp_port_messages", "PTP messages received/transmitted by linuxptp", labels=["interface", "port_identity", "direction", "message_type"], ) have_packet_counters = False for block in blocks_for(blocks, "PORT_STATS_NP"): identity = as_str(block.get("portIdentity")) interface = port_to_interface.get(identity, identity) for key, value in block.items(): if not (key.startswith("rx_") or key.startswith("tx_")): continue direction, message_type = key.split("_", 1) packet_counter.add_metric( [interface, identity, direction, message_type.lower()], as_float(value, 0.0), ) have_packet_counters = True if have_packet_counters: yield packet_counter service_counter = CounterMetricFamily( "ptp_port_service_events", "linuxptp port service timeout/mismatch counters", labels=["interface", "port_identity", "event"], ) have_service_counters = False for block in blocks_for(blocks, "PORT_SERVICE_STATS_NP"): identity = as_str(block.get("portIdentity")) interface = port_to_interface.get(identity, identity) for key, value in block.items(): if key.startswith("_") or key == "portIdentity": continue if isinstance(value, (int, float)): service_counter.add_metric( [interface, identity, key.lower()], as_float(value, 0.0), ) have_service_counters = True if have_service_counters: yield service_counter # -------- systemd service health -------- services: List[Tuple[str, Optional[str]]] = [ ("ptp4l", self.ptp4l_service), ("phc2sys", self.phc2sys_service), ] service_up = GaugeMetricFamily( "ptp_service_up", "Whether the configured local PTP-related systemd service is active", labels=["component", "service"], ) have_services = False for component, service in services: if service: service_up.add_metric([component, service], systemd_active(service)) have_services = True if have_services: yield service_up # -------- phc2sys servo data from journal -------- if self.phc2sys_service: try: rows = read_phc2sys_journal(self.phc2sys_service) offset = GaugeMetricFamily( "ptp_phc2sys_offset_nanoseconds", "Latest phc2sys servo offset from journal output", labels=["clock"], ) frequency = GaugeMetricFamily( "ptp_phc2sys_frequency_ppb", "Latest phc2sys frequency correction from journal output", labels=["clock"], ) delay = GaugeMetricFamily( "ptp_phc2sys_delay_nanoseconds", "Latest phc2sys clock-read delay from journal output", labels=["clock"], ) for clock, row in rows.items(): offset.add_metric([clock], row["offset_ns"]) frequency.add_metric([clock], row["frequency_ppb"]) if "delay_ns" in row: delay.add_metric([clock], row["delay_ns"]) if rows: yield offset yield frequency yield delay except Exception: # Do not fail the PTP scrape merely because journal access is unavailable. pass # -------- chrony data on the grandmaster -------- if self.chrony: try: tracking = parse_chrony_tracking( run_command(["chronyc", "tracking"], timeout=self.command_timeout) ) chrony_fields = [ ("ptp_gm_chrony_stratum", "stratum", "Chrony stratum"), ( "ptp_gm_chrony_system_time_offset_seconds", "system_time_offset_seconds", "CLOCK_REALTIME offset from chrony's reference; positive means local clock ahead", ), ( "ptp_gm_chrony_last_offset_seconds", "last_offset_seconds", "Chrony last measured offset", ), ( "ptp_gm_chrony_rms_offset_seconds", "rms_offset_seconds", "Chrony RMS offset", ), ( "ptp_gm_chrony_frequency_ppm", "frequency_ppm", "Chrony frequency correction in ppm", ), ( "ptp_gm_chrony_residual_frequency_ppm", "residual_frequency_ppm", "Chrony residual frequency in ppm", ), ("ptp_gm_chrony_skew_ppm", "skew_ppm", "Chrony frequency skew in ppm"), ( "ptp_gm_chrony_root_delay_seconds", "root_delay_seconds", "Chrony root delay", ), ( "ptp_gm_chrony_root_dispersion_seconds", "root_dispersion_seconds", "Chrony root dispersion", ), ] for metric_name, field, help_text in chrony_fields: if field in tracking: g = GaugeMetricFamily(metric_name, help_text) g.add_metric([], float(tracking[field])) yield g leap = tracking.get("leap_status") if leap is not None: info = InfoMetricFamily( "ptp_gm_chrony_tracking", "Chrony tracking metadata", labels=["reference_id", "leap_status"], ) info.add_metric( [ str(tracking.get("reference_id", "")), str(leap), ], {}, ) yield info sources = parse_chrony_sources( run_command(["chronyc", "sources", "-n"], timeout=self.command_timeout) ) source_metric = GaugeMetricFamily( "ptp_gm_chrony_source", "Chrony source state. State '*' is selected, '+' is combined, '?' unreachable, etc.", labels=["source", "mode", "state"], ) for source in sources: source_metric.add_metric( [source["source"], source["mode"], source["state"]], 1.0, ) if sources: yield source_metric except Exception: if scrape_success: # Core PTP data is still valid; expose chrony failure separately. pass # -------- PMC dataset health -------- dataset_metric = GaugeMetricFamily( "ptp_exporter_pmc_dataset_success", "1 if this linuxptp management dataset returned at least one response", labels=["dataset", "error"], ) for dataset in PMC_DATASETS: ok = dataset_status.get(dataset, 0.0) error = "none" if ok else dataset_errors.get(dataset, "unknown") dataset_metric.add_metric([dataset, error], ok) yield dataset_metric # -------- exporter self-health -------- success = GaugeMetricFamily( "ptp_exporter_scrape_success", "1 if the core linuxptp PMC scrape succeeded", labels=["stage"], ) success.add_metric([error_stage or "ok"], scrape_success) yield success duration = GaugeMetricFamily( "ptp_exporter_scrape_duration_seconds", "Time spent collecting PTP metrics", ) duration.add_metric([], time.monotonic() - start) yield duration def build_arg_parser() -> argparse.ArgumentParser: p = argparse.ArgumentParser(description="Prometheus exporter for linuxptp") p.add_argument( "--role", choices=["grandmaster", "boundary", "client"], required=True, help="Logical role of this Linux node", ) p.add_argument("--listen", default="0.0.0.0") p.add_argument("--port", type=int, default=9559) p.add_argument("--uds", default="/var/run/ptp4lro", help="ptp4l read-only UDS") p.add_argument("--pmc-binary", default="/usr/sbin/pmc") p.add_argument( "--pmc-runtime-dir", default="/run/ptp-exporter", help="Writable directory used for pmc's client-side Unix sockets", ) p.add_argument("--timeout", type=float, default=3.0) p.add_argument( "--chrony", action="store_true", help="Also export chrony tracking/source metrics; intended for the GNSS grandmaster", ) p.add_argument("--ptp4l-service", default=None) p.add_argument("--phc2sys-service", default=None) return p def main() -> None: args = build_arg_parser().parse_args() collector = LinuxPTPCollector( role=args.role, uds=args.uds, pmc_binary=args.pmc_binary, command_timeout=args.timeout, chrony=args.chrony, ptp4l_service=args.ptp4l_service, phc2sys_service=args.phc2sys_service, pmc_runtime_dir=args.pmc_runtime_dir, ) REGISTRY.register(collector) start_http_server(args.port, addr=args.listen) print( f"linux PTP exporter listening on {args.listen}:{args.port} " f"(role={args.role}, uds={args.uds})", flush=True, ) while True: time.sleep(3600) if __name__ == "__main__": main()