Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 17 additions & 3 deletions openpilot/selfdrive/test/process_replay/migration.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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]]
Expand Down Expand Up @@ -387,15 +389,15 @@ 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

# 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())
Expand Down Expand Up @@ -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)
1 change: 1 addition & 0 deletions openpilot/tools/jotpluggler/app.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
5 changes: 4 additions & 1 deletion openpilot/tools/jotpluggler/app.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};

Expand Down Expand Up @@ -411,6 +412,7 @@ std::vector<RouteSeries> decode_can_messages(const std::vector<CanMessageData> &
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);
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);

Expand Down
4 changes: 2 additions & 2 deletions openpilot/tools/jotpluggler/layout.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
}

Expand Down
3 changes: 3 additions & 0 deletions openpilot/tools/jotpluggler/main.cc
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ void print_usage(const char *argv0) {
<< " --output <png>\n"
<< " --show\n"
<< " --sync-load\n"
<< " --no-migration\n"
<< "\n"
<< "Examples:\n"
<< " " << argv0 << "\n"
Expand Down Expand Up @@ -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;
Expand Down
10 changes: 6 additions & 4 deletions openpilot/tools/jotpluggler/runtime.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::mutex> 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();
}
Expand All @@ -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);
Expand Down Expand Up @@ -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<bool> active{false};
std::atomic<bool> completed{false};
Expand All @@ -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 {
Expand Down
19 changes: 17 additions & 2 deletions openpilot/tools/jotpluggler/sketch_layout.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -1551,12 +1552,17 @@ SeriesAccumulator extract_segment_series(const std::vector<Event> &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<int, SegmentLogs> &segments,
const SchemaIndex &schema,
const dbc::Database *can_dbc,
LogSelector selector,
bool skip_raw_can,
bool migrate,
LoadStats *stats) {
struct SegmentResult {
SeriesAccumulator series;
Expand Down Expand Up @@ -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<std::mutex> lock(error_mutex);
if (first_error.empty()) {
Expand Down Expand Up @@ -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{};

Expand Down Expand Up @@ -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),
Expand Down
2 changes: 1 addition & 1 deletion openpilot/tools/replay/SConscript
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading
Loading