From fddb9fb31c406081a33d9f86d933d2ce880156a0 Mon Sep 17 00:00:00 2001 From: Adeeb Shihadeh Date: Sat, 20 Jun 2026 16:35:28 -0700 Subject: [PATCH] remove statsd (#38204) --- common/hardware/hw.py | 7 -- selfdrive/test/test_onroad.py | 1 - system/athena/athenad.py | 30 ----- system/hardware/hardwared.py | 21 ---- system/hardware/power_monitoring.py | 2 - system/loggerd/config.py | 4 - system/manager/process_config.py | 1 - system/statsd.py | 183 ---------------------------- 8 files changed, 249 deletions(-) delete mode 100755 system/statsd.py diff --git a/common/hardware/hw.py b/common/hardware/hw.py index 2c649db9e..1041a17c1 100644 --- a/common/hardware/hw.py +++ b/common/hardware/hw.py @@ -44,13 +44,6 @@ class Paths: else: return "/persist/" - @staticmethod - def stats_root() -> str: - if PC: - return str(Path(Paths.comma_home()) / "stats") - else: - return "/data/stats/" - @staticmethod def config_root() -> str: if PC: diff --git a/selfdrive/test/test_onroad.py b/selfdrive/test/test_onroad.py index 99d0279c4..974859a23 100644 --- a/selfdrive/test/test_onroad.py +++ b/selfdrive/test/test_onroad.py @@ -63,7 +63,6 @@ PROCS = { "system.micd": 5.0, "system.timed": 0, "selfdrive.pandad.pandad": 0, - "system.statsd": 1.0, "system.loggerd.uploader": 15.0, "system.loggerd.deleter": 1.0, "./pandad": 19.0, diff --git a/system/athena/athenad.py b/system/athena/athenad.py index 697819d0a..a322652e3 100755 --- a/system/athena/athenad.py +++ b/system/athena/athenad.py @@ -12,7 +12,6 @@ import random import select import socket import sys -import tempfile import threading import time from dataclasses import asdict, dataclass, replace @@ -186,7 +185,6 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: threading.Thread(target=upload_handler, args=(end_event,), name='upload_handler3'), threading.Thread(target=upload_handler, args=(end_event,), name='upload_handler4'), threading.Thread(target=log_handler, args=(end_event,), name='log_handler'), - threading.Thread(target=stat_handler, args=(end_event,), name='stat_handler'), ] + [ threading.Thread(target=jsonrpc_handler, args=(end_event,), name=f'worker_{x}') for x in range(HANDLER_THREADS) @@ -695,34 +693,6 @@ def log_handler(end_event: threading.Event) -> None: cloudlog.exception("athena.log_handler.exception") -def stat_handler(end_event: threading.Event) -> None: - STATS_DIR = Paths.stats_root() - last_scan = 0.0 - - while not end_event.is_set(): - curr_scan = time.monotonic() - try: - if curr_scan - last_scan > 10: - stat_filenames = list(filter(lambda name: not name.startswith(tempfile.gettempprefix()), os.listdir(STATS_DIR))) - if len(stat_filenames) > 0: - stat_path = os.path.join(STATS_DIR, stat_filenames[0]) - with open(stat_path) as f: - jsonrpc = { - "method": "storeStats", - "params": { - "stats": f.read() - }, - "jsonrpc": "2.0", - "id": stat_filenames[0] - } - send_queue_push(json.dumps(jsonrpc), SEND_PRIORITY_LOW) - os.remove(stat_path) - last_scan = curr_scan - except Exception: - cloudlog.exception("athena.stat_handler.exception") - time.sleep(0.1) - - def ws_proxy_recv(ws: WebSocket, local_sock: socket.socket, ssock: socket.socket, end_event: threading.Event, global_end_event: threading.Event) -> None: while not (end_event.is_set() or global_end_event.is_set()): try: diff --git a/system/hardware/hardwared.py b/system/hardware/hardwared.py index f8a9d3add..70f31d5f7 100755 --- a/system/hardware/hardwared.py +++ b/system/hardware/hardwared.py @@ -19,7 +19,6 @@ from openpilot.common.realtime import DT_HW from openpilot.selfdrive.selfdrived.alertmanager import set_offroad_alert from openpilot.common.hardware import HARDWARE, TICI, PC from openpilot.system.loggerd.config import get_available_percent -from openpilot.system.statsd import statlog from openpilot.common.swaglog import cloudlog from openpilot.system.hardware.power_monitoring import PowerMonitoring from openpilot.system.hardware.fan_controller import FanController @@ -361,11 +360,9 @@ def hardware_thread(end_event, hw_queue) -> None: msg.deviceState.offroadPowerUsageUwh = power_monitor.get_power_used() msg.deviceState.carBatteryCapacityUwh = max(0, power_monitor.get_car_battery_capacity()) current_power_draw = HARDWARE.get_current_power_draw() - statlog.sample("power_draw", current_power_draw) msg.deviceState.powerDrawW = current_power_draw som_power_draw = HARDWARE.get_som_power_draw() - statlog.sample("som_power_draw", som_power_draw) msg.deviceState.somPowerDrawW = som_power_draw # Check if we need to shut down @@ -383,24 +380,6 @@ def hardware_thread(end_event, hw_queue) -> None: msg.deviceState.thermalStatus = thermal_status pm.send("deviceState", msg) - # Log to statsd - statlog.gauge("free_space_percent", msg.deviceState.freeSpacePercent) - statlog.gauge("gpu_usage_percent", msg.deviceState.gpuUsagePercent) - statlog.gauge("memory_usage_percent", msg.deviceState.memoryUsagePercent) - for i, usage in enumerate(msg.deviceState.cpuUsagePercent): - statlog.gauge(f"cpu{i}_usage_percent", usage) - for i, temp in enumerate(msg.deviceState.cpuTempC): - statlog.gauge(f"cpu{i}_temperature", temp) - for i, temp in enumerate(msg.deviceState.gpuTempC): - statlog.gauge(f"gpu{i}_temperature", temp) - statlog.gauge("memory_temperature", msg.deviceState.memoryTempC) - for i, temp in enumerate(msg.deviceState.pmicTempC): - statlog.gauge(f"pmic{i}_temperature", temp) - for i, temp in enumerate(last_hw_state.modem_temps): - statlog.gauge(f"modem_temperature{i}", temp) - statlog.gauge("fan_speed_percent_desired", msg.deviceState.fanSpeedPercentDesired) - statlog.gauge("screen_brightness_percent", msg.deviceState.screenBrightnessPercent) - # report to server once every 10 minutes, or every 1s when thermally blocked rising_edge_started = should_start and not should_start_prev status_packet_interval = 1. if show_alert else 600. diff --git a/system/hardware/power_monitoring.py b/system/hardware/power_monitoring.py index ca43d9007..72a8c6848 100644 --- a/system/hardware/power_monitoring.py +++ b/system/hardware/power_monitoring.py @@ -4,7 +4,6 @@ import threading from openpilot.common.params import Params from openpilot.common.hardware import HARDWARE from openpilot.common.swaglog import cloudlog -from openpilot.system.statsd import statlog CAR_VOLTAGE_LOW_PASS_K = 0.011 # LPF gain for 45s tau (dt/tau / (dt/tau + 1)) @@ -50,7 +49,6 @@ class PowerMonitoring: # Low-pass battery voltage self.car_voltage_instant_mV = voltage self.car_voltage_mV = ((voltage * CAR_VOLTAGE_LOW_PASS_K) + (self.car_voltage_mV * (1 - CAR_VOLTAGE_LOW_PASS_K))) - statlog.gauge("car_voltage", self.car_voltage_mV / 1e3) # Cap the car battery power and save it in a param every 10-ish seconds self.car_battery_capacity_uWh = max(self.car_battery_capacity_uWh, 0) diff --git a/system/loggerd/config.py b/system/loggerd/config.py index 0a8ee9839..c2d213e90 100644 --- a/system/loggerd/config.py +++ b/system/loggerd/config.py @@ -5,10 +5,6 @@ from openpilot.common.hardware.hw import Paths CAMERA_FPS = 20 SEGMENT_LENGTH = 60 -STATS_DIR_FILE_LIMIT = 10000 -STATS_SOCKET = "ipc:///tmp/stats" -STATS_FLUSH_TIME_S = 60 - def get_available_percent(default: float) -> float: try: statvfs = os.statvfs(Paths.log_root()) diff --git a/system/manager/process_config.py b/system/manager/process_config.py index a099bc753..f610d1c78 100644 --- a/system/manager/process_config.py +++ b/system/manager/process_config.py @@ -116,7 +116,6 @@ procs = [ PythonProcess("tombstoned", "system.tombstoned", always_run, enabled=not PC), PythonProcess("updated", "system.updated.updated", only_offroad, enabled=not PC), PythonProcess("uploader", "system.loggerd.uploader", always_run), - PythonProcess("statsd", "system.statsd", always_run), PythonProcess("feedbackd", "selfdrive.ui.feedback.feedbackd", only_onroad), # debug procs diff --git a/system/statsd.py b/system/statsd.py deleted file mode 100755 index 250780e01..000000000 --- a/system/statsd.py +++ /dev/null @@ -1,183 +0,0 @@ -#!/usr/bin/env python3 -import os -import zmq -import time -import uuid -from pathlib import Path -from collections import defaultdict -from datetime import datetime, UTC -from typing import NoReturn - -from openpilot.common.params import Params -from cereal.messaging import SubMaster -from openpilot.common.hardware.hw import Paths -from openpilot.common.swaglog import cloudlog -from openpilot.common.hardware import HARDWARE -from openpilot.common.utils import atomic_write -from openpilot.system.version import get_build_metadata -from openpilot.system.loggerd.config import STATS_DIR_FILE_LIMIT, STATS_SOCKET, STATS_FLUSH_TIME_S - - -class METRIC_TYPE: - GAUGE = 'g' - SAMPLE = 'sa' - -class StatLog: - def __init__(self): - self.pid = None - self.zctx = None - self.sock = None - - def connect(self) -> None: - self.zctx = zmq.Context() - self.sock = self.zctx.socket(zmq.PUSH) - self.sock.setsockopt(zmq.LINGER, 10) - self.sock.connect(STATS_SOCKET) - self.pid = os.getpid() - - def __del__(self): - if self.sock is not None: - self.sock.close() - if self.zctx is not None: - self.zctx.term() - - def _send(self, metric: str) -> None: - if os.getpid() != self.pid: - self.connect() - - try: - self.sock.send_string(metric, zmq.NOBLOCK) - except zmq.error.Again: - # drop :/ - pass - - def gauge(self, name: str, value: float) -> None: - self._send(f"{name}:{value}|{METRIC_TYPE.GAUGE}") - - # Samples will be recorded in a buffer and at aggregation time, - # statistical properties will be logged (mean, count, percentiles, ...) - def sample(self, name: str, value: float): - self._send(f"{name}:{value}|{METRIC_TYPE.SAMPLE}") - - -def main() -> NoReturn: - dongle_id = Params().get("DongleId") - def get_influxdb_line(measurement: str, value: float | dict[str, float], timestamp: datetime, tags: dict) -> str: - res = f"{measurement}" - for k, v in tags.items(): - res += f",{k}={str(v)}" - res += " " - - if isinstance(value, float): - value = {'value': value} - - for k, v in value.items(): - res += f"{k}={v}," - - res += f"dongle_id=\"{dongle_id}\" {int(timestamp.timestamp() * 1e9)}\n" - return res - - # open statistics socket - ctx = zmq.Context.instance() - sock = ctx.socket(zmq.PULL) - sock.bind(STATS_SOCKET) - - STATS_DIR = Paths.stats_root() - - # initialize stats directory - Path(STATS_DIR).mkdir(parents=True, exist_ok=True) - - build_metadata = get_build_metadata() - - # initialize tags - tags = { - 'started': False, - 'version': build_metadata.openpilot.version, - 'branch': build_metadata.channel, - 'dirty': build_metadata.openpilot.is_dirty, - 'origin': build_metadata.openpilot.git_normalized_origin, - 'deviceType': HARDWARE.get_device_type(), - } - - # subscribe to deviceState for started state - sm = SubMaster(['deviceState']) - - idx = 0 - boot_uid = str(uuid.uuid4())[:8] - last_flush_time = time.monotonic() - gauges = {} - samples: dict[str, list[float]] = defaultdict(list) - try: - while True: - started_prev = sm['deviceState'].started - sm.update() - - # Update metrics - while True: - try: - metric = sock.recv_string(zmq.NOBLOCK) - try: - metric_type = metric.split('|')[1] - metric_name = metric.split(':')[0] - metric_value = float(metric.split('|')[0].split(':')[1]) - - if metric_type == METRIC_TYPE.GAUGE: - gauges[metric_name] = metric_value - elif metric_type == METRIC_TYPE.SAMPLE: - samples[metric_name].append(metric_value) - else: - cloudlog.event("unknown metric type", metric_type=metric_type) - except Exception: - cloudlog.event("malformed metric", metric=metric) - except zmq.error.Again: - break - - # flush when started state changes or after FLUSH_TIME_S - if (time.monotonic() > last_flush_time + STATS_FLUSH_TIME_S) or (sm['deviceState'].started != started_prev): - result = "" - current_time = datetime.now(UTC) - tags['started'] = sm['deviceState'].started - - for key, value in gauges.items(): - result += get_influxdb_line(f"gauge.{key}", value, current_time, tags) - - for key, values in samples.items(): - values.sort() - sample_count = len(values) - sample_sum = sum(values) - - stats = { - 'count': sample_count, - 'min': values[0], - 'max': values[-1], - 'mean': sample_sum / sample_count, - } - for percentile in [0.05, 0.5, 0.95]: - value = values[int(round(percentile * (sample_count - 1)))] - stats[f"p{int(percentile * 100)}"] = value - - result += get_influxdb_line(f"sample.{key}", stats, current_time, tags) - - # clear intermediate data - gauges.clear() - samples.clear() - last_flush_time = time.monotonic() - - # check that we aren't filling up the drive - if len(os.listdir(STATS_DIR)) < STATS_DIR_FILE_LIMIT: - if len(result) > 0: - stats_path = os.path.join(STATS_DIR, f"{boot_uid}_{idx}") - with atomic_write(stats_path) as f: - f.write(result) - idx += 1 - else: - cloudlog.error("stats dir full") - finally: - sock.close() - ctx.term() - - -if __name__ == "__main__": - main() -else: - statlog = StatLog()