From 67a7c148b45860d574d48af2a8de30a4a1ea45d1 Mon Sep 17 00:00:00 2001 From: firestar5683 <168790843+firestar5683@users.noreply.github.com> Date: Sat, 19 Sep 2026 17:27:00 -0700 Subject: [PATCH] Prime native replay camera and initial cache before starting source clock --- tools/replay/camera.cc | 9 ++++++ tools/replay/camera.h | 1 + tools/replay/replay.cc | 35 ++++++++++++++++++++++++ tools/replay/replay.h | 2 ++ tools/replay/startup_prime.h | 15 ++++++++++ tools/replay/tests/test_startup_prime.cc | 14 ++++++++++ 6 files changed, 76 insertions(+) create mode 100644 tools/replay/startup_prime.h create mode 100644 tools/replay/tests/test_startup_prime.cc diff --git a/tools/replay/camera.cc b/tools/replay/camera.cc index 73243ed20d..f6b7dee375 100644 --- a/tools/replay/camera.cc +++ b/tools/replay/camera.cc @@ -105,6 +105,15 @@ VisionBuf *CameraServer::getFrame(Camera &cam, FrameReader *fr, int32_t segment_ return nullptr; } +bool CameraServer::primeFrame(CameraType type, FrameReader *fr, const Event *event) { + // Before first publication the camera queue is empty; decode without sending a future frame. + waitForSent(); + capnp::FlatArrayMessageReader reader(event->data); + auto evt = reader.getRoot(); + auto eidx = capnp::AnyStruct::Reader(evt).getPointerSection()[0].getAs(); + return getFrame(cameras_[type], fr, eidx.getSegmentId(), eidx.getFrameId()) != nullptr; +} + void CameraServer::pushFrame(CameraType type, FrameReader *fr, const Event *event) { auto &cam = cameras_[type]; if (cam.width != fr->width || cam.height != fr->height) { diff --git a/tools/replay/camera.h b/tools/replay/camera.h index 21c3d98dcf..130282578b 100644 --- a/tools/replay/camera.h +++ b/tools/replay/camera.h @@ -16,6 +16,7 @@ class CameraServer { public: CameraServer(std::pair camera_size[MAX_CAMERAS] = nullptr); ~CameraServer(); + bool primeFrame(CameraType type, FrameReader* fr, const Event *event); void pushFrame(CameraType type, FrameReader* fr, const Event *event); void waitForSent(); diff --git a/tools/replay/replay.cc b/tools/replay/replay.cc index cc105dd10e..bf9c52f0d2 100644 --- a/tools/replay/replay.cc +++ b/tools/replay/replay.cc @@ -2,6 +2,9 @@ #include #include +#include +#include +#include "tools/replay/startup_prime.h" #include "cereal/services.h" #include "common/params.h" #include "tools/replay/util.h" @@ -18,6 +21,8 @@ Replay::Replay(const std::string &route, std::vector allow, std::ve SubMaster *sm, uint32_t flags, const std::string &data_dir, bool auto_source) : sm_(sm), flags_(flags), seg_mgr_(std::make_unique(route, flags, data_dir, auto_source)) { std::signal(SIGUSR1, interrupt_sleep_handler); + const char *prime = std::getenv("ROADSCORE_REPLAY_PRIME"); + startup_prime_ = prime && std::string(prime) == "1"; if (!(flags_ & REPLAY_FLAG_ALL_SERVICES)) { block.insert(block.end(), {"bookmarkButton", "uiDebug", "userBookmark"}); @@ -91,7 +96,9 @@ void Replay::interruptStream(const std::function &update_fn) { } { interrupt_requested_ = true; + const uint64_t lock_start = nanos_since_boot(); std::unique_lock lock(stream_lock_); + if (startup_prime_) rInfo("REPLAY_STREAM_INTERRUPT_WAIT seconds=%.6f", (nanos_since_boot() - lock_start) / 1e9); events_ready_ = update_fn(); interrupt_requested_ = user_paused_; } @@ -125,6 +132,15 @@ void Replay::seekTo(double seconds, bool relative) { void Replay::checkSeekProgress() { if (!seg_mgr_->getEventData()->isSegmentLoaded(current_segment_.load())) return; + if (startup_prime_ && !startup_primed_) { + std::vector available, loaded; + for (const auto &[n, unused] : seg_mgr_->route_.segments()) available.push_back(n); + for (const auto &[n, unused] : seg_mgr_->getEventData()->segments) loaded.push_back(n); + if (!startup_cache_ready(available, loaded, current_segment_.load(), seg_mgr_->segment_cache_limit_)) { + rInfo("REPLAY_STARTUP_CACHE_WAIT loaded=%zu cache=%d", loaded.size(), seg_mgr_->segment_cache_limit_); + return; + } + } double seek_to = seeking_to_.exchange(-1.0, std::memory_order_acquire); if (seek_to >= 0 && onSeekedTo) { onSeekedTo(seek_to); @@ -270,11 +286,29 @@ void Replay::streamThread() { continue; } + if (startup_prime_ && !startup_primed_) { + const uint64_t prime_start = nanos_since_boot(); + if (camera_server_) { + auto camera = std::find_if(first, events.cend(), [](const Event &e) { + return e.which == cereal::Event::ROAD_ENCODE_IDX && e.eidx_segnum >= 0; + }); + if (camera == events.cend()) throw std::runtime_error("Startup prime has no road camera event"); + auto segment = event_data_->segments.find(camera->eidx_segnum); + if (segment == event_data_->segments.end() || !segment->second->frames[RoadCam] || + !camera_server_->primeFrame(RoadCam, segment->second->frames[RoadCam].get(), &*camera)) { + throw std::runtime_error("Startup road camera prime failed"); + } + } + startup_primed_ = true; + rInfo("REPLAY_STARTUP_PRIMED seconds=%.6f", (nanos_since_boot() - prime_start) / 1e9); + } auto it = publishEvents(first, events.cend()); // Ensure frames are sent before unlocking to prevent race conditions if (camera_server_) { + const uint64_t drain_start = nanos_since_boot(); camera_server_->waitForSent(); + if (startup_prime_) rInfo("REPLAY_CAMERA_DRAIN seconds=%.6f", (nanos_since_boot() - drain_start) / 1e9); } if (it == events.cend() && !hasFlag(REPLAY_FLAG_NO_LOOP)) { @@ -317,6 +351,7 @@ std::vector::const_iterator Replay::publishEvents(std::vector::con // - A negative time_diff may indicate slow execution or system wake-up, // - A time_diff exceeding 1 second suggests a skipped segment. if ((time_diff < -1e9 || time_diff >= 1e9) || speed_ != prev_replay_speed) { + if (startup_prime_) rWarning("REPLAY_CLOCK_REANCHOR lag_seconds=%.6f", -time_diff / 1e9); evt_start_ts = evt.mono_time; loop_start_ts = current_nanos; prev_replay_speed = speed_; diff --git a/tools/replay/replay.h b/tools/replay/replay.h index 5e868d2427..a9d77516fb 100644 --- a/tools/replay/replay.h +++ b/tools/replay/replay.h @@ -85,6 +85,8 @@ private: std::mutex stream_lock_; bool user_paused_ = false; std::condition_variable stream_cv_; + bool startup_prime_ = false; + std::atomic startup_primed_ = false; std::atomic current_segment_ = 0; std::atomic seeking_to_ = -1.0; std::atomic exit_ = false; diff --git a/tools/replay/startup_prime.h b/tools/replay/startup_prime.h new file mode 100644 index 0000000000..c273165b95 --- /dev/null +++ b/tools/replay/startup_prime.h @@ -0,0 +1,15 @@ +#pragma once +#include +#include + +inline bool startup_cache_ready(const std::vector &route, const std::vector &loaded, + int current, int cache_limit) { + auto cur = std::lower_bound(route.begin(), route.end(), current); + if (cur == route.end() || *cur != current || cache_limit <= 0) return false; + auto begin = cur - std::min(cache_limit / 2, cur - route.begin()); + auto end = begin + std::min(cache_limit, route.end() - begin); + begin = end - std::min(cache_limit, end - route.begin()); + return std::all_of(begin, end, [&](int segment) { + return std::binary_search(loaded.begin(), loaded.end(), segment); + }); +} diff --git a/tools/replay/tests/test_startup_prime.cc b/tools/replay/tests/test_startup_prime.cc new file mode 100644 index 0000000000..d3b473fa3d --- /dev/null +++ b/tools/replay/tests/test_startup_prime.cc @@ -0,0 +1,14 @@ +#include +#include "tools/replay/startup_prime.h" +int main() { + std::vector route{0,1,2,3,4,5,6,7}; + assert(!startup_cache_ready(route,{2},2,5)); + assert(!startup_cache_ready(route,{2,3,4},2,5)); + assert(startup_cache_ready(route,{0,1,2,3,4},2,5)); + assert(!startup_cache_ready(route,{0,1,2,3,4},7,5)); + assert(startup_cache_ready(route,{3,4,5,6,7},7,5)); + assert(startup_cache_ready({0,1},{0,1},0,5)); + assert(!startup_cache_ready({0,2,3},{0,2},2,5)); + assert(startup_cache_ready({0,2,3},{0,2,3},2,5)); + assert(!startup_cache_ready(route,{0,1,2,3,4},2,0)); +}