mirror of
https://github.com/sunnypilot/sunnypilot.git
synced 2026-09-29 21:33:42 +08:00
Merge remote-tracking branch 'sunnypilot/sunnypilot/master' into hkg-angle-steering-2025
# Conflicts: # opendbc_repo
This commit is contained in:
@@ -25,7 +25,5 @@ source .venv/bin/activate
|
||||
|
||||
if [[ "$(uname)" == 'Darwin' ]]; then
|
||||
touch "$ROOT"/.env
|
||||
echo "# msgq doesn't work on mac" >> "$ROOT"/.env
|
||||
echo "export ZMQ=1" >> "$ROOT"/.env
|
||||
echo "export OBJC_DISABLE_INITIALIZE_FORK_SAFETY=YES" >> "$ROOT"/.env
|
||||
fi
|
||||
|
||||
@@ -9,7 +9,7 @@ from openpilot.selfdrive.test.process_replay.migration import migrate_all
|
||||
from openpilot.tools.lib.logreader import _LogFileReader, LogReader
|
||||
|
||||
|
||||
def flatten_dict(d: dict, sep: str = "/", prefix: str = None) -> dict:
|
||||
def flatten_dict(d: dict, sep: str = "/", prefix: str | None = None) -> dict:
|
||||
result = {}
|
||||
stack: list[tuple] = [(d, prefix)]
|
||||
|
||||
|
||||
@@ -7,7 +7,7 @@ import dearpygui.dearpygui as dpg
|
||||
import multiprocessing
|
||||
import uuid
|
||||
import signal
|
||||
import yaml # type: ignore
|
||||
import yaml
|
||||
from openpilot.common.swaglog import cloudlog
|
||||
from openpilot.common.basedir import BASEDIR
|
||||
from openpilot.tools.jotpluggler.data import DataManager
|
||||
|
||||
@@ -9,7 +9,7 @@ from abc import ABC, abstractmethod
|
||||
class ViewPanel(ABC):
|
||||
"""Abstract base class for all view panels that can be displayed in a plot container"""
|
||||
|
||||
def __init__(self, panel_id: str = None):
|
||||
def __init__(self, panel_id: str | None = None):
|
||||
self.panel_id = panel_id or str(uuid.uuid4())
|
||||
self.title = "Untitled Panel"
|
||||
|
||||
|
||||
+15
-15
@@ -1,6 +1,7 @@
|
||||
import os
|
||||
import subprocess
|
||||
import json
|
||||
import logging
|
||||
from collections.abc import Iterator
|
||||
from collections import OrderedDict
|
||||
|
||||
@@ -9,12 +10,12 @@ from openpilot.tools.lib.filereader import FileReader, resolve_name
|
||||
from openpilot.tools.lib.exceptions import DataUnreadableError
|
||||
from openpilot.tools.lib.vidindex import hevc_index
|
||||
|
||||
logger = logging.getLogger("tools")
|
||||
|
||||
HEVC_SLICE_B = 0
|
||||
HEVC_SLICE_P = 1
|
||||
HEVC_SLICE_I = 2
|
||||
|
||||
|
||||
class LRUCache:
|
||||
def __init__(self, capacity: int):
|
||||
self._cache: OrderedDict = OrderedDict()
|
||||
@@ -32,7 +33,6 @@ class LRUCache:
|
||||
def __contains__(self, key):
|
||||
return key in self._cache
|
||||
|
||||
|
||||
def assert_hvec(fn: str) -> None:
|
||||
with FileReader(fn) as f:
|
||||
header = f.read(4)
|
||||
@@ -42,10 +42,11 @@ def assert_hvec(fn: str) -> None:
|
||||
if 'hevc' not in fn:
|
||||
raise NotImplementedError(fn)
|
||||
|
||||
def decompress_video_data(rawdat, w, h, pix_fmt="rgb24", vid_fmt='hevc') -> np.ndarray:
|
||||
def decompress_video_data(rawdat, w, h, pix_fmt="rgb24", vid_fmt='hevc', hwaccel="auto", loglevel="info") -> np.ndarray:
|
||||
threads = os.getenv("FFMPEG_THREADS", "0")
|
||||
args = ["ffmpeg", "-v", "quiet",
|
||||
args = ["ffmpeg", "-v", loglevel,
|
||||
"-threads", threads,
|
||||
"-hwaccel", hwaccel,
|
||||
"-c:v", "hevc",
|
||||
"-vsync", "0",
|
||||
"-f", vid_fmt,
|
||||
@@ -98,15 +99,15 @@ def get_video_index(fn):
|
||||
'probe': probe
|
||||
}
|
||||
|
||||
|
||||
class FfmpegDecoder:
|
||||
def __init__(self, fn: str, index_data: dict|None = None,
|
||||
pix_fmt: str = "rgb24"):
|
||||
pix_fmt: str = "rgb24", hwaccel="auto", loglevel="quiet"):
|
||||
self.fn = fn
|
||||
self.index, self.prefix, self.w, self.h = get_index_data(fn, index_data)
|
||||
self.frame_count = len(self.index) - 1 # sentinel row at the end
|
||||
self.iframes = np.where(self.index[:, 0] == HEVC_SLICE_I)[0]
|
||||
self.pix_fmt = pix_fmt
|
||||
self.loglevel, self.hwaccel = loglevel, hwaccel
|
||||
|
||||
def _gop_bounds(self, frame_idx: int):
|
||||
f_b = frame_idx
|
||||
@@ -118,7 +119,7 @@ class FfmpegDecoder:
|
||||
return f_b, f_e, self.index[f_b, 1], self.index[f_e, 1]
|
||||
|
||||
def _decode_gop(self, raw: bytes) -> Iterator[np.ndarray]:
|
||||
yield from decompress_video_data(raw, self.w, self.h, self.pix_fmt)
|
||||
yield from decompress_video_data(raw, self.w, self.h, pix_fmt=self.pix_fmt, hwaccel=self.hwaccel, loglevel=self.loglevel)
|
||||
|
||||
def get_gop_start(self, frame_idx: int):
|
||||
return self.iframes[np.searchsorted(self.iframes, frame_idx, side="right") - 1]
|
||||
@@ -133,7 +134,7 @@ class FfmpegDecoder:
|
||||
f.seek(off_b)
|
||||
raw = self.prefix + f.read(off_e - off_b)
|
||||
# number of frames to discard inside this GOP before the wanted one
|
||||
for i, frm in enumerate(decompress_video_data(raw, self.w, self.h, self.pix_fmt)):
|
||||
for i, frm in enumerate(decompress_video_data(raw, self.w, self.h, self.pix_fmt, hwaccel=self.hwaccel, loglevel=self.loglevel)):
|
||||
fidx = f_b + i
|
||||
if fidx >= end_fidx:
|
||||
return
|
||||
@@ -141,17 +142,16 @@ class FfmpegDecoder:
|
||||
yield fidx, frm
|
||||
fidx += 1
|
||||
|
||||
def FrameIterator(fn: str, index_data: dict|None=None,
|
||||
pix_fmt: str = "rgb24",
|
||||
start_fidx:int=0, end_fidx=None, frame_skip:int=1) -> Iterator[np.ndarray]:
|
||||
dec = FfmpegDecoder(fn, pix_fmt=pix_fmt, index_data=index_data)
|
||||
def FrameIterator(fn: str, index_data: dict|None=None, pix_fmt: str = "rgb24",
|
||||
start_fidx:int=0, end_fidx=None, frame_skip:int=1, hwaccel="auto", loglevel="quiet") -> Iterator[np.ndarray]:
|
||||
dec = FfmpegDecoder(fn, pix_fmt=pix_fmt, index_data=index_data, hwaccel=hwaccel, loglevel=loglevel)
|
||||
for _, frame in dec.get_iterator(start_fidx=start_fidx, end_fidx=end_fidx, frame_skip=frame_skip):
|
||||
yield frame
|
||||
|
||||
class FrameReader:
|
||||
def __init__(self, fn: str, index_data: dict|None = None,
|
||||
cache_size: int = 30, pix_fmt: str = "rgb24"):
|
||||
self.decoder = FfmpegDecoder(fn, index_data, pix_fmt)
|
||||
def __init__(self, fn: str, index_data: dict|None = None, cache_size: int = 30,
|
||||
pix_fmt: str = "rgb24", hwaccel="auto", loglevel="quiet"):
|
||||
self.decoder = FfmpegDecoder(fn, index_data=index_data, pix_fmt=pix_fmt, hwaccel=hwaccel, loglevel=loglevel)
|
||||
self.iframes = self.decoder.iframes
|
||||
self._cache: LRUCache = LRUCache(cache_size)
|
||||
self.w, self.h, self.frame_count, = self.decoder.w, self.decoder.h, self.decoder.frame_count
|
||||
|
||||
@@ -13,7 +13,6 @@ import warnings
|
||||
import zstandard as zstd
|
||||
|
||||
from collections.abc import Iterable, Iterator
|
||||
from typing import cast
|
||||
from urllib.parse import parse_qs, urlparse
|
||||
|
||||
from cereal import log as capnp_log
|
||||
@@ -180,7 +179,7 @@ def auto_source(identifier: str, sources: list[Source], default_mode: ReadMode)
|
||||
|
||||
# We've found all files, return them
|
||||
if len(needed_seg_idxs) == 0:
|
||||
return cast(list[str], list(valid_files.values()))
|
||||
return list(valid_files.values())
|
||||
else:
|
||||
raise FileNotFoundError(f"Did not find {fn} for seg idxs {needed_seg_idxs} of {sr.route_name}")
|
||||
|
||||
@@ -245,7 +244,7 @@ class LogReader:
|
||||
return identifiers
|
||||
|
||||
def __init__(self, identifier: str | list[str], default_mode: ReadMode = ReadMode.RLOG,
|
||||
sources: list[Source] = None, sort_by_time=False, only_union_types=False):
|
||||
sources: list[Source] | None = None, sort_by_time=False, only_union_types=False):
|
||||
if sources is None:
|
||||
sources = [internal_source, comma_api_source, openpilotci_source, comma_car_segments_source]
|
||||
|
||||
|
||||
+14
-13
@@ -23,7 +23,6 @@ class FileName:
|
||||
|
||||
class Route:
|
||||
def __init__(self, name, data_dir=None):
|
||||
self._metadata = None
|
||||
self._name = RouteName(name)
|
||||
self.files = None
|
||||
if data_dir is not None:
|
||||
@@ -32,13 +31,6 @@ class Route:
|
||||
self._segments = self._get_segments_remote()
|
||||
self.max_seg_number = self._segments[-1].name.segment_num
|
||||
|
||||
@property
|
||||
def metadata(self):
|
||||
if not self._metadata:
|
||||
api = CommaApi(get_token())
|
||||
self._metadata = api.get('v1/route/' + self.name.canonical_name)
|
||||
return self._metadata
|
||||
|
||||
@property
|
||||
def name(self):
|
||||
return self._name
|
||||
@@ -90,7 +82,6 @@ class Route:
|
||||
url if fn in FileName.DCAMERA else segments[segment_name].dcamera_path,
|
||||
url if fn in FileName.ECAMERA else segments[segment_name].ecamera_path,
|
||||
url if fn in FileName.QCAMERA else segments[segment_name].qcamera_path,
|
||||
self.metadata['url'],
|
||||
)
|
||||
else:
|
||||
segments[segment_name] = Segment(
|
||||
@@ -101,7 +92,6 @@ class Route:
|
||||
url if fn in FileName.DCAMERA else None,
|
||||
url if fn in FileName.ECAMERA else None,
|
||||
url if fn in FileName.QCAMERA else None,
|
||||
self.metadata['url'],
|
||||
)
|
||||
|
||||
return sorted(segments.values(), key=lambda seg: seg.name.segment_num)
|
||||
@@ -167,7 +157,7 @@ class Route:
|
||||
except StopIteration:
|
||||
qcamera_path = None
|
||||
|
||||
segments.append(Segment(segment, log_path, qlog_path, camera_path, dcamera_path, ecamera_path, qcamera_path, self.metadata['url']))
|
||||
segments.append(Segment(segment, log_path, qlog_path, camera_path, dcamera_path, ecamera_path, qcamera_path))
|
||||
|
||||
if len(segments) == 0:
|
||||
raise ValueError(f'Could not find segments for route {self.name.canonical_name} in data directory {data_dir}')
|
||||
@@ -175,10 +165,9 @@ class Route:
|
||||
|
||||
|
||||
class Segment:
|
||||
def __init__(self, name, log_path, qlog_path, camera_path, dcamera_path, ecamera_path, qcamera_path, url):
|
||||
def __init__(self, name, log_path, qlog_path, camera_path, dcamera_path, ecamera_path, qcamera_path):
|
||||
self._events = None
|
||||
self._name = SegmentName(name)
|
||||
self.url = f'{url}/{self._name.segment_num}'
|
||||
self.log_path = log_path
|
||||
self.qlog_path = qlog_path
|
||||
self.camera_path = camera_path
|
||||
@@ -190,6 +179,18 @@ class Segment:
|
||||
def name(self):
|
||||
return self._name
|
||||
|
||||
@staticmethod
|
||||
@cache
|
||||
def _get_route_metadata(route_name: str):
|
||||
api = CommaApi(get_token())
|
||||
return api.get(f'v1/route/{route_name}')
|
||||
|
||||
@property
|
||||
def url(self):
|
||||
route_name = self._name.route_name.canonical_name
|
||||
metadata = self._get_route_metadata(route_name)
|
||||
return f'{metadata["url"]}/{self._name.segment_num}'
|
||||
|
||||
@property
|
||||
def events(self):
|
||||
if not self._events:
|
||||
|
||||
@@ -2,11 +2,13 @@ import http.server
|
||||
import os
|
||||
import shutil
|
||||
import socket
|
||||
import tempfile
|
||||
import pytest
|
||||
|
||||
from openpilot.selfdrive.test.helpers import http_server_context
|
||||
from openpilot.system.hardware.hw import Paths
|
||||
from openpilot.tools.lib.url_file import URLFile
|
||||
from openpilot.tools.lib.url_file import URLFile, prune_cache
|
||||
import openpilot.tools.lib.url_file as url_file_module
|
||||
|
||||
|
||||
class CachingTestRequestHandler(http.server.BaseHTTPRequestHandler):
|
||||
@@ -128,3 +130,35 @@ class TestFileDownload:
|
||||
CachingTestRequestHandler.FILE_EXISTS = True
|
||||
length = URLFile(file_url).get_length()
|
||||
assert length == 4
|
||||
|
||||
|
||||
class TestCache:
|
||||
def test_prune_cache(self, monkeypatch):
|
||||
with tempfile.TemporaryDirectory() as tmpdir:
|
||||
monkeypatch.setattr(Paths, 'download_cache_root', staticmethod(lambda: tmpdir + "/"))
|
||||
|
||||
# setup test files and manifest
|
||||
manifest_lines = []
|
||||
for i in range(3):
|
||||
fname = f"hash_{i}"
|
||||
with open(tmpdir + "/" + fname, "wb") as f:
|
||||
f.truncate(1000)
|
||||
manifest_lines.append(f"{fname} {1000 + i}")
|
||||
with open(tmpdir + "/manifest.txt", "w") as f:
|
||||
f.write('\n'.join(manifest_lines))
|
||||
|
||||
# under limit, shouldn't prune
|
||||
assert len(os.listdir(tmpdir)) == 4
|
||||
prune_cache()
|
||||
assert len(os.listdir(tmpdir)) == 4
|
||||
|
||||
# set a tiny cache limit to force eviction (1.5 chunks worth)
|
||||
monkeypatch.setattr(url_file_module, 'CACHE_SIZE', url_file_module.CHUNK_SIZE + url_file_module.CHUNK_SIZE // 2)
|
||||
|
||||
# prune_cache should evict oldest files to get under limit
|
||||
prune_cache()
|
||||
remaining = os.listdir(tmpdir)
|
||||
# should have evicted at least one file + manifest
|
||||
assert len(remaining) < 4
|
||||
# newest file should remain
|
||||
assert manifest_lines[2].split()[0] in remaining
|
||||
|
||||
+36
-7
@@ -1,8 +1,9 @@
|
||||
import re
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import socket
|
||||
from hashlib import sha256
|
||||
import time
|
||||
from hashlib import md5
|
||||
from urllib3 import PoolManager, Retry
|
||||
from urllib3.response import BaseHTTPResponse
|
||||
from urllib3.util import Timeout
|
||||
@@ -14,14 +15,41 @@ from urllib3.exceptions import MaxRetryError
|
||||
# Cache chunk size
|
||||
K = 1000
|
||||
CHUNK_SIZE = 1000 * K
|
||||
CACHE_SIZE = 10 * 1024 * 1024 * 1024 # total cache size in GB
|
||||
|
||||
logging.getLogger("urllib3").setLevel(logging.WARNING)
|
||||
|
||||
|
||||
def hash_256(link: str) -> str:
|
||||
return sha256((link.split("?")[0]).encode('utf-8')).hexdigest()
|
||||
def hash_url(link: str) -> str:
|
||||
return md5((link.split("?")[0]).encode('utf-8')).hexdigest()
|
||||
|
||||
|
||||
def prune_cache(new_entry: str | None = None) -> None:
|
||||
"""Evicts oldest cache files (LRU) until cache is under the size limit."""
|
||||
# we use a manifest to avoid tons of os.stat syscalls (slow)
|
||||
manifest = {}
|
||||
manifest_path = Paths.download_cache_root() + "manifest.txt"
|
||||
if os.path.exists(manifest_path):
|
||||
with open(manifest_path) as f:
|
||||
manifest = {parts[0]: int(parts[1]) for line in f if (parts := line.strip().split()) and len(parts) == 2}
|
||||
|
||||
if new_entry:
|
||||
manifest[new_entry] = int(time.time()) # noqa: TID251
|
||||
|
||||
# evict the least recently used files until under limit
|
||||
sorted_items = sorted(manifest.items(), key=lambda x: x[1])
|
||||
while len(manifest) * CHUNK_SIZE > CACHE_SIZE and sorted_items:
|
||||
key, _ = sorted_items.pop(0)
|
||||
try:
|
||||
os.remove(Paths.download_cache_root() + key)
|
||||
except OSError:
|
||||
pass
|
||||
manifest.pop(key, None)
|
||||
|
||||
# write out manifest
|
||||
with atomic_write(manifest_path, mode="w", overwrite=True) as f:
|
||||
f.write('\n'.join(f"{k} {v}" for k, v in manifest.items()))
|
||||
|
||||
class URLFileException(Exception):
|
||||
pass
|
||||
|
||||
@@ -77,7 +105,7 @@ class URLFile:
|
||||
if self._length is not None:
|
||||
return self._length
|
||||
|
||||
file_length_path = os.path.join(Paths.download_cache_root(), hash_256(self._url) + "_length")
|
||||
file_length_path = os.path.join(Paths.download_cache_root(), hash_url(self._url) + "_length")
|
||||
if not self._force_download and os.path.exists(file_length_path):
|
||||
with open(file_length_path) as file_length:
|
||||
content = file_length.read()
|
||||
@@ -103,7 +131,7 @@ class URLFile:
|
||||
while True:
|
||||
self._pos = position
|
||||
chunk_number = self._pos / CHUNK_SIZE
|
||||
file_name = hash_256(self._url) + "_" + str(chunk_number)
|
||||
file_name = hash_url(self._url) + "_" + str(chunk_number)
|
||||
full_path = os.path.join(Paths.download_cache_root(), str(file_name))
|
||||
data = None
|
||||
# If we don't have a file, download it
|
||||
@@ -111,6 +139,7 @@ class URLFile:
|
||||
data = self.read_aux(ll=CHUNK_SIZE)
|
||||
with atomic_write(full_path, mode="wb", overwrite=True) as new_cached_file:
|
||||
new_cached_file.write(data)
|
||||
prune_cache(file_name)
|
||||
else:
|
||||
with open(full_path, "rb") as cached_file:
|
||||
data = cached_file.read()
|
||||
@@ -164,7 +193,7 @@ class URLFile:
|
||||
return parts
|
||||
|
||||
def seek(self, pos: int) -> None:
|
||||
self._pos = pos
|
||||
self._pos = int(pos)
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
|
||||
@@ -140,7 +140,7 @@ def get_ue(dat: bytes, start_idx: int, skip_bits: int) -> tuple[int, int]:
|
||||
j -= 1
|
||||
|
||||
if prefix_val == 1 and prefix_len - 1 == suffix_len:
|
||||
val = 2**(prefix_len-1) - 1 + suffix_val
|
||||
val = int(2**(prefix_len-1) - 1 + suffix_val)
|
||||
size = prefix_len + suffix_len
|
||||
return val, size
|
||||
i += 1
|
||||
|
||||
+4
-13
@@ -1,24 +1,14 @@
|
||||
#include "tools/replay/camera.h"
|
||||
|
||||
#include <cassert>
|
||||
#include <algorithm>
|
||||
|
||||
#include <capnp/dynamic.h>
|
||||
|
||||
#include "third_party/linux/include/msm_media_info.h"
|
||||
#include "system/camerad/cameras/nv12_info.h"
|
||||
#include "tools/replay/util.h"
|
||||
|
||||
const int BUFFER_COUNT = 40;
|
||||
|
||||
std::tuple<size_t, size_t, size_t> get_nv12_info(int width, int height) {
|
||||
int nv12_width = VENUS_Y_STRIDE(COLOR_FMT_NV12, width);
|
||||
int nv12_height = VENUS_Y_SCANLINES(COLOR_FMT_NV12, height);
|
||||
assert(nv12_width == VENUS_UV_STRIDE(COLOR_FMT_NV12, width));
|
||||
assert(nv12_height / 2 == VENUS_UV_SCANLINES(COLOR_FMT_NV12, height));
|
||||
size_t nv12_buffer_size = 2346 * nv12_width; // comes from v4l2_format.fmt.pix_mp.plane_fmt[0].sizeimage
|
||||
return {nv12_width, nv12_height, nv12_buffer_size};
|
||||
}
|
||||
|
||||
CameraServer::CameraServer(std::pair<int, int> camera_size[MAX_CAMERAS]) {
|
||||
for (int i = 0; i < MAX_CAMERAS; ++i) {
|
||||
std::tie(cameras_[i].width, cameras_[i].height) = camera_size[i];
|
||||
@@ -50,9 +40,10 @@ void CameraServer::startVipcServer() {
|
||||
|
||||
if (cam.width > 0 && cam.height > 0) {
|
||||
rInfo("camera[%d] frame size %dx%d", cam.type, cam.width, cam.height);
|
||||
auto [nv12_width, nv12_height, nv12_buffer_size] = get_nv12_info(cam.width, cam.height);
|
||||
auto [stride, y_height, uv_height_, buffer_size] = get_nv12_info(cam.width, cam.height);
|
||||
(void)uv_height_; // unused in replay
|
||||
vipc_server_->create_buffers_with_sizes(cam.stream_type, BUFFER_COUNT, cam.width, cam.height,
|
||||
nv12_buffer_size, nv12_width, nv12_width * nv12_height);
|
||||
buffer_size, stride, stride * y_height);
|
||||
if (!cam.thread.joinable()) {
|
||||
cam.thread = std::thread(&CameraServer::cameraThread, this, std::ref(cam));
|
||||
}
|
||||
|
||||
@@ -10,8 +10,6 @@
|
||||
#include "tools/replay/framereader.h"
|
||||
#include "tools/replay/logreader.h"
|
||||
|
||||
std::tuple<size_t, size_t, size_t> get_nv12_info(int width, int height);
|
||||
|
||||
class CameraServer {
|
||||
public:
|
||||
CameraServer(std::pair<int, int> camera_size[MAX_CAMERAS] = nullptr);
|
||||
|
||||
@@ -17,7 +17,7 @@ std::string cacheFilePath(const std::string &url) {
|
||||
}
|
||||
|
||||
std::string FileReader::read(const std::string &file, std::atomic<bool> *abort) {
|
||||
const bool is_remote = file.find("https://") == 0;
|
||||
const bool is_remote = (file.find("https://") == 0) || (file.find("http://") == 0);
|
||||
const std::string local_file = is_remote ? cacheFilePath(file) : file;
|
||||
std::string result;
|
||||
|
||||
|
||||
@@ -72,7 +72,7 @@ FrameReader::~FrameReader() {
|
||||
}
|
||||
|
||||
bool FrameReader::load(CameraType type, const std::string &url, bool no_hw_decoder, std::atomic<bool> *abort, bool local_cache, int chunk_size, int retries) {
|
||||
auto local_file_path = url.find("https://") == 0 ? cacheFilePath(url) : url;
|
||||
auto local_file_path = (url.find("https://") == 0 || url.find("http://") == 0) ? cacheFilePath(url) : url;
|
||||
if (!util::file_exists(local_file_path)) {
|
||||
FileReader f(local_cache, chunk_size, retries);
|
||||
if (f.read(url, abort).empty()) {
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
import itertools
|
||||
from typing import Any
|
||||
|
||||
import matplotlib.pyplot as plt
|
||||
import numpy as np
|
||||
import pygame
|
||||
import pyray as rl
|
||||
|
||||
from matplotlib.backends.backend_agg import FigureCanvasAgg
|
||||
|
||||
@@ -18,21 +17,25 @@ YELLOW = (255, 255, 0)
|
||||
BLACK = (0, 0, 0)
|
||||
WHITE = (255, 255, 255)
|
||||
|
||||
|
||||
class UIParams:
|
||||
lidar_x, lidar_y, lidar_zoom = 384, 960, 6
|
||||
lidar_car_x, lidar_car_y = lidar_x / 2., lidar_y / 1.1
|
||||
lidar_car_x, lidar_car_y = lidar_x / 2.0, lidar_y / 1.1
|
||||
car_hwidth = 1.7272 / 2 * lidar_zoom
|
||||
car_front = 2.6924 * lidar_zoom
|
||||
car_back = 1.8796 * lidar_zoom
|
||||
car_color = 110
|
||||
|
||||
|
||||
UP = UIParams
|
||||
|
||||
METER_WIDTH = 20
|
||||
|
||||
|
||||
class Calibration:
|
||||
def __init__(self, num_px, rpy, intrinsic, calib_scale):
|
||||
self.intrinsic = intrinsic
|
||||
self.extrinsics_matrix = get_view_frame_from_calib_frame(rpy[0], rpy[1], rpy[2], 0.0)[:,:3]
|
||||
self.extrinsics_matrix = get_view_frame_from_calib_frame(rpy[0], rpy[1], rpy[2], 0.0)[:, :3]
|
||||
self.zoom = calib_scale
|
||||
|
||||
def car_space_to_ff(self, x, y, z):
|
||||
@@ -47,19 +50,18 @@ class Calibration:
|
||||
return pts / self.zoom
|
||||
|
||||
|
||||
_COLOR_CACHE : dict[tuple[int, int, int], Any] = {}
|
||||
_COLOR_CACHE: dict[tuple[int, int, int], int] = {
|
||||
(255, 0, 0): 1, # RED
|
||||
(0, 255, 0): 2, # GREEN
|
||||
(0, 0, 255): 3, # BLUE
|
||||
(255, 255, 0): 4, # YELLOW
|
||||
(0, 0, 0): 0, # BLACK
|
||||
(255, 255, 255): 255, # WHITE
|
||||
}
|
||||
|
||||
|
||||
def find_color(lidar_surface, color):
|
||||
if color in _COLOR_CACHE:
|
||||
return _COLOR_CACHE[color]
|
||||
tcolor = 0
|
||||
ret = 255
|
||||
for x in lidar_surface.get_palette():
|
||||
if x[0:3] == color:
|
||||
ret = tcolor
|
||||
break
|
||||
tcolor += 1
|
||||
_COLOR_CACHE[color] = ret
|
||||
return ret
|
||||
return _COLOR_CACHE.get(color, 255)
|
||||
|
||||
|
||||
def to_topdown_pt(y, x):
|
||||
@@ -91,13 +93,7 @@ def draw_path(path, color, img, calibration, top_down, lid_color=None, z_off=0):
|
||||
|
||||
|
||||
def init_plots(arr, name_to_arr_idx, plot_xlims, plot_ylims, plot_names, plot_colors, plot_styles):
|
||||
color_palette = { "r": (1, 0, 0),
|
||||
"g": (0, 1, 0),
|
||||
"b": (0, 0, 1),
|
||||
"k": (0, 0, 0),
|
||||
"y": (1, 1, 0),
|
||||
"p": (0, 1, 1),
|
||||
"m": (1, 0, 1)}
|
||||
color_palette = {"r": (1, 0, 0), "g": (0, 1, 0), "b": (0, 0, 1), "k": (0, 0, 0), "y": (1, 1, 0), "p": (0, 1, 1), "m": (1, 0, 1)}
|
||||
|
||||
dpi = 90
|
||||
fig = plt.figure(figsize=(575 / dpi, 600 / dpi), dpi=dpi)
|
||||
@@ -107,7 +103,7 @@ def init_plots(arr, name_to_arr_idx, plot_xlims, plot_ylims, plot_names, plot_co
|
||||
|
||||
axs = []
|
||||
for pn in range(len(plot_ylims)):
|
||||
ax = fig.add_subplot(len(plot_ylims), 1, len(axs)+1)
|
||||
ax = fig.add_subplot(len(plot_ylims), 1, len(axs) + 1)
|
||||
ax.set_xlim(plot_xlims[pn][0], plot_xlims[pn][1])
|
||||
ax.set_ylim(plot_ylims[pn][0], plot_ylims[pn][1])
|
||||
ax.patch.set_facecolor((0.4, 0.4, 0.4))
|
||||
@@ -116,15 +112,11 @@ def init_plots(arr, name_to_arr_idx, plot_xlims, plot_ylims, plot_names, plot_co
|
||||
plots, idxs, plot_select = [], [], []
|
||||
for i, pl_list in enumerate(plot_names):
|
||||
for j, item in enumerate(pl_list):
|
||||
plot, = axs[i].plot(arr[:, name_to_arr_idx[item]],
|
||||
label=item,
|
||||
color=color_palette[plot_colors[i][j]],
|
||||
linestyle=plot_styles[i][j])
|
||||
(plot,) = axs[i].plot(arr[:, name_to_arr_idx[item]], label=item, color=color_palette[plot_colors[i][j]], linestyle=plot_styles[i][j])
|
||||
plots.append(plot)
|
||||
idxs.append(name_to_arr_idx[item])
|
||||
plot_select.append(i)
|
||||
axs[i].set_title(", ".join(f"{nm} ({cl})"
|
||||
for (nm, cl) in zip(pl_list, plot_colors[i], strict=False)), fontsize=10)
|
||||
axs[i].set_title(", ".join(f"{nm} ({cl})" for (nm, cl) in zip(pl_list, plot_colors[i], strict=False)), fontsize=10)
|
||||
axs[i].tick_params(axis="x", colors="white")
|
||||
axs[i].tick_params(axis="y", colors="white")
|
||||
axs[i].title.set_color("white")
|
||||
@@ -134,6 +126,12 @@ def init_plots(arr, name_to_arr_idx, plot_xlims, plot_ylims, plot_names, plot_co
|
||||
|
||||
canvas.draw()
|
||||
|
||||
# Pre-create texture for plots (reuse each frame to avoid log spam)
|
||||
w, h = canvas.get_width_height()
|
||||
plot_image = rl.gen_image_color(w, h, rl.BLACK)
|
||||
plot_texture = rl.load_texture_from_image(plot_image)
|
||||
rl.unload_image(plot_image)
|
||||
|
||||
def draw_plots(arr):
|
||||
for ax in axs:
|
||||
ax.draw_artist(ax.patch)
|
||||
@@ -141,17 +139,13 @@ def init_plots(arr, name_to_arr_idx, plot_xlims, plot_ylims, plot_names, plot_co
|
||||
plots[i].set_ydata(arr[:, idxs[i]])
|
||||
axs[plot_select[i]].draw_artist(plots[i])
|
||||
|
||||
raw_data = canvas.buffer_rgba()
|
||||
plot_surface = pygame.image.frombuffer(raw_data, canvas.get_width_height(), "RGBA").convert()
|
||||
return plot_surface
|
||||
raw_data = np.ascontiguousarray(canvas.buffer_rgba(), dtype=np.uint8)
|
||||
rl.update_texture(plot_texture, rl.ffi.cast("void *", raw_data.ctypes.data))
|
||||
return plot_texture
|
||||
|
||||
return draw_plots
|
||||
|
||||
|
||||
def pygame_modules_have_loaded():
|
||||
return pygame.display.get_init() and pygame.font.get_init()
|
||||
|
||||
|
||||
def plot_model(m, img, calibration, top_down):
|
||||
if calibration is None or top_down is None:
|
||||
return
|
||||
@@ -166,7 +160,7 @@ def plot_model(m, img, calibration, top_down):
|
||||
|
||||
_, py_top = to_topdown_pt(x + x_std, y)
|
||||
px, py_bottom = to_topdown_pt(x - x_std, y)
|
||||
top_down[1][int(round(px - 4)):int(round(px + 4)), py_top:py_bottom] = find_color(top_down[0], YELLOW)
|
||||
top_down[1][int(round(px - 4)) : int(round(px + 4)), py_top:py_bottom] = find_color(top_down[0], YELLOW)
|
||||
|
||||
for path, prob, _ in zip(m.laneLines, m.laneLineProbs, m.laneLineStds, strict=True):
|
||||
color = (0, int(255 * prob), 0)
|
||||
@@ -202,22 +196,15 @@ def maybe_update_radar_points(lt, lid_overlay):
|
||||
# negative here since radar is left positive
|
||||
px, py = to_topdown_pt(pt[0], -pt[1])
|
||||
if px != -1:
|
||||
lid_overlay[px - 4:px + 4, py - 4:py + 4] = 0
|
||||
lid_overlay[px - 2:px + 2, py - 2:py + 2] = 255
|
||||
lid_overlay[px - 4 : px + 4, py - 4 : py + 4] = 0
|
||||
lid_overlay[px - 2 : px + 2, py - 2 : py + 2] = 255
|
||||
|
||||
|
||||
def get_blank_lid_overlay(UP):
|
||||
lid_overlay = np.zeros((UP.lidar_x, UP.lidar_y), 'uint8')
|
||||
# Draw the car.
|
||||
lid_overlay[int(round(UP.lidar_car_x - UP.car_hwidth)):int(
|
||||
round(UP.lidar_car_x + UP.car_hwidth)), int(round(UP.lidar_car_y -
|
||||
UP.car_front))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x - UP.car_hwidth)):int(
|
||||
round(UP.lidar_car_x + UP.car_hwidth)), int(round(UP.lidar_car_y +
|
||||
UP.car_back))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x - UP.car_hwidth)), int(
|
||||
round(UP.lidar_car_y - UP.car_front)):int(round(
|
||||
UP.lidar_car_y + UP.car_back))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x + UP.car_hwidth)), int(
|
||||
round(UP.lidar_car_y - UP.car_front)):int(round(
|
||||
UP.lidar_car_y + UP.car_back))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x - UP.car_hwidth)) : int(round(UP.lidar_car_x + UP.car_hwidth)), int(round(UP.lidar_car_y - UP.car_front))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x - UP.car_hwidth)) : int(round(UP.lidar_car_x + UP.car_hwidth)), int(round(UP.lidar_car_y + UP.car_back))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x - UP.car_hwidth)), int(round(UP.lidar_car_y - UP.car_front)) : int(round(UP.lidar_car_y + UP.car_back))] = UP.car_color
|
||||
lid_overlay[int(round(UP.lidar_car_x + UP.car_hwidth)), int(round(UP.lidar_car_y - UP.car_front)) : int(round(UP.lidar_car_y + UP.car_back))] = UP.car_color
|
||||
return lid_overlay
|
||||
|
||||
@@ -1,11 +1,13 @@
|
||||
#include <getopt.h>
|
||||
|
||||
#include <iomanip>
|
||||
#include <iostream>
|
||||
#include <map>
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
#include "common/prefix.h"
|
||||
#include "common/timing.h"
|
||||
#include "tools/replay/consoleui.h"
|
||||
#include "tools/replay/replay.h"
|
||||
#include "tools/replay/util.h"
|
||||
@@ -31,6 +33,7 @@ Options:
|
||||
--no-hw-decoder Disable HW video decoding
|
||||
--no-vipc Do not output video
|
||||
--all Output all messages including bookmarkButton, uiDebug, userBookmark
|
||||
--benchmark Run in benchmark mode (process all events then exit with stats)
|
||||
-h, --help Show this help message
|
||||
)";
|
||||
|
||||
@@ -66,6 +69,7 @@ bool parseArgs(int argc, char *argv[], ReplayConfig &config) {
|
||||
{"no-hw-decoder", no_argument, nullptr, 0},
|
||||
{"no-vipc", no_argument, nullptr, 0},
|
||||
{"all", no_argument, nullptr, 0},
|
||||
{"benchmark", no_argument, nullptr, 0},
|
||||
{"help", no_argument, nullptr, 'h'},
|
||||
{nullptr, 0, nullptr, 0}, // Terminating entry
|
||||
};
|
||||
@@ -79,6 +83,7 @@ bool parseArgs(int argc, char *argv[], ReplayConfig &config) {
|
||||
{"no-hw-decoder", REPLAY_FLAG_NO_HW_DECODER},
|
||||
{"no-vipc", REPLAY_FLAG_NO_VIPC},
|
||||
{"all", REPLAY_FLAG_ALL_SERVICES},
|
||||
{"benchmark", REPLAY_FLAG_BENCHMARK},
|
||||
};
|
||||
|
||||
if (argc == 1) {
|
||||
@@ -149,6 +154,28 @@ int main(int argc, char *argv[]) {
|
||||
return 1;
|
||||
}
|
||||
|
||||
if (config.flags & REPLAY_FLAG_BENCHMARK) {
|
||||
replay.start(config.start_seconds);
|
||||
replay.waitForFinished();
|
||||
|
||||
const auto &stats = replay.getBenchmarkStats();
|
||||
uint64_t process_start = stats.process_start_ts;
|
||||
|
||||
std::cout << "\n===== REPLAY BENCHMARK RESULTS =====\n";
|
||||
std::cout << "Route: " << replay.route().name() << "\n\n";
|
||||
|
||||
std::cout << "TIMELINE:\n";
|
||||
std::cout << " t=0 ms process start\n";
|
||||
for (const auto &[ts, event] : stats.timeline) {
|
||||
double ms = (ts - process_start) / 1e6;
|
||||
std::cout << " t=" << std::fixed << std::setprecision(0) << ms << " ms"
|
||||
<< std::string(std::max(1, 8 - static_cast<int>(std::to_string(static_cast<int>(ms)).length())), ' ')
|
||||
<< event << "\n";
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
ConsoleUI console_ui(&replay);
|
||||
replay.start(config.start_seconds);
|
||||
return console_ui.exec();
|
||||
|
||||
+72
-4
@@ -2,6 +2,8 @@
|
||||
|
||||
#include <capnp/dynamic.h>
|
||||
#include <csignal>
|
||||
#include <iomanip>
|
||||
#include <sstream>
|
||||
#include "cereal/services.h"
|
||||
#include "common/params.h"
|
||||
#include "tools/replay/util.h"
|
||||
@@ -19,6 +21,14 @@ Replay::Replay(const std::string &route, std::vector<std::string> allow, std::ve
|
||||
: sm_(sm), flags_(flags), seg_mgr_(std::make_unique<SegmentManager>(route, flags, data_dir, auto_source)) {
|
||||
std::signal(SIGUSR1, interrupt_sleep_handler);
|
||||
|
||||
if (flags_ & REPLAY_FLAG_BENCHMARK) {
|
||||
benchmark_stats_.process_start_ts = nanos_since_boot();
|
||||
seg_mgr_->setBenchmarkCallback([this](int seg_num, const std::string& event) {
|
||||
benchmark_stats_.timeline.emplace_back(nanos_since_boot(),
|
||||
"segment " + std::to_string(seg_num) + " " + event);
|
||||
});
|
||||
}
|
||||
|
||||
if (!(flags_ & REPLAY_FLAG_ALL_SERVICES)) {
|
||||
block.insert(block.end(), {"bookmarkButton", "uiDebug", "userBookmark"});
|
||||
}
|
||||
@@ -78,8 +88,13 @@ Replay::~Replay() {
|
||||
|
||||
bool Replay::load() {
|
||||
rInfo("loading route %s", seg_mgr_->route_.name().c_str());
|
||||
|
||||
if (!seg_mgr_->load()) return false;
|
||||
|
||||
if (hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
benchmark_stats_.timeline.emplace_back(nanos_since_boot(), "route metadata loaded");
|
||||
}
|
||||
|
||||
min_seconds_ = seg_mgr_->route_.segments().begin()->first * 60;
|
||||
max_seconds_ = (seg_mgr_->route_.segments().rbegin()->first + 1) * 60;
|
||||
return true;
|
||||
@@ -257,8 +272,13 @@ void Replay::streamThread() {
|
||||
stream_thread_id = pthread_self();
|
||||
std::unique_lock lk(stream_lock_);
|
||||
|
||||
int last_processed_segment = -1;
|
||||
uint64_t segment_start_time = 0;
|
||||
bool streaming_started = false;
|
||||
|
||||
while (true) {
|
||||
stream_cv_.wait(lk, [this]() { return exit_ || (events_ready_ && !interrupt_requested_); });
|
||||
|
||||
if (exit_) break;
|
||||
|
||||
event_data_ = seg_mgr_->getEventData();
|
||||
@@ -270,14 +290,19 @@ void Replay::streamThread() {
|
||||
continue;
|
||||
}
|
||||
|
||||
auto it = publishEvents(first, events.cend());
|
||||
if (!streaming_started && hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
benchmark_stats_.timeline.emplace_back(nanos_since_boot(), "streaming started");
|
||||
streaming_started = true;
|
||||
}
|
||||
|
||||
auto it = publishEvents(first, events.cend(), last_processed_segment, segment_start_time);
|
||||
|
||||
// Ensure frames are sent before unlocking to prevent race conditions
|
||||
if (camera_server_) {
|
||||
camera_server_->waitForSent();
|
||||
}
|
||||
|
||||
if (it == events.cend() && !hasFlag(REPLAY_FLAG_NO_LOOP)) {
|
||||
if (it == events.cend() && !hasFlag(REPLAY_FLAG_NO_LOOP) && !hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
int last_segment = seg_mgr_->route_.segments().rbegin()->first;
|
||||
if (event_data_->isSegmentLoaded(last_segment)) {
|
||||
rInfo("reaches the end of route, restart from beginning");
|
||||
@@ -285,12 +310,28 @@ void Replay::streamThread() {
|
||||
seekTo(minSeconds(), false);
|
||||
stream_lock_.lock();
|
||||
}
|
||||
} else if (it == events.cend() && hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
// Exit benchmark mode after first segment completes
|
||||
exit_ = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if (hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
benchmark_stats_.timeline.emplace_back(nanos_since_boot(), "benchmark done");
|
||||
|
||||
{
|
||||
std::unique_lock lock(benchmark_lock_);
|
||||
benchmark_done_ = true;
|
||||
}
|
||||
benchmark_cv_.notify_one();
|
||||
}
|
||||
}
|
||||
|
||||
std::vector<Event>::const_iterator Replay::publishEvents(std::vector<Event>::const_iterator first,
|
||||
std::vector<Event>::const_iterator last) {
|
||||
std::vector<Event>::const_iterator last,
|
||||
int &last_processed_segment,
|
||||
uint64_t &segment_start_time) {
|
||||
uint64_t evt_start_ts = cur_mono_time_;
|
||||
uint64_t loop_start_ts = nanos_since_boot();
|
||||
double prev_replay_speed = speed_;
|
||||
@@ -304,6 +345,23 @@ std::vector<Event>::const_iterator Replay::publishEvents(std::vector<Event>::con
|
||||
seg_mgr_->setCurrentSegment(segment);
|
||||
}
|
||||
|
||||
// Track segment completion for benchmark timeline
|
||||
if (hasFlag(REPLAY_FLAG_BENCHMARK) && segment != last_processed_segment) {
|
||||
if (last_processed_segment >= 0 && segment_start_time > 0) {
|
||||
uint64_t processing_time_ns = nanos_since_boot() - segment_start_time;
|
||||
double processing_time_ms = processing_time_ns / 1e6;
|
||||
double realtime_factor = 60.0 / (processing_time_ns / 1e9); // 60s per segment
|
||||
|
||||
std::ostringstream oss;
|
||||
oss << "segment " << last_processed_segment << " done publishing ("
|
||||
<< std::fixed << std::setprecision(0) << processing_time_ms << " ms, "
|
||||
<< std::fixed << std::setprecision(0) << realtime_factor << "x realtime)";
|
||||
benchmark_stats_.timeline.emplace_back(nanos_since_boot(), oss.str());
|
||||
}
|
||||
segment_start_time = nanos_since_boot();
|
||||
last_processed_segment = segment;
|
||||
}
|
||||
|
||||
cur_mono_time_ = evt.mono_time;
|
||||
cur_which_ = evt.which;
|
||||
|
||||
@@ -320,7 +378,8 @@ std::vector<Event>::const_iterator Replay::publishEvents(std::vector<Event>::con
|
||||
evt_start_ts = evt.mono_time;
|
||||
loop_start_ts = current_nanos;
|
||||
prev_replay_speed = speed_;
|
||||
} else if (time_diff > 0) {
|
||||
} else if (time_diff > 0 && !hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
// Skip sleep in benchmark mode for maximum throughput
|
||||
precise_nano_sleep(time_diff, interrupt_requested_);
|
||||
}
|
||||
|
||||
@@ -338,3 +397,12 @@ std::vector<Event>::const_iterator Replay::publishEvents(std::vector<Event>::con
|
||||
|
||||
return first;
|
||||
}
|
||||
|
||||
void Replay::waitForFinished() {
|
||||
if (!hasFlag(REPLAY_FLAG_BENCHMARK)) {
|
||||
return;
|
||||
}
|
||||
|
||||
std::unique_lock lock(benchmark_lock_);
|
||||
benchmark_cv_.wait(lock, [this]() { return benchmark_done_; });
|
||||
}
|
||||
|
||||
+16
-1
@@ -24,6 +24,12 @@ enum REPLAY_FLAGS {
|
||||
REPLAY_FLAG_NO_HW_DECODER = 0x0100,
|
||||
REPLAY_FLAG_NO_VIPC = 0x0400,
|
||||
REPLAY_FLAG_ALL_SERVICES = 0x0800,
|
||||
REPLAY_FLAG_BENCHMARK = 0x1000,
|
||||
};
|
||||
|
||||
struct BenchmarkStats {
|
||||
uint64_t process_start_ts = 0;
|
||||
std::vector<std::pair<uint64_t, std::string>> timeline;
|
||||
};
|
||||
|
||||
class Replay {
|
||||
@@ -57,6 +63,8 @@ public:
|
||||
inline const std::optional<Timeline::Entry> findAlertAtTime(double sec) const { return timeline_.findAlertAtTime(sec); }
|
||||
const std::shared_ptr<SegmentManager::EventData> getEventData() const { return seg_mgr_->getEventData(); }
|
||||
void installEventFilter(std::function<bool(const Event *)> filter) { event_filter_ = filter; }
|
||||
void waitForFinished();
|
||||
const BenchmarkStats &getBenchmarkStats() const { return benchmark_stats_; }
|
||||
|
||||
// Event callback functions
|
||||
std::function<void()> onSegmentsMerged = nullptr;
|
||||
@@ -72,7 +80,9 @@ private:
|
||||
void handleSegmentMerge();
|
||||
void interruptStream(const std::function<bool()>& update_fn);
|
||||
std::vector<Event>::const_iterator publishEvents(std::vector<Event>::const_iterator first,
|
||||
std::vector<Event>::const_iterator last);
|
||||
std::vector<Event>::const_iterator last,
|
||||
int &last_processed_segment,
|
||||
uint64_t &segment_start_time);
|
||||
void publishMessage(const Event *e);
|
||||
void publishFrame(const Event *e);
|
||||
void checkSeekProgress();
|
||||
@@ -107,4 +117,9 @@ private:
|
||||
std::function<bool(const Event *)> event_filter_ = nullptr;
|
||||
|
||||
std::shared_ptr<SegmentManager::EventData> event_data_ = std::make_shared<SegmentManager::EventData>();
|
||||
|
||||
BenchmarkStats benchmark_stats_;
|
||||
std::condition_variable benchmark_cv_;
|
||||
std::mutex benchmark_lock_;
|
||||
bool benchmark_done_ = false;
|
||||
};
|
||||
|
||||
@@ -118,9 +118,15 @@ void SegmentManager::loadSegmentsInRange(SegmentMap::iterator begin, SegmentMap:
|
||||
for (auto it = first; it != last; ++it) {
|
||||
auto &segment_ptr = it->second;
|
||||
if (!segment_ptr) {
|
||||
if (onBenchmarkEvent_) {
|
||||
onBenchmarkEvent_(it->first, "loading");
|
||||
}
|
||||
segment_ptr = std::make_shared<Segment>(
|
||||
it->first, route_.at(it->first), flags_, filters_,
|
||||
[this](int seg_num, bool success) {
|
||||
if (onBenchmarkEvent_) {
|
||||
onBenchmarkEvent_(seg_num, success ? "loaded" : "load failed");
|
||||
}
|
||||
std::unique_lock lock(mutex_);
|
||||
needs_update_ = true;
|
||||
cv_.notify_one();
|
||||
|
||||
@@ -27,6 +27,7 @@ public:
|
||||
bool load();
|
||||
void setCurrentSegment(int seg_num);
|
||||
void setCallback(const std::function<void()> &callback) { onSegmentMergedCallback_ = callback; }
|
||||
void setBenchmarkCallback(const std::function<void(int, const std::string&)> &callback) { onBenchmarkEvent_ = callback; }
|
||||
void setFilters(const std::vector<bool> &filters) { filters_ = filters; }
|
||||
const std::shared_ptr<EventData> getEventData() const { return std::atomic_load(&event_data_); }
|
||||
bool hasSegment(int n) const { return segments_.find(n) != segments_.end(); }
|
||||
@@ -52,5 +53,6 @@ private:
|
||||
SegmentMap segments_;
|
||||
std::shared_ptr<EventData> event_data_;
|
||||
std::function<void()> onSegmentMergedCallback_ = nullptr;
|
||||
std::function<void(int, const std::string&)> onBenchmarkEvent_ = nullptr;
|
||||
std::set<int> merged_segments_;
|
||||
};
|
||||
|
||||
+143
-101
@@ -5,57 +5,88 @@ import sys
|
||||
|
||||
import cv2
|
||||
import numpy as np
|
||||
import pygame
|
||||
import pyray as rl
|
||||
|
||||
import cereal.messaging as messaging
|
||||
from openpilot.common.basedir import BASEDIR
|
||||
from openpilot.common.transformations.camera import DEVICE_CAMERAS
|
||||
from openpilot.tools.replay.lib.ui_helpers import (UP,
|
||||
BLACK, GREEN,
|
||||
YELLOW, Calibration,
|
||||
get_blank_lid_overlay, init_plots,
|
||||
maybe_update_radar_points, plot_lead,
|
||||
plot_model,
|
||||
pygame_modules_have_loaded)
|
||||
from openpilot.tools.replay.lib.ui_helpers import (
|
||||
UP,
|
||||
BLACK,
|
||||
GREEN,
|
||||
YELLOW,
|
||||
Calibration,
|
||||
get_blank_lid_overlay,
|
||||
init_plots,
|
||||
maybe_update_radar_points,
|
||||
plot_lead,
|
||||
plot_model,
|
||||
)
|
||||
from msgq.visionipc import VisionIpcClient, VisionStreamType
|
||||
|
||||
os.environ['BASEDIR'] = BASEDIR
|
||||
|
||||
ANGLE_SCALE = 5.0
|
||||
|
||||
|
||||
def ui_thread(addr):
|
||||
cv2.setNumThreads(1)
|
||||
pygame.init()
|
||||
pygame.font.init()
|
||||
assert pygame_modules_have_loaded()
|
||||
|
||||
disp_info = pygame.display.Info()
|
||||
max_height = disp_info.current_h
|
||||
# Get monitor info before creating window
|
||||
rl.set_config_flags(rl.ConfigFlags.FLAG_MSAA_4X_HINT)
|
||||
rl.init_window(1, 1, "")
|
||||
max_height = rl.get_monitor_height(0)
|
||||
rl.close_window()
|
||||
|
||||
hor_mode = os.getenv("HORIZONTAL") is not None
|
||||
hor_mode = True if max_height < 960+300 else hor_mode
|
||||
hor_mode = True if max_height < 960 + 300 else hor_mode
|
||||
|
||||
if hor_mode:
|
||||
size = (640+384+640, 960)
|
||||
size = (640 + 384 + 640, 960)
|
||||
write_x = 5
|
||||
write_y = 680
|
||||
else:
|
||||
size = (640+384, 960+300)
|
||||
size = (640 + 384, 960 + 300)
|
||||
write_x = 645
|
||||
write_y = 970
|
||||
|
||||
pygame.display.set_caption("openpilot debug UI")
|
||||
screen = pygame.display.set_mode(size, pygame.DOUBLEBUF)
|
||||
rl.set_trace_log_level(rl.TraceLogLevel.LOG_ERROR)
|
||||
rl.set_config_flags(rl.ConfigFlags.FLAG_MSAA_4X_HINT)
|
||||
rl.init_window(size[0], size[1], "openpilot debug UI")
|
||||
rl.set_target_fps(60)
|
||||
|
||||
alert1_font = pygame.font.SysFont("arial", 30)
|
||||
alert2_font = pygame.font.SysFont("arial", 20)
|
||||
info_font = pygame.font.SysFont("arial", 15)
|
||||
# Load font
|
||||
font_path = os.path.join(BASEDIR, "selfdrive/assets/fonts/JetBrainsMono-Medium.ttf")
|
||||
font = rl.load_font_ex(font_path, 32, None, 0)
|
||||
|
||||
camera_surface = pygame.surface.Surface((640, 480), 0, 24).convert()
|
||||
top_down_surface = pygame.surface.Surface((UP.lidar_x, UP.lidar_y), 0, 8)
|
||||
# Create textures for camera and top-down view
|
||||
camera_image = rl.gen_image_color(640, 480, rl.BLACK)
|
||||
camera_texture = rl.load_texture_from_image(camera_image)
|
||||
rl.unload_image(camera_image)
|
||||
|
||||
sm = messaging.SubMaster(['carState', 'longitudinalPlan', 'carControl', 'radarState', 'liveCalibration', 'controlsState',
|
||||
'selfdriveState', 'liveTracks', 'modelV2', 'liveParameters', 'roadCameraState'], addr=addr)
|
||||
# lid_overlay array is (lidar_x, lidar_y) = (384, 960)
|
||||
# pygame treats first axis as width, so texture is 384 wide x 960 tall
|
||||
# For raylib, we need to transpose to get (height, width) = (960, 384) for the RGBA array
|
||||
top_down_image = rl.gen_image_color(UP.lidar_x, UP.lidar_y, rl.BLACK)
|
||||
top_down_texture = rl.load_texture_from_image(top_down_image)
|
||||
rl.unload_image(top_down_image)
|
||||
|
||||
sm = messaging.SubMaster(
|
||||
[
|
||||
'carState',
|
||||
'longitudinalPlan',
|
||||
'carControl',
|
||||
'radarState',
|
||||
'liveCalibration',
|
||||
'controlsState',
|
||||
'selfdriveState',
|
||||
'liveTracks',
|
||||
'modelV2',
|
||||
'liveParameters',
|
||||
'roadCameraState',
|
||||
],
|
||||
addr=addr,
|
||||
)
|
||||
|
||||
img = np.zeros((480, 640, 3), dtype='uint8')
|
||||
imgff = None
|
||||
@@ -65,74 +96,78 @@ def ui_thread(addr):
|
||||
lid_overlay_blank = get_blank_lid_overlay(UP)
|
||||
|
||||
# plots
|
||||
name_to_arr_idx = { "gas": 0,
|
||||
"computer_gas": 1,
|
||||
"user_brake": 2,
|
||||
"computer_brake": 3,
|
||||
"v_ego": 4,
|
||||
"v_pid": 5,
|
||||
"angle_steers_des": 6,
|
||||
"angle_steers": 7,
|
||||
"angle_steers_k": 8,
|
||||
"steer_torque": 9,
|
||||
"v_override": 10,
|
||||
"v_cruise": 11,
|
||||
"a_ego": 12,
|
||||
"a_target": 13}
|
||||
name_to_arr_idx = {
|
||||
"gas": 0,
|
||||
"computer_gas": 1,
|
||||
"user_brake": 2,
|
||||
"computer_brake": 3,
|
||||
"v_ego": 4,
|
||||
"v_pid": 5,
|
||||
"angle_steers_des": 6,
|
||||
"angle_steers": 7,
|
||||
"angle_steers_k": 8,
|
||||
"steer_torque": 9,
|
||||
"v_override": 10,
|
||||
"v_cruise": 11,
|
||||
"a_ego": 12,
|
||||
"a_target": 13,
|
||||
}
|
||||
|
||||
plot_arr = np.zeros((100, len(name_to_arr_idx.values())))
|
||||
|
||||
plot_xlims = [(0, plot_arr.shape[0]), (0, plot_arr.shape[0]), (0, plot_arr.shape[0]), (0, plot_arr.shape[0])]
|
||||
plot_ylims = [(-0.1, 1.1), (-ANGLE_SCALE, ANGLE_SCALE), (0., 75.), (-3.0, 2.0)]
|
||||
plot_names = [["gas", "computer_gas", "user_brake", "computer_brake"],
|
||||
["angle_steers", "angle_steers_des", "angle_steers_k", "steer_torque"],
|
||||
["v_ego", "v_override", "v_pid", "v_cruise"],
|
||||
["a_ego", "a_target"]]
|
||||
plot_colors = [["b", "b", "g", "r", "y"],
|
||||
["b", "g", "y", "r"],
|
||||
["b", "g", "r", "y"],
|
||||
["b", "r"]]
|
||||
plot_styles = [["-", "-", "-", "-", "-"],
|
||||
["-", "-", "-", "-"],
|
||||
["-", "-", "-", "-"],
|
||||
["-", "-"]]
|
||||
plot_ylims = [(-0.1, 1.1), (-ANGLE_SCALE, ANGLE_SCALE), (0.0, 75.0), (-3.0, 2.0)]
|
||||
plot_names = [
|
||||
["gas", "computer_gas", "user_brake", "computer_brake"],
|
||||
["angle_steers", "angle_steers_des", "angle_steers_k", "steer_torque"],
|
||||
["v_ego", "v_override", "v_pid", "v_cruise"],
|
||||
["a_ego", "a_target"],
|
||||
]
|
||||
plot_colors = [["b", "b", "g", "r", "y"], ["b", "g", "y", "r"], ["b", "g", "r", "y"], ["b", "r"]]
|
||||
plot_styles = [["-", "-", "-", "-", "-"], ["-", "-", "-", "-"], ["-", "-", "-", "-"], ["-", "-"]]
|
||||
|
||||
draw_plots = init_plots(plot_arr, name_to_arr_idx, plot_xlims, plot_ylims, plot_names, plot_colors, plot_styles)
|
||||
|
||||
# Palette for converting lid_overlay grayscale indices to RGBA colors
|
||||
palette = np.zeros((256, 4), dtype=np.uint8)
|
||||
palette[:, 3] = 255 # alpha
|
||||
palette[1] = [255, 0, 0, 255] # RED
|
||||
palette[2] = [0, 255, 0, 255] # GREEN
|
||||
palette[3] = [0, 0, 255, 255] # BLUE
|
||||
palette[4] = [255, 255, 0, 255] # YELLOW
|
||||
palette[110] = [110, 110, 110, 255] # car_color (gray)
|
||||
palette[255] = [255, 255, 255, 255] # WHITE
|
||||
|
||||
vipc_client = VisionIpcClient("camerad", VisionStreamType.VISION_STREAM_ROAD, True)
|
||||
while True:
|
||||
for event in pygame.event.get():
|
||||
if event.type == pygame.QUIT:
|
||||
pygame.quit()
|
||||
sys.exit()
|
||||
|
||||
screen.fill((64, 64, 64))
|
||||
lid_overlay = lid_overlay_blank.copy()
|
||||
top_down = top_down_surface, lid_overlay
|
||||
|
||||
while not rl.window_should_close():
|
||||
# ***** frame *****
|
||||
if not vipc_client.is_connected():
|
||||
vipc_client.connect(True)
|
||||
vipc_client.connect(False)
|
||||
|
||||
rl.begin_drawing()
|
||||
rl.clear_background(rl.Color(64, 64, 64, 255))
|
||||
|
||||
yuv_img_raw = vipc_client.recv()
|
||||
if yuv_img_raw is None or not yuv_img_raw.data.any():
|
||||
rl.draw_text_ex(font, "waiting for frames", rl.Vector2(200, 200), 30, 0, rl.WHITE)
|
||||
rl.end_drawing()
|
||||
continue
|
||||
|
||||
lid_overlay = lid_overlay_blank.copy()
|
||||
top_down = top_down_texture, lid_overlay
|
||||
|
||||
sm.update(0)
|
||||
|
||||
camera = DEVICE_CAMERAS[("tici", str(sm['roadCameraState'].sensor))]
|
||||
|
||||
imgff = np.frombuffer(yuv_img_raw.data, dtype=np.uint8).reshape((len(yuv_img_raw.data) // vipc_client.stride, vipc_client.stride))
|
||||
num_px = vipc_client.width * vipc_client.height
|
||||
rgb = cv2.cvtColor(imgff[:vipc_client.height * 3 // 2, :vipc_client.width], cv2.COLOR_YUV2RGB_NV12)
|
||||
rgb = cv2.cvtColor(imgff[: vipc_client.height * 3 // 2, : vipc_client.width], cv2.COLOR_YUV2RGB_NV12)
|
||||
|
||||
qcam = "QCAM" in os.environ
|
||||
bb_scale = (528 if qcam else camera.fcam.width) / 640.
|
||||
calib_scale = camera.fcam.width / 640.
|
||||
zoom_matrix = np.asarray([
|
||||
[bb_scale, 0., 0.],
|
||||
[0., bb_scale, 0.],
|
||||
[0., 0., 1.]])
|
||||
bb_scale = (528 if qcam else camera.fcam.width) / 640.0
|
||||
calib_scale = camera.fcam.width / 640.0
|
||||
zoom_matrix = np.asarray([[bb_scale, 0.0, 0.0], [0.0, bb_scale, 0.0], [0.0, 0.0, 1.0]])
|
||||
cv2.warpAffine(rgb, zoom_matrix[:2], (img.shape[1], img.shape[0]), dst=img, flags=cv2.WARP_INVERSE_MAP)
|
||||
|
||||
intrinsic_matrix = camera.fcam.intrinsics
|
||||
@@ -151,11 +186,11 @@ def ui_thread(addr):
|
||||
plot_arr[-1, name_to_arr_idx['angle_steers_k']] = angle_steers_k
|
||||
plot_arr[-1, name_to_arr_idx['gas']] = sm['carState'].gasDEPRECATED
|
||||
# TODO gas is deprecated
|
||||
plot_arr[-1, name_to_arr_idx['computer_gas']] = np.clip(sm['carControl'].actuators.accel/4.0, 0.0, 1.0)
|
||||
plot_arr[-1, name_to_arr_idx['computer_gas']] = np.clip(sm['carControl'].actuators.accel / 4.0, 0.0, 1.0)
|
||||
plot_arr[-1, name_to_arr_idx['user_brake']] = sm['carState'].brake
|
||||
plot_arr[-1, name_to_arr_idx['steer_torque']] = sm['carControl'].actuators.torque * ANGLE_SCALE
|
||||
# TODO brake is deprecated
|
||||
plot_arr[-1, name_to_arr_idx['computer_brake']] = np.clip(-sm['carControl'].actuators.accel/4.0, 0.0, 1.0)
|
||||
plot_arr[-1, name_to_arr_idx['computer_brake']] = np.clip(-sm['carControl'].actuators.accel / 4.0, 0.0, 1.0)
|
||||
plot_arr[-1, name_to_arr_idx['v_ego']] = sm['carState'].vEgo
|
||||
plot_arr[-1, name_to_arr_idx['v_cruise']] = sm['carState'].cruiseState.speed
|
||||
plot_arr[-1, name_to_arr_idx['a_ego']] = sm['carState'].aEgo
|
||||
@@ -177,56 +212,63 @@ def ui_thread(addr):
|
||||
calibration = Calibration(num_px, rpyCalib, intrinsic_matrix, calib_scale)
|
||||
|
||||
# *** blits ***
|
||||
pygame.surfarray.blit_array(camera_surface, img.swapaxes(0, 1))
|
||||
screen.blit(camera_surface, (0, 0))
|
||||
# Update camera texture from numpy array
|
||||
img_rgba = cv2.cvtColor(img, cv2.COLOR_RGB2RGBA)
|
||||
rl.update_texture(camera_texture, rl.ffi.cast("void *", img_rgba.ctypes.data))
|
||||
rl.draw_texture(camera_texture, 0, 0, rl.WHITE)
|
||||
|
||||
# display alerts
|
||||
alert_line1 = alert1_font.render(sm['selfdriveState'].alertText1, True, (255, 0, 0))
|
||||
alert_line2 = alert2_font.render(sm['selfdriveState'].alertText2, True, (255, 0, 0))
|
||||
screen.blit(alert_line1, (180, 150))
|
||||
screen.blit(alert_line2, (180, 190))
|
||||
rl.draw_text_ex(font, sm['selfdriveState'].alertText1, rl.Vector2(180, 150), 30, 0, rl.RED)
|
||||
rl.draw_text_ex(font, sm['selfdriveState'].alertText2, rl.Vector2(180, 190), 20, 0, rl.RED)
|
||||
|
||||
# draw plots (texture is reused internally)
|
||||
plot_texture = draw_plots(plot_arr)
|
||||
if hor_mode:
|
||||
screen.blit(draw_plots(plot_arr), (640+384, 0))
|
||||
rl.draw_texture(plot_texture, 640 + 384, 0, rl.WHITE)
|
||||
else:
|
||||
screen.blit(draw_plots(plot_arr), (0, 600))
|
||||
rl.draw_texture(plot_texture, 0, 600, rl.WHITE)
|
||||
|
||||
pygame.surfarray.blit_array(*top_down)
|
||||
screen.blit(top_down[0], (640, 0))
|
||||
# Convert lid_overlay to RGBA and update top_down texture
|
||||
# lid_overlay is (384, 960), need to transpose to (960, 384) for row-major RGBA buffer
|
||||
lid_rgba = palette[lid_overlay.T]
|
||||
rl.update_texture(top_down_texture, rl.ffi.cast("void *", np.ascontiguousarray(lid_rgba).ctypes.data))
|
||||
rl.draw_texture(top_down_texture, 640, 0, rl.WHITE)
|
||||
|
||||
SPACING = 25
|
||||
|
||||
lines = [
|
||||
info_font.render("ENABLED", True, GREEN if sm['selfdriveState'].enabled else BLACK),
|
||||
info_font.render("SPEED: " + str(round(sm['carState'].vEgo, 1)) + " m/s", True, YELLOW),
|
||||
info_font.render("LONG CONTROL STATE: " + str(sm['controlsState'].longControlState), True, YELLOW),
|
||||
info_font.render("LONG MPC SOURCE: " + str(sm['longitudinalPlan'].longitudinalPlanSource), True, YELLOW),
|
||||
("ENABLED", GREEN if sm['selfdriveState'].enabled else BLACK),
|
||||
("SPEED: " + str(round(sm['carState'].vEgo, 1)) + " m/s", YELLOW),
|
||||
("LONG CONTROL STATE: " + str(sm['controlsState'].longControlState), YELLOW),
|
||||
("LONG MPC SOURCE: " + str(sm['longitudinalPlan'].longitudinalPlanSource), YELLOW),
|
||||
None,
|
||||
info_font.render("ANGLE OFFSET (AVG): " + str(round(sm['liveParameters'].angleOffsetAverageDeg, 2)) + " deg", True, YELLOW),
|
||||
info_font.render("ANGLE OFFSET (INSTANT): " + str(round(sm['liveParameters'].angleOffsetDeg, 2)) + " deg", True, YELLOW),
|
||||
info_font.render("STIFFNESS: " + str(round(sm['liveParameters'].stiffnessFactor * 100., 2)) + " %", True, YELLOW),
|
||||
info_font.render("STEER RATIO: " + str(round(sm['liveParameters'].steerRatio, 2)), True, YELLOW)
|
||||
("ANGLE OFFSET (AVG): " + str(round(sm['liveParameters'].angleOffsetAverageDeg, 2)) + " deg", YELLOW),
|
||||
("ANGLE OFFSET (INSTANT): " + str(round(sm['liveParameters'].angleOffsetDeg, 2)) + " deg", YELLOW),
|
||||
("STIFFNESS: " + str(round(sm['liveParameters'].stiffnessFactor * 100.0, 2)) + " %", YELLOW),
|
||||
("STEER RATIO: " + str(round(sm['liveParameters'].steerRatio, 2)), YELLOW),
|
||||
]
|
||||
|
||||
for i, line in enumerate(lines):
|
||||
if line is not None:
|
||||
screen.blit(line, (write_x, write_y + i * SPACING))
|
||||
color = rl.Color(line[1][0], line[1][1], line[1][2], 255)
|
||||
rl.draw_text_ex(font, line[0], rl.Vector2(write_x, write_y + i * SPACING), 20, 0, color)
|
||||
|
||||
rl.end_drawing()
|
||||
|
||||
rl.unload_texture(camera_texture)
|
||||
rl.unload_texture(top_down_texture)
|
||||
rl.unload_font(font)
|
||||
rl.close_window()
|
||||
|
||||
# this takes time...vsync or something
|
||||
pygame.display.flip()
|
||||
|
||||
def get_arg_parser():
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Show replay data in a UI.",
|
||||
formatter_class=argparse.ArgumentDefaultsHelpFormatter)
|
||||
parser = argparse.ArgumentParser(description="Show replay data in a UI.", formatter_class=argparse.ArgumentDefaultsHelpFormatter)
|
||||
|
||||
parser.add_argument("ip_address", nargs="?", default="127.0.0.1",
|
||||
help="The ip address on which to receive zmq messages.")
|
||||
parser.add_argument("ip_address", nargs="?", default="127.0.0.1", help="The ip address on which to receive zmq messages.")
|
||||
|
||||
parser.add_argument("--frame-address", default=None,
|
||||
help="The frame address (fully qualified ZMQ endpoint for frames) on which to receive zmq messages.")
|
||||
parser.add_argument("--frame-address", default=None, help="The frame address (fully qualified ZMQ endpoint for frames) on which to receive zmq messages.")
|
||||
return parser
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
args = get_arg_parser().parse_args(sys.argv[1:])
|
||||
|
||||
|
||||
@@ -154,6 +154,8 @@ size_t getRemoteFileSize(const std::string &url, std::atomic<bool> *abort) {
|
||||
curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, dumy_write_cb);
|
||||
curl_easy_setopt(curl, CURLOPT_HEADER, 1);
|
||||
curl_easy_setopt(curl, CURLOPT_NOBODY, 1);
|
||||
curl_easy_setopt(curl, CURLOPT_NOSIGNAL, 1);
|
||||
curl_easy_setopt(curl, CURLOPT_FOLLOWLOCATION, 1L);
|
||||
|
||||
CURLM *cm = curl_multi_init();
|
||||
curl_multi_add_handle(cm, curl);
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
#!/usr/bin/env python3
|
||||
# type: ignore
|
||||
|
||||
import os
|
||||
import time
|
||||
|
||||
Reference in New Issue
Block a user