# File: main.py # Copyright (C) 2026 Erick Ahmed # SPDX-License-Identifier: AGPL-3.0-or-later """CANveyor dashboard entry point. Handles raw log ingestion, J1939 decoding, precomputation of statistical figures, and exposes a Dash application for browsing the processed data. """ import os from concurrent.futures import ThreadPoolExecutor from pathlib import Path from typing import Dict, List, Tuple import dash import dash_bootstrap_components as dbc import numpy as np import pandas as pd import polars as pl from dash import dcc, html, Input, Output, State from decoder import decode_j1939_frames from logs.view import get_logs_table_component, prepare_logs_data from parser import parse_csv, parse_log from stats.correlation import calculate_correlation, plot_correlation_heatmap from stats.entropy import calculate_byte_entropy, plot_entropy_heatmap from stats.frequency import calculate_frequency, plot_frequency from stats.id_viewer import _format_can_id_vec, plot_bits from stats.utils.loader import load_data from vehicle import get_vehicle_module RAW_LOG_DIR = "data/logs" CSV_DIR = "data/csv" PARQUET_DIR = "data/parquet" PAGE_SIZE = 25_000 BYTE_COLS = [f"b{i}" for i in range(8)] BUS_OPTIONS = [ {"label": "Bus 1", "value": "Bus 1"}, {"label": "Bus 2", "value": "Bus 2"}, ] DATA: Dict[str, Dict[str, pd.DataFrame]] = {} VEHICLE_META: Dict[str, Dict[str, str]] = {} PRECOMPUTED_FIGURES: Dict[str, object] = {} DATA_BY_ID: Dict[Tuple[str, str], Dict[str, Tuple[pd.DataFrame, List[str]]]] = {} CORR_CACHE: Dict[Tuple, object] = {} PREPARED_LOGS_CACHE: Dict[Tuple[str, str], pd.DataFrame] = {} def parse_vehicle_from_filename(filename: str) -> Tuple[str, str, str]: """Derive (vehicle, brand, model) from a log file name.""" stem = Path(filename).stem if "-" in stem: brand, model_part = stem.split("-", 1) else: brand, model_part = stem, "Unknown" model = model_part.replace("_", " ") vehicle = f"{brand} {model}".strip() return vehicle, brand, model def _vehicle_paths(vehicle: str) -> Dict[str, str]: """Return all intermediate file paths for a given vehicle.""" return { "bus1_csv": f"{CSV_DIR}/{vehicle}_bus1.csv", "bus2_csv": f"{CSV_DIR}/{vehicle}_bus2.csv", "bus1_parquet": f"{PARQUET_DIR}/{vehicle}_bus1.parquet", "bus2_parquet": f"{PARQUET_DIR}/{vehicle}_bus2.parquet", "bus1_decoded": f"{PARQUET_DIR}/{vehicle}_bus1_decoded.parquet", "bus2_decoded": f"{PARQUET_DIR}/{vehicle}_bus2_decoded.parquet", } def run_pipeline() -> None: """Parse raw log files, convert to parquet, and decode J1939 frames.""" for directory in (RAW_LOG_DIR, CSV_DIR, PARQUET_DIR): os.makedirs(directory, exist_ok=True) for log_file in Path(RAW_LOG_DIR).glob("*.txt"): vehicle, _, _ = parse_vehicle_from_filename(log_file.name) paths = _vehicle_paths(vehicle) if Path(paths["bus1_decoded"]).exists() and Path(paths["bus2_decoded"]).exists(): continue print(f"Parsing raw log: {log_file.name}...") parse_log(str(log_file), paths["bus1_csv"], paths["bus2_csv"]) print("Converting to parquet...") parse_csv(paths["bus1_csv"]).sink_parquet(paths["bus1_parquet"]) parse_csv(paths["bus2_csv"]).sink_parquet(paths["bus2_parquet"]) print("Decoding J1939...") df1 = pl.read_parquet(paths["bus1_parquet"]) df2 = pl.read_parquet(paths["bus2_parquet"]) decode_j1939_frames(df1).write_parquet(paths["bus1_decoded"]) decode_j1939_frames(df2).write_parquet(paths["bus2_decoded"]) def load_vehicle_data() -> None: """Load all decoded parquet files into the in-memory DATA store.""" for log_file in Path(RAW_LOG_DIR).glob("*.txt"): vehicle, brand, model = parse_vehicle_from_filename(log_file.name) VEHICLE_META[vehicle] = {"brand": brand, "model": model} paths = _vehicle_paths(vehicle) if Path(paths["bus1_decoded"]).exists() and Path(paths["bus2_decoded"]).exists(): DATA[vehicle] = { "Bus 1": load_data(paths["bus1_decoded"]), "Bus 2": load_data(paths["bus2_decoded"]), } def _can_id_column(df: pd.DataFrame) -> str: return "ID" if "ID" in df.columns else "Identifier" def _downsample_unchanged(group: pd.DataFrame, byte_cols: List[str]) -> pd.DataFrame: """Keep only rows where at least one byte changed vs. the previous row.""" if group.empty or not byte_cols: return group arr = group[byte_cols].to_numpy(dtype=np.float32, copy=False) if len(arr) <= 1: return group changed = np.any(arr[1:] != arr[:-1], axis=1) keep = np.concatenate(([True], changed)) return group.iloc[keep] def process_bus_data( vehicle: str, bus: str, df: pd.DataFrame ) -> Tuple[str, str, Dict[str, object], Dict[str, Tuple[pd.DataFrame, List[str]]]]: """Compute per-bus figures and ID-grouped, downsampled frames.""" precomp = { f"{vehicle}_{bus}_freq": plot_frequency( calculate_frequency(df), title=f"{vehicle} {bus} Frequency" ), f"{vehicle}_{bus}_entropy": plot_entropy_heatmap( calculate_byte_entropy(df), title=f"{vehicle} {bus} Byte-Level Entropy" ), } df = df.assign(Formatted_ID=_format_can_id_vec(df[_can_id_column(df)])) df = df.sort_values(["Formatted_ID", "Timestamp"], kind="stable") grouped: Dict[str, Tuple[pd.DataFrame, List[str]]] = {} for can_id, group in df.groupby(by="Formatted_ID"): byte_cols = [c for c in BYTE_COLS if c in group.columns] group = _downsample_unchanged(group, byte_cols) grouped[can_id] = (group, byte_cols) return vehicle, bus, precomp, grouped def precompute_all() -> None: """Run :func:`process_bus_data` across every vehicle/bus pair in parallel.""" with ThreadPoolExecutor() as executor: futures = [ executor.submit(process_bus_data, vehicle, bus, df) for vehicle, buses in DATA.items() for bus, df in buses.items() ] for future in futures: v, b, precomp, grouped = future.result() PRECOMPUTED_FIGURES.update(precomp) DATA_BY_ID[(v, b)] = grouped run_pipeline() print("Loading data into memory...") load_vehicle_data() precompute_all() app = dash.Dash(__name__, external_stylesheets=[dbc.themes.BOOTSTRAP]) app.config.suppress_callback_exceptions = True def _vehicle_dropdown(dropdown_id: str) -> dcc.Dropdown: return dcc.Dropdown( id=dropdown_id, options=[{"label": v, "value": v} for v in DATA.keys()], value=list(DATA.keys())[0] if DATA else None, clearable=False, ) def _bus_dropdown(dropdown_id: str) -> dcc.Dropdown: return dcc.Dropdown( id=dropdown_id, options=BUS_OPTIONS, value="Bus 1", clearable=False, ) def _label(text: str) -> html.Label: return html.Label(text, className="mt-2") app.layout = dbc.Container( [ html.H1("CANveyor", className="my-4"), dbc.Tabs( [ dbc.Tab( label="Overview", tab_id="overview", children=[html.Div(id="overview-content")], ), dbc.Tab( label="Vehicles", tab_id="vehicles", children=[ dbc.Row( [ dbc.Col(_label("Vehicle:"), width="auto"), dbc.Col( _vehicle_dropdown("vehicles-vehicle-selector"), width=3, className="me-4", ), ], className="mb-3 mt-3", align="end", ), html.Div(id="vehicles-content", className="mt-3"), ], ), dbc.Tab( label="Logs", tab_id="logs", children=[ dbc.Row( [ dbc.Col(_label("Vehicle:"), width="auto"), dbc.Col( _vehicle_dropdown("logs-vehicle-selector"), width=3, className="me-4", ), dbc.Col(_label("Bus:"), width="auto"), dbc.Col( _bus_dropdown("logs-bus-selector"), width=2, ), ], className="mb-3 mt-3", align="end", ), get_logs_table_component(), ], ), dbc.Tab( label="Statistics", tab_id="statistics", children=[ dbc.Row( [ dbc.Col(_label("Vehicle:"), width="auto"), dbc.Col( _vehicle_dropdown("vehicle-selector"), width=3, className="me-4", ), dbc.Col(_label("Bus:"), width="auto"), dbc.Col( _bus_dropdown("bus-selector"), width=2, ), ], className="mb-3 mt-3", align="end", ), dbc.Tabs( [ dbc.Tab(label="Frequency", tab_id="freq"), dbc.Tab(label="ID Viewer", tab_id="id_viewer"), dbc.Tab(label="Correlation", tab_id="corr"), dbc.Tab(label="Entropy", tab_id="entropy"), ], id="tabs", active_tab="freq", ), html.Div(id="tab-content", className="mt-3"), ], ), ], id="main-tabs", active_tab="statistics", ), ], fluid=True, ) def get_prepared_logs(vehicle: str, bus: str) -> pd.DataFrame: """Lazily prepare and cache log table data for a vehicle/bus pair.""" cache_key = (vehicle, bus) if cache_key not in PREPARED_LOGS_CACHE: PREPARED_LOGS_CACHE[cache_key] = prepare_logs_data(DATA[vehicle][bus]) return PREPARED_LOGS_CACHE[cache_key] def build_page_buttons( current_page: int, total_pages: int, max_buttons: int = 15 ) -> List: """Build the pagination button list with ellipses where appropriate.""" buttons: List = [] if total_pages <= 1: return buttons half = max_buttons // 2 start = max(0, current_page - half) end = min(total_pages, start + max_buttons) if end - start < max_buttons: start = max(0, end - max_buttons) if start > 0: buttons.append( dbc.Button( "1", id={"type": "page-btn", "index": 0}, color="secondary", outline=True, size="sm", className="me-1", ) ) if start > 1: buttons.append(html.Span("…", className="mx-1 align-middle")) for i in range(start, end): is_current = i == current_page buttons.append( dbc.Button( str(i + 1), id={"type": "page-btn", "index": i}, size="sm", color="primary" if is_current else "secondary", outline=not is_current, className="me-1", disabled=is_current, ) ) if end < total_pages: if end < total_pages - 1: buttons.append(html.Span("…", className="mx-1 align-middle")) buttons.append( dbc.Button( str(total_pages), id={"type": "page-btn", "index": total_pages - 1}, color="secondary", outline=True, size="sm", className="me-1", ) ) return buttons @app.callback( Output("logs-table", "data"), Output("logs-table", "columns"), Output("logs-info-text", "children"), Output("logs-page-nav", "children"), Output("logs-current-page", "data"), Input("logs-vehicle-selector", "value"), Input("logs-bus-selector", "value"), Input("logs-prev-btn", "n_clicks"), Input("logs-next-btn", "n_clicks"), Input({"type": "page-btn", "index": dash.ALL}, "n_clicks"), State("logs-current-page", "data"), ) def update_logs_table(vehicle, bus, prev_clicks, next_clicks, page_btn_clicks, current_page): if not vehicle or not bus or vehicle not in DATA or bus not in DATA[vehicle]: return [], [], "No data available", [], 0 prepared_df = get_prepared_logs(vehicle, bus) total_rows = len(prepared_df) if total_rows == 0: return [], [], "No data available", [], 0 total_pages = max(1, (total_rows + PAGE_SIZE - 1) // PAGE_SIZE) ctx = dash.callback_context triggered_id = ctx.triggered_id current_page = current_page if current_page is not None else 0 if triggered_id in ("logs-vehicle-selector", "logs-bus-selector"): current_page = 0 elif triggered_id == "logs-prev-btn": current_page = max(0, current_page - 1) elif triggered_id == "logs-next-btn": current_page = current_page + 1 elif isinstance(triggered_id, dict) and triggered_id.get("type") == "page-btn": if ctx.triggered and ctx.triggered[0]["value"]: current_page = triggered_id["index"] current_page = max(0, min(current_page, total_pages - 1)) start_idx = current_page * PAGE_SIZE end_idx = min(start_idx + PAGE_SIZE, total_rows) page_data = prepared_df.iloc[start_idx:end_idx].to_dict("records") columns = [{"name": c, "id": c} for c in prepared_df.columns] info_text = ( f"Page {current_page + 1} of {total_pages} | " f"Showing rows {start_idx + 1:,}–{end_idx:,} " f"of {total_rows:,} total frames" ) page_buttons = build_page_buttons(current_page, total_pages) return page_data, columns, info_text, page_buttons, current_page @app.callback( Output("tab-content", "children"), Input("tabs", "active_tab"), Input("vehicle-selector", "value"), Input("bus-selector", "value"), ) def render_content(tab, vehicle, bus): if not vehicle or not bus or vehicle not in DATA or bus not in DATA[vehicle]: return html.Div("No data available") df = DATA[vehicle][bus] if tab == "freq": return dcc.Graph( figure=PRECOMPUTED_FIGURES[f"{vehicle}_{bus}_freq"], style={"height": "80vh"}, ) if tab == "id_viewer": ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys()) return html.Div( [ html.Label("CAN ID:"), dcc.Dropdown( id="id-selector", options=[{"label": i, "value": i} for i in ids], value=ids[0] if ids else None, clearable=False, style={"width": "50%", "marginBottom": "10px"}, ), dcc.Graph(id="id-viewer-graph", style={"height": "70vh"}), ] ) if tab == "corr": ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys()) return html.Div( [ dbc.Row( [ dbc.Col(html.Label("Method:"), width=1, className="mt-2"), dbc.Col( dcc.Dropdown( id="corr-method", options=[ {"label": "Pearson", "value": "pearson"}, {"label": "Spearman", "value": "spearman"}, ], value="pearson", clearable=False, ), width=2, ), dbc.Col(html.Label("Target ID:"), width=1, className="mt-2"), dbc.Col( dcc.Dropdown( id="corr-target", options=[{"label": "All IDs (Max Corr)", "value": "all"}] + [{"label": i, "value": i} for i in ids], value="all", clearable=True, ), width=4, ), ], className="mb-3", ), dcc.Graph(id="corr-graph", style={"height": "80vh"}), ] ) if tab == "entropy": return dcc.Graph( figure=PRECOMPUTED_FIGURES[f"{vehicle}_{bus}_entropy"], style={"height": "80vh"}, ) return html.Div("Tab not found") @app.callback( Output("id-viewer-graph", "figure"), Input("id-selector", "value"), Input("vehicle-selector", "value"), Input("bus-selector", "value"), Input("tabs", "active_tab"), ) def update_id_viewer(selected_id, vehicle, bus, tab): if tab != "id_viewer" or not selected_id or not vehicle or not bus: return dash.no_update grouped_data = DATA_BY_ID.get((vehicle, bus), {}) if selected_id not in grouped_data: return dash.no_update filtered_df, byte_cols = grouped_data[selected_id] return plot_bits( filtered_df, byte_cols, selected_id, title=f"{vehicle} {bus} Byte Visualization", ) @app.callback( Output("corr-graph", "figure"), Input("corr-method", "value"), Input("corr-target", "value"), Input("vehicle-selector", "value"), Input("bus-selector", "value"), Input("tabs", "active_tab"), ) def update_corr(method, target, vehicle, bus, tab): if tab != "corr" or not vehicle or not bus: return dash.no_update target_id = None if target == "all" or not target else target cache_key = (vehicle, bus, method, target_id) if cache_key not in CORR_CACHE: df = DATA[vehicle][bus] CORR_CACHE[cache_key] = calculate_correlation( df, method=method, target_id=target_id ) corr_df = CORR_CACHE[cache_key] title = f"{vehicle} {bus} Correlation" if target_id: title += f" ({target_id})" return plot_correlation_heatmap(corr_df, target_id=target_id, title=title) @app.callback( Output("vehicles-content", "children"), Input("vehicles-vehicle-selector", "value"), ) def render_vehicles(vehicle): if not vehicle or vehicle not in DATA: return html.Div("No data available", className="text-muted") brand = VEHICLE_META.get(vehicle, {}).get("brand", "") vehicle_module = get_vehicle_module(brand) dfs = list(DATA[vehicle].values()) if not dfs: return html.Div("No data available", className="text-muted") df = pd.concat(dfs, ignore_index=True) if "Timestamp" in df.columns: df = df.sort_values("Timestamp", kind="stable").reset_index(drop=True) cards = [] for nid, frame_def in vehicle_module.DECODER_RULES.items(): decoded = vehicle_module.decode_dataframe(df, frame_def.can_id) for item in frame_def.signals: if hasattr(item, "plot_func") and callable(item.plot_func): title = f"{frame_def.can_id} - {item.name}" fig = item.plot_func(decoded, frame_def.color) else: sig = item if getattr(sig, "skip_plot", False): continue unit_str = f" ({sig.unit})" if sig.unit else "" title = f"{sig.name}{unit_str}" fig = vehicle_module.plot_signal( decoded, sig.name, title=title, color=frame_def.color ) cards.append( dbc.Col( dbc.Card( [ dbc.CardBody( [ dcc.Graph( figure=fig, config={"displayModeBar": False}, style={"height": "280px"}, ) ], className="p-2", ), ], className="shadow-sm border-0 h-100", ), xs=12, sm=6, md=4, lg=3, className="mb-3", ) ) if not cards: return html.Div( "No decoded signals available. Add rules in the vehicle module.", className="text-muted", ) return dbc.Row(cards) if __name__ == "__main__": app.run(debug=False)