diff --git a/openpilot/selfdrive/test/process_replay/migration.py b/openpilot/selfdrive/test/process_replay/migration.py index e38d1e1042bdc6..13ab47f31a0bb5 100644 --- a/openpilot/selfdrive/test/process_replay/migration.py +++ b/openpilot/selfdrive/test/process_replay/migration.py @@ -4,6 +4,8 @@ import capnp import functools import traceback +import sys +import argparse from openpilot.cereal import messaging, log from opendbc.car.structs import car @@ -17,7 +19,7 @@ from openpilot.selfdrive.test.process_replay.vision_meta import meta_from_encode_index from openpilot.selfdrive.controls.lib.drive_helpers import CONTROL_N, get_accel_from_plan, should_stop from openpilot.system.manager.process_config import managed_processes -from openpilot.tools.lib.logreader import LogIterable +from openpilot.tools.lib.logreader import LogIterable, LogReader, save_log MessageWithIndex = tuple[int, capnp.lib.capnp._DynamicStructReader] MigrationOps = tuple[list[tuple[int, capnp.lib.capnp._DynamicStructReader]], list[capnp.lib.capnp._DynamicStructReader], list[int]] @@ -387,7 +389,7 @@ def migrate_cameraStates(msgs): encode_id = frame_to_encode_id[msg.which()].get(camera_state.frameId) if encode_id is None: - print(f"Missing encoded frame for camera feed {msg.which()} with frameId: {camera_state.frameId}") + print(f"Missing encoded frame for camera feed {msg.which()} with frameId: {camera_state.frameId}", file=sys.stderr) if len(frame_to_encode_id[msg.which()]) != 0: del_ops.append(index) continue @@ -395,7 +397,7 @@ def migrate_cameraStates(msgs): # fallback mechanism for logs without encodeIdx (e.g. logs from before 2022 with dcamera recording disabled) # try to fake encode_id by subtracting lowest frameId encode_id = camera_state.frameId - min_frame_id[msg.which()] - print(f"Faking encodeId to {encode_id} for camera feed {msg.which()} with frameId: {camera_state.frameId}") + print(f"Faking encodeId to {encode_id} for camera feed {msg.which()} with frameId: {camera_state.frameId}", file=sys.stderr) new_msg = messaging.new_message(msg.which()) new_camera_state = getattr(new_msg, new_msg.which()) @@ -527,3 +529,15 @@ def migrate_driverMonitoringState(msgs): ops.append((index, as_reader(new_msg))) return ops, [], [] + + +if __name__ == '__main__': + parser = argparse.ArgumentParser(description="Migrate logs") + parser.add_argument("input_path", help="Segment identifier or path to file") + parser.add_argument("output_path", help="Path to output file") + args = parser.parse_args() + + output_path = args.output_path if args.output_path != "-" else "/dev/stdout" + mlr = migrate_all(LogReader(args.input_path)) + print(f"Saving migrated log to {output_path}", file=sys.stderr) + save_log(output_path, mlr) diff --git a/openpilot/tools/jotpluggler/app.cc b/openpilot/tools/jotpluggler/app.cc index e6ba696bae8c95..0d9375b9f6c84d 100644 --- a/openpilot/tools/jotpluggler/app.cc +++ b/openpilot/tools/jotpluggler/app.cc @@ -1815,6 +1815,7 @@ int run(const Options &options) { .route_name = options.route_name, .data_dir = options.data_dir, .dbc_override = {}, + .migrate = options.migrate, .stream_source = StreamSourceConfig{.kind = is_local_stream_address(options.stream_address) ? StreamSourceKind::CerealLocal : StreamSourceKind::CerealRemote, diff --git a/openpilot/tools/jotpluggler/app.h b/openpilot/tools/jotpluggler/app.h index b7754f51b2346a..c9f0192c237ffc 100644 --- a/openpilot/tools/jotpluggler/app.h +++ b/openpilot/tools/jotpluggler/app.h @@ -34,6 +34,7 @@ struct Options { bool show = false; bool sync_load = false; bool stream = false; + bool migrate = true; double stream_buffer_seconds = 30.0; }; @@ -411,6 +412,7 @@ std::vector decode_can_messages(const std::vector & RouteData load_route_data(const std::string &route_name, const std::string &data_dir = {}, const std::string &dbc_name = {}, + bool migrate = true, const RouteLoadProgressCallback &progress = {}); RouteIdentifier parse_route_identifier(std::string_view route_name); void rebuild_gps_trace(RouteData *route_data); @@ -484,6 +486,7 @@ struct AppSession { std::string route_name; std::string data_dir; std::string dbc_override; + bool migrate = true; StreamSourceConfig stream_source; double stream_buffer_seconds = 30.0; SessionDataMode data_mode = SessionDataMode::Route; @@ -846,7 +849,7 @@ class AsyncRouteLoader { AsyncRouteLoader(const AsyncRouteLoader &) = delete; AsyncRouteLoader &operator=(const AsyncRouteLoader &) = delete; - void start(const std::string &route_name, const std::string &data_dir, const std::string &dbc_name); + void start(const std::string &route_name, const std::string &data_dir, const std::string &dbc_name, bool migrate); RouteLoadSnapshot snapshot() const; bool consume(RouteData *route_data, std::string *error_text); diff --git a/openpilot/tools/jotpluggler/layout.cc b/openpilot/tools/jotpluggler/layout.cc index 6bc3a6168a4fc0..091aea5abebb98 100644 --- a/openpilot/tools/jotpluggler/layout.cc +++ b/openpilot/tools/jotpluggler/layout.cc @@ -547,7 +547,7 @@ bool save_layout(AppSession *session, UiState *state, const std::string &layout_ void rebuild_session_route_data(AppSession *session, UiState *state, const RouteLoadProgressCallback &progress) { - apply_route_data(session, state, load_route_data(session->route_name, session->data_dir, session->dbc_override, progress)); + apply_route_data(session, state, load_route_data(session->route_name, session->data_dir, session->dbc_override, session->migrate, progress)); } void stop_stream_session(AppSession *session, UiState *state, bool preserve_data) { @@ -627,7 +627,7 @@ void start_async_route_load(AppSession *session, UiState *state) { return; } apply_route_data(session, state, RouteData{}); - session->route_loader->start(session->route_name, session->data_dir, session->dbc_override); + session->route_loader->start(session->route_name, session->data_dir, session->dbc_override, session->migrate); state->status_text = session->route_name.empty() ? "Ready" : "Loading route " + session->route_name; } diff --git a/openpilot/tools/jotpluggler/main.cc b/openpilot/tools/jotpluggler/main.cc index 22bc29664c2616..89cb8989e810e3 100644 --- a/openpilot/tools/jotpluggler/main.cc +++ b/openpilot/tools/jotpluggler/main.cc @@ -22,6 +22,7 @@ void print_usage(const char *argv0) { << " --output \n" << " --show\n" << " --sync-load\n" + << " --no-migration\n" << "\n" << "Examples:\n" << " " << argv0 << "\n" @@ -94,6 +95,8 @@ int main(int argc, char *argv[]) { options.show = true; } else if (arg == "--sync-load") { options.sync_load = true; + } else if (arg == "--no-migration") { + options.migrate = false; } else if (arg == "--help" || arg == "-h") { print_usage(argv[0]); return 0; diff --git a/openpilot/tools/jotpluggler/runtime.cc b/openpilot/tools/jotpluggler/runtime.cc index a1e47c7e8ea04f..e7e8f3708dc230 100644 --- a/openpilot/tools/jotpluggler/runtime.cc +++ b/openpilot/tools/jotpluggler/runtime.cc @@ -260,13 +260,14 @@ struct AsyncRouteLoader::Impl { join(); } - void start(const std::string &route_name_value, const std::string &data_dir_value, const std::string &dbc_name_value) { + void start(const std::string &route_name_value, const std::string &data_dir_value, const std::string &dbc_name_value, bool migrate_value) { join(); { std::lock_guard lock(mutex); route_name = route_name_value; data_dir = data_dir_value; dbc_name = dbc_name_value; + migrate = migrate_value; result.reset(); error_text.clear(); } @@ -282,7 +283,7 @@ struct AsyncRouteLoader::Impl { worker = std::thread([this]() { try { - RouteData route_data = load_route_data(route_name, data_dir, dbc_name, [this](const RouteLoadProgress &progress) { + RouteData route_data = load_route_data(route_name, data_dir, dbc_name, migrate, [this](const RouteLoadProgress &progress) { total_segments.store(progress.total_segments > 0 ? progress.total_segments : progress.segment_count); segments_downloaded.store(progress.segments_downloaded); segments_parsed.store(progress.segments_parsed); @@ -346,6 +347,7 @@ struct AsyncRouteLoader::Impl { std::string route_name; std::string data_dir; std::string dbc_name; + bool migrate = true; std::string error_text; std::atomic active{false}; std::atomic completed{false}; @@ -361,8 +363,8 @@ AsyncRouteLoader::AsyncRouteLoader(bool enable_terminal_progress) AsyncRouteLoader::~AsyncRouteLoader() = default; -void AsyncRouteLoader::start(const std::string &route_name, const std::string &data_dir, const std::string &dbc_name) { - impl_->start(route_name, data_dir, dbc_name); +void AsyncRouteLoader::start(const std::string &route_name, const std::string &data_dir, const std::string &dbc_name, bool migrate) { + impl_->start(route_name, data_dir, dbc_name, migrate); } RouteLoadSnapshot AsyncRouteLoader::snapshot() const { diff --git a/openpilot/tools/jotpluggler/sketch_layout.cc b/openpilot/tools/jotpluggler/sketch_layout.cc index ddb38fd473208c..f64c1f735a4c5b 100644 --- a/openpilot/tools/jotpluggler/sketch_layout.cc +++ b/openpilot/tools/jotpluggler/sketch_layout.cc @@ -26,6 +26,7 @@ #include "json11/json11.hpp" #include "tools/replay/logreader.h" #include "tools/replay/py_downloader.h" +#include "tools/replay/py_process.h" namespace fs = std::filesystem; @@ -1551,12 +1552,17 @@ SeriesAccumulator extract_segment_series(const std::vector &events, return merged; } +std::string migrate_log(const std::string &log_path) { + return PyProcess::runModule("openpilot.selfdrive.test.process_replay.migration", {log_path, "-"}, nullptr, false); +} + LoadedRouteArtifacts load_route_series_parallel( const std::map &segments, const SchemaIndex &schema, const dbc::Database *can_dbc, LogSelector selector, bool skip_raw_can, + bool migrate, LoadStats *stats) { struct SegmentResult { SeriesAccumulator series; @@ -1600,8 +1606,16 @@ LoadedRouteArtifacts load_route_series_parallel( continue; } + std::string migrated_data; + if (migrate) { + migrated_data = migrate_log(log_path); + } + LogReader reader; - if (!reader.load(log_path, nullptr, true)) { + const bool loaded = !migrated_data.empty() + ? reader.load(migrated_data.data(), migrated_data.size()) + : reader.load(log_path, nullptr, true); + if (!loaded) { segment_stats.failed = true; std::lock_guard lock(error_mutex); if (first_error.empty()) { @@ -1849,6 +1863,7 @@ SketchLayout load_sketch_layout(const fs::path &layout_path) { RouteData load_route_data(const std::string &route_name, const std::string &data_dir, const std::string &dbc_name, + bool migrate, const RouteLoadProgressCallback &progress) { if (route_name.empty()) return RouteData{}; @@ -1876,7 +1891,7 @@ RouteData load_route_data(const std::string &route_name, const SchemaIndex &schema = SchemaIndex::instance(); LoadedRouteArtifacts artifacts = load_route_series_parallel(segments, schema, can_dbc ? &*can_dbc : nullptr, - route.selector, can_dbc.has_value(), &stats); + route.selector, can_dbc.has_value(), migrate, &stats); RouteData route_data = build_route_data(std::move(artifacts.series), std::move(artifacts.can_messages), std::move(artifacts.logs), diff --git a/openpilot/tools/replay/SConscript b/openpilot/tools/replay/SConscript index 9e060bc82946fb..6dfb45e88682a6 100644 --- a/openpilot/tools/replay/SConscript +++ b/openpilot/tools/replay/SConscript @@ -7,7 +7,7 @@ base_frameworks = ['VideoToolbox', 'CoreMedia', 'CoreFoundation', 'CoreVideo'] i base_libs = [common, messaging, cereal, visionipc, 'm', 'pthread'] replay_lib_src = ["replay.cc", "consoleui.cc", "camera.cc", "filereader.cc", "logreader.cc", "framereader.cc", - "route.cc", "util.cc", "seg_mgr.cc", "timeline.cc", "py_downloader.cc"] + "route.cc", "util.cc", "seg_mgr.cc", "timeline.cc", "py_process.cc", "py_downloader.cc"] if arch != "Darwin": replay_lib_src.append("#openpilot/system/loggerd/encoder/v4l_decoder.cc") replay_lib = replay_env.Library("replay", replay_lib_src, LIBS=base_libs, FRAMEWORKS=base_frameworks) diff --git a/openpilot/tools/replay/py_downloader.cc b/openpilot/tools/replay/py_downloader.cc index a265dfd6a37d89..d1126971362e28 100644 --- a/openpilot/tools/replay/py_downloader.cc +++ b/openpilot/tools/replay/py_downloader.cc @@ -1,138 +1,23 @@ #include "tools/replay/py_downloader.h" -#include -#include #include -#include -#include #include -#include "tools/replay/util.h" +#include "tools/replay/py_process.h" namespace { static std::mutex handler_mutex; static DownloadProgressHandler progress_handler = nullptr; -// Run a Python command and capture stdout. Stderr is left attached to the parent. -// Returns stdout content. If abort is signaled, kills the child process. -std::string runPython(const std::vector &args, std::atomic *abort = nullptr) { - // Build argv for execvp - std::vector argv; - argv.push_back("python3"); - argv.push_back("-m"); - argv.push_back("openpilot.tools.lib.file_downloader"); - for (const auto &a : args) { - argv.push_back(a.c_str()); - } - argv.push_back(nullptr); - - int stdout_pipe[2]; - if (pipe(stdout_pipe) != 0) { - rWarning("py_downloader: pipe() failed"); - return {}; - } - - pid_t pid = fork(); - if (pid < 0) { - rWarning("py_downloader: fork() failed"); - close(stdout_pipe[0]); close(stdout_pipe[1]); - return {}; - } - - if (pid == 0) { - // Child process — detach from controlling terminal so Python - // cannot corrupt terminal settings needed by ncurses in the parent. - setsid(); - int devnull = open("/dev/null", O_RDONLY); - if (devnull >= 0) { - dup2(devnull, STDIN_FILENO); - if (devnull > STDERR_FILENO) close(devnull); - } - - // Clear OPENPILOT_PREFIX so the Python process uses default paths - // (e.g. ~/.comma/auth.json). The prefix is only for IPC in the parent. - unsetenv("OPENPILOT_PREFIX"); - - close(stdout_pipe[0]); - dup2(stdout_pipe[1], STDOUT_FILENO); - close(stdout_pipe[1]); - - execvp("python3", const_cast(argv.data())); - _exit(127); - } - - // Parent process - close(stdout_pipe[1]); - - std::string stdout_data; - char buf[4096]; - - // Use select() so abort can interrupt while waiting for Python output. - fd_set rfds; - bool stdout_open = true; - - while (stdout_open) { - if (abort && *abort) { - kill(pid, SIGTERM); - break; - } - - FD_ZERO(&rfds); - FD_SET(stdout_pipe[0], &rfds); - - struct timeval tv = {0, 100000}; // 100ms timeout - int ret = select(stdout_pipe[0] + 1, &rfds, nullptr, nullptr, &tv); - if (ret < 0) break; - - if (FD_ISSET(stdout_pipe[0], &rfds)) { - ssize_t n = read(stdout_pipe[0], buf, sizeof(buf)); - if (n <= 0) { - stdout_open = false; - } else { - stdout_data.append(buf, n); - } - } - } - - // Drain remaining pipe data to prevent child from blocking on write - while (true) { - ssize_t n = read(stdout_pipe[0], buf, sizeof(buf)); - if (n <= 0) break; - stdout_data.append(buf, n); - } - close(stdout_pipe[0]); - - int status; - waitpid(pid, &status, 0); - - const bool aborted = abort && *abort; - const bool expected_sigterm = aborted && WIFSIGNALED(status) && WTERMSIG(status) == SIGTERM; - bool failed = aborted || - (WIFEXITED(status) && WEXITSTATUS(status) != 0) || - WIFSIGNALED(status); - if (failed) { - if (expected_sigterm) { - // Route/camera teardown cancels outstanding downloader subprocesses. - // Keep that expected shutdown path quiet. - } else if (WIFEXITED(status) && WEXITSTATUS(status) != 0) { - rWarning("py_downloader: process exited with code %d", WEXITSTATUS(status)); - } else if (WIFSIGNALED(status)) { - rWarning("py_downloader: process killed by signal %d", WTERMSIG(status)); - } +// Run the file_downloader module and notify the progress handler on failure. +std::string runDownloader(const std::vector &args, std::atomic *abort = nullptr) { + std::string result = PyProcess::runModule("openpilot.tools.lib.file_downloader", args, abort); + if (result.empty()) { std::lock_guard lk(handler_mutex); - if (progress_handler) { - progress_handler(0, 0, false); - } - return {}; + if (progress_handler) progress_handler(0, 0, false); } - - // Trim trailing newline - while (!stdout_data.empty() && (stdout_data.back() == '\n' || stdout_data.back() == '\r')) { - stdout_data.pop_back(); - } - - return stdout_data; + return result; } } // namespace @@ -149,19 +34,19 @@ std::string download(const std::string &url, bool use_cache, std::atomic * if (!use_cache) { args.push_back("--no-cache"); } - return runPython(args, abort); + return runDownloader(args, abort); } std::string decompress(const std::string &path, std::atomic *abort) { - return runPython({"decompress", path}, abort); + return runDownloader({"decompress", path}, abort); } std::string getRouteFiles(const std::string &route) { - return runPython({"route-files", route}); + return runDownloader({"route-files", route}); } std::string getDevices() { - return runPython({"devices"}); + return runDownloader({"devices"}); } std::string getDeviceRoutes(const std::string &dongle_id, int64_t start_ms, int64_t end_ms, bool preserved) { @@ -178,7 +63,7 @@ std::string getDeviceRoutes(const std::string &dongle_id, int64_t start_ms, int6 args.push_back(std::to_string(end_ms)); } } - return runPython(args); + return runDownloader(args); } } // namespace PyDownloader diff --git a/openpilot/tools/replay/py_process.cc b/openpilot/tools/replay/py_process.cc new file mode 100644 index 00000000000000..b2c76083ad1915 --- /dev/null +++ b/openpilot/tools/replay/py_process.cc @@ -0,0 +1,129 @@ +#include "tools/replay/py_process.h" + +#include +#include +#include +#include + +#include "tools/replay/util.h" + +namespace PyProcess { + +std::string runModule(const std::string &module, const std::vector &args, + std::atomic *abort, bool trim) { + // Build argv for execvp + std::vector argv; + argv.push_back("python3"); + argv.push_back("-m"); + argv.push_back(module.c_str()); + for (const auto &a : args) { + argv.push_back(a.c_str()); + } + argv.push_back(nullptr); + + int stdout_pipe[2]; + if (pipe(stdout_pipe) != 0) { + rWarning("py_process: pipe() failed"); + return {}; + } + + pid_t pid = fork(); + if (pid < 0) { + rWarning("py_process: fork() failed"); + close(stdout_pipe[0]); close(stdout_pipe[1]); + return {}; + } + + if (pid == 0) { + // Child process — detach from controlling terminal so Python + // cannot corrupt terminal settings needed by ncurses in the parent. + setsid(); + int devnull = open("/dev/null", O_RDONLY); + if (devnull >= 0) { + dup2(devnull, STDIN_FILENO); + if (devnull > STDERR_FILENO) close(devnull); + } + + // Clear OPENPILOT_PREFIX so the Python process uses default paths + // (e.g. ~/.comma/auth.json). The prefix is only for IPC in the parent. + unsetenv("OPENPILOT_PREFIX"); + + close(stdout_pipe[0]); + dup2(stdout_pipe[1], STDOUT_FILENO); + close(stdout_pipe[1]); + + execvp("python3", const_cast(argv.data())); + _exit(127); + } + + // Parent process + close(stdout_pipe[1]); + + std::string stdout_data; + char buf[4096]; + + // Use select() so abort can interrupt while waiting for Python output. + fd_set rfds; + bool stdout_open = true; + + while (stdout_open) { + if (abort && *abort) { + kill(pid, SIGTERM); + break; + } + + FD_ZERO(&rfds); + FD_SET(stdout_pipe[0], &rfds); + + struct timeval tv = {0, 100000}; // 100ms timeout + int ret = select(stdout_pipe[0] + 1, &rfds, nullptr, nullptr, &tv); + if (ret < 0) break; + + if (FD_ISSET(stdout_pipe[0], &rfds)) { + ssize_t n = read(stdout_pipe[0], buf, sizeof(buf)); + if (n <= 0) { + stdout_open = false; + } else { + stdout_data.append(buf, n); + } + } + } + + // Drain remaining pipe data to prevent child from blocking on write + while (true) { + ssize_t n = read(stdout_pipe[0], buf, sizeof(buf)); + if (n <= 0) break; + stdout_data.append(buf, n); + } + close(stdout_pipe[0]); + + int status; + waitpid(pid, &status, 0); + + const bool aborted = abort && *abort; + const bool expected_sigterm = aborted && WIFSIGNALED(status) && WTERMSIG(status) == SIGTERM; + bool failed = aborted || + (WIFEXITED(status) && WEXITSTATUS(status) != 0) || + WIFSIGNALED(status); + if (failed) { + if (expected_sigterm) { + // Caller signaled abort; expected shutdown path. + } else if (WIFEXITED(status) && WEXITSTATUS(status) != 0) { + rWarning("py_process: %s exited with code %d", module.c_str(), WEXITSTATUS(status)); + } else if (WIFSIGNALED(status)) { + rWarning("py_process: %s killed by signal %d", module.c_str(), WTERMSIG(status)); + } + return {}; + } + + // Trim trailing newline + if (trim) { + while (!stdout_data.empty() && (stdout_data.back() == '\n' || stdout_data.back() == '\r')) { + stdout_data.pop_back(); + } + } + + return stdout_data; +} + +} // namespace PyProcess diff --git a/openpilot/tools/replay/py_process.h b/openpilot/tools/replay/py_process.h new file mode 100644 index 00000000000000..d8ec67c27def65 --- /dev/null +++ b/openpilot/tools/replay/py_process.h @@ -0,0 +1,14 @@ +#pragma once + +#include +#include +#include + +namespace PyProcess { + +// Run a Python command and capture stdout. Stderr is left attached to the parent. +// Returns stdout content. If abort is signaled, kills the child process. +std::string runModule(const std::string &module, const std::vector &args, + std::atomic *abort = nullptr, bool trim = true); + +} // namespace PyProcess