diff --git a/logs/view.py b/logs/view.py new file mode 100644 index 0000000..88b377f --- /dev/null +++ b/logs/view.py @@ -0,0 +1,83 @@ +# File: logs/view.py +# Copyright (C) 2026 Erick Ahmed +# SPDX-License-Identifier: AGPL-3.0-or-later + +import pandas as pd +from dash import html, dash_table, dcc +import dash_bootstrap_components as dbc + +PAGE_SIZE = 25000 + +def prepare_logs_data(df: pd.DataFrame) -> pd.DataFrame: + if df is None or df.empty: + return pd.DataFrame() + + df = df.copy() + + if 'j1939_metadata' in df.columns: + df['Priority'] = df['j1939_metadata'].apply(lambda x: x.get('Priority') if isinstance(x, dict) else None) + df['PF'] = df['j1939_metadata'].apply(lambda x: x.get('PF') if isinstance(x, dict) else None) + df['PS'] = df['j1939_metadata'].apply(lambda x: x.get('PS') if isinstance(x, dict) else None) + df['SA'] = df['j1939_metadata'].apply(lambda x: x.get('SA') if isinstance(x, dict) else None) + df['DA'] = df['j1939_metadata'].apply(lambda x: x.get('DA') if isinstance(x, dict) else None) + df['PGN'] = df['j1939_metadata'].apply(lambda x: x.get('PGN') if isinstance(x, dict) else None) + else: + for col in ['Priority', 'PF', 'PS', 'SA', 'DA', 'PGN']: + df[col] = None + + for i in range(8): + col = f'b{i}' + if col in df.columns: + df[col] = df[col].apply(lambda x: f"{int(x):02X}" if pd.notna(x) else "") + else: + df[col] = "" + + if 'ID' in df.columns: + df['ID'] = df['ID'].astype(str) + + display_cols = ['Timestamp', 'ID', 'DLC', 'b0', 'b1', 'b2', 'b3', 'b4', 'b5', 'b6', 'b7', 'Priority', 'PF', 'PS', 'SA', 'DA', 'PGN'] + display_df = df[[c for c in display_cols if c in df.columns]] + + return display_df.fillna("") + +def get_logs_table_component(): + return html.Div([ + html.Div(id='logs-info-text', className="text-muted mb-2"), + dash_table.DataTable( + id='logs-table', + virtualization=True, + page_action='none', + style_table={'overflowX': 'auto', 'height': '70vh', 'overflowY': 'auto'}, + style_header={ + 'backgroundColor': '#1a1a1a', + 'color': 'white', + 'fontWeight': 'bold', + 'textAlign': 'center', + 'position': 'sticky', + 'top': 0 + }, + style_data={ + 'backgroundColor': '#f8f9fa', + 'color': '#2a2a2a', + 'textAlign': 'center' + }, + style_data_conditional=[ + { + 'if': {'row_index': 'odd'}, + 'backgroundColor': 'rgb(240, 240, 240)' + } + ], + style_cell={ + 'minWidth': '80px', + 'padding': '5px', + 'textAlign': 'center', + 'fontFamily': 'Segoe UI, Arial, sans-serif' + } + ), + html.Div([ + dbc.Button("Prev", id='logs-prev-btn', color="secondary", outline=True, size="sm", className="me-2"), + html.Div(id='logs-page-nav', className="d-inline-block", style={'verticalAlign': 'middle'}), + dbc.Button("Next", id='logs-next-btn', color="secondary", outline=True, size="sm", className="ms-2"), + ], className="d-flex justify-content-center align-items-center mt-3"), + dcc.Store(id='logs-current-page', data=0), + ]) diff --git a/main.py b/main.py index 68cd788..8b01a27 100644 --- a/main.py +++ b/main.py @@ -8,7 +8,7 @@ from concurrent.futures import ThreadPoolExecutor import polars as pl import dash -from dash import dcc, html, Input, Output +from dash import dcc, html, Input, Output, State import dash_bootstrap_components as dbc import numpy as np @@ -19,56 +19,84 @@ from stats.id_viewer import _format_can_id_vec, plot_bits from stats.frequency import calculate_frequency, plot_frequency from stats.correlation import calculate_correlation, plot_correlation_heatmap from stats.entropy import calculate_byte_entropy, plot_entropy_heatmap +from logs.view import get_logs_table_component, prepare_logs_data -RAW_LOG = "data/logs/rawlog.txt" -BUS1_CSV = "data/csv/bus1.csv" -BUS2_CSV = "data/csv/bus2.csv" -BUS1_PARQUET = "data/parquet/bus1.parquet" -BUS2_PARQUET = "data/parquet/bus2.parquet" -BUS1_DECODED = "data/parquet/bus1_decoded.parquet" -BUS2_DECODED = "data/parquet/bus2_decoded.parquet" +RAW_LOG_DIR = "data/logs" +PAGE_SIZE = 25000 + +def parse_vehicle_from_filename(filename: str): + 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 run_pipeline(): - os.makedirs("data/logs", exist_ok=True) + os.makedirs(RAW_LOG_DIR, exist_ok=True) os.makedirs("data/csv", exist_ok=True) os.makedirs("data/parquet", exist_ok=True) - if not Path(BUS1_DECODED).exists() or not Path(BUS2_DECODED).exists(): - print("Parsing raw log...") - parse_log(RAW_LOG, BUS1_CSV, BUS2_CSV) - print("Converting to parquet...") - parse_csv(BUS1_CSV).sink_parquet(BUS1_PARQUET) - parse_csv(BUS2_CSV).sink_parquet(BUS2_PARQUET) + for log_file in Path(RAW_LOG_DIR).glob("*.txt"): + vehicle, brand, model = parse_vehicle_from_filename(log_file.name) - print("Decoding J1939...") - df1 = pl.read_parquet(BUS1_PARQUET) - df2 = pl.read_parquet(BUS2_PARQUET) - dec1 = decode_j1939_frames(df1) - dec2 = decode_j1939_frames(df2) - dec1.write_parquet(BUS1_DECODED) - dec2.write_parquet(BUS2_DECODED) + bus1_csv = f"data/csv/{vehicle}_bus1.csv" + bus2_csv = f"data/csv/{vehicle}_bus2.csv" + bus1_parquet = f"data/parquet/{vehicle}_bus1.parquet" + bus2_parquet = f"data/parquet/{vehicle}_bus2.parquet" + bus1_decoded = f"data/parquet/{vehicle}_bus1_decoded.parquet" + bus2_decoded = f"data/parquet/{vehicle}_bus2_decoded.parquet" + + if not Path(bus1_decoded).exists() or not Path(bus2_decoded).exists(): + print(f"Parsing raw log: {log_file.name}...") + parse_log(str(log_file), bus1_csv, bus2_csv) + + print("Converting to parquet...") + parse_csv(bus1_csv).sink_parquet(bus1_parquet) + parse_csv(bus2_csv).sink_parquet(bus2_parquet) + + print("Decoding J1939...") + df1 = pl.read_parquet(bus1_parquet) + df2 = pl.read_parquet(bus2_parquet) + dec1 = decode_j1939_frames(df1) + dec2 = decode_j1939_frames(df2) + dec1.write_parquet(bus1_decoded) + dec2.write_parquet(bus2_decoded) run_pipeline() print("Loading data into memory...") -DATA = { - "Bus 1": load_data(BUS1_DECODED), - "Bus 2": load_data(BUS2_DECODED) -} +DATA = {} +VEHICLE_META = {} + +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} + + bus1_decoded = f"data/parquet/{vehicle}_bus1_decoded.parquet" + bus2_decoded = f"data/parquet/{vehicle}_bus2_decoded.parquet" + + if Path(bus1_decoded).exists() and Path(bus2_decoded).exists(): + DATA[vehicle] = { + "Bus 1": load_data(bus1_decoded), + "Bus 2": load_data(bus2_decoded) + } PRECOMPUTED_FIGURES = {} DATA_BY_ID = {} CORR_CACHE = {} +PREPARED_LOGS_CACHE = {} -def process_bus_data(bus, df): +def process_bus_data(vehicle, bus, df): precomp = {} - precomp[f"{bus}_freq"] = plot_frequency(calculate_frequency(df), title=f"{bus} Frequency") - precomp[f"{bus}_entropy"] = plot_entropy_heatmap(calculate_byte_entropy(df), title=f"{bus} Byte-Level Entropy") + precomp[f"{vehicle}_{bus}_freq"] = plot_frequency(calculate_frequency(df), title=f"{vehicle} {bus} Frequency") + precomp[f"{vehicle}_{bus}_entropy"] = plot_entropy_heatmap(calculate_byte_entropy(df), title=f"{vehicle} {bus} Byte-Level Entropy") can_id_col = 'ID' if 'ID' in df.columns else 'Identifier' formatted = _format_can_id_vec(df[can_id_col]) df = df.assign(Formatted_ID=formatted) - df = df.sort_values(['Formatted_ID', 'Timestamp'], kind='stable') grouped = {} @@ -82,15 +110,17 @@ def process_bus_data(bus, df): group = group.iloc[keep] grouped[can_id] = (group, byte_cols) - return precomp, grouped + return vehicle, bus, precomp, grouped with ThreadPoolExecutor() as executor: - futures = {executor.submit(process_bus_data, bus, df): bus for bus, df in DATA.items()} + futures = [] + for vehicle, buses in DATA.items(): + for bus, df in buses.items(): + futures.append(executor.submit(process_bus_data, vehicle, bus, df)) for future in futures: - bus = futures[future] - precomp, grouped = future.result() + v, b, precomp, grouped = future.result() PRECOMPUTED_FIGURES.update(precomp) - DATA_BY_ID[bus] = grouped + DATA_BY_ID[(v, b)] = grouped app = dash.Dash(__name__, external_stylesheets=[dbc.themes.BOOTSTRAP]) app.config.suppress_callback_exceptions = True @@ -101,16 +131,42 @@ app.layout = dbc.Container([ dbc.Tab(label="Overview", tab_id="overview", children=[ html.Div(id="overview-content") ]), - dbc.Tab(label="Statistics", tab_id="statistics", children=[ + dbc.Tab(label="Logs", tab_id="logs", children=[ dbc.Row([ - dbc.Col(html.Label("Select Bus:"), width=1, className="mt-2"), + dbc.Col(html.Label("Select Vehicle:", className="mt-2"), width="auto"), dbc.Col(dcc.Dropdown( - id='bus-selector', - options=[{'label': k, 'value': k} for k in DATA.keys()], + id='logs-vehicle-selector', + options=[{'label': v, 'value': v} for v in DATA.keys()], + value=list(DATA.keys())[0] if DATA else None, + clearable=False + ), width=3, className="me-4"), + dbc.Col(html.Label("Select Bus:", className="mt-2"), width="auto"), + dbc.Col(dcc.Dropdown( + id='logs-bus-selector', + options=[{'label': 'Bus 1', 'value': 'Bus 1'}, {'label': 'Bus 2', 'value': 'Bus 2'}], value='Bus 1', clearable=False ), width=2), - ], className="mb-3 mt-3"), + ], className="mb-3 mt-3", align="end"), + get_logs_table_component() + ]), + dbc.Tab(label="Statistics", tab_id="statistics", children=[ + dbc.Row([ + dbc.Col(html.Label("Select Vehicle:", className="mt-2"), width="auto"), + dbc.Col(dcc.Dropdown( + id='vehicle-selector', + options=[{'label': v, 'value': v} for v in DATA.keys()], + value=list(DATA.keys())[0] if DATA else None, + clearable=False + ), width=3, className="me-4"), + dbc.Col(html.Label("Select Bus:", className="mt-2"), width="auto"), + dbc.Col(dcc.Dropdown( + id='bus-selector', + options=[{'label': 'Bus 1', 'value': 'Bus 1'}, {'label': 'Bus 2', 'value': 'Bus 2'}], + value='Bus 1', + clearable=False + ), 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"), @@ -122,19 +178,134 @@ app.layout = dbc.Container([ ], id="main-tabs", active_tab="statistics") ], fluid=True) +def get_prepared_logs(vehicle, bus): + cache_key = (vehicle, bus) + if cache_key not in PREPARED_LOGS_CACHE: + df = DATA[vehicle][bus] + PREPARED_LOGS_CACHE[cache_key] = prepare_logs_data(df) + return PREPARED_LOGS_CACHE[cache_key] + +def build_page_buttons(current_page: int, total_pages: int, max_buttons: int = 15) -> list: + 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": i, "id": i} for i 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, bus): - df = DATA[bus] +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"{bus}_freq"], style={'height': '80vh'}) + return dcc.Graph(figure=PRECOMPUTED_FIGURES[f"{vehicle}_{bus}_freq"], style={'height': '80vh'}) elif tab == 'id_viewer': - ids = sorted(DATA_BY_ID[bus].keys()) + ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys()) return html.Div([ html.Label("Select CAN ID:"), dcc.Dropdown( @@ -148,7 +319,7 @@ def render_content(tab, bus): ]) elif tab == 'corr': - ids = sorted(DATA_BY_ID[bus].keys()) + ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys()) return html.Div([ dbc.Row([ dbc.Col(html.Label("Method:"), width=1, className="mt-2"), @@ -170,53 +341,55 @@ def render_content(tab, bus): ]) elif tab == 'entropy': - return dcc.Graph(figure=PRECOMPUTED_FIGURES[f"{bus}_entropy"], style={'height': '80vh'}) + 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, bus, tab): - if tab != 'id_viewer' or not selected_id: +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(bus, {}) + 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"{bus} Byte Visualization") + 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, bus, tab): - if tab != 'corr': +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 = (bus, method, target_id) + cache_key = (vehicle, bus, method, target_id) if cache_key not in CORR_CACHE: - df = DATA[bus] + df = DATA[vehicle][bus] corr_df = calculate_correlation(df, method=method, target_id=target_id) CORR_CACHE[cache_key] = corr_df else: corr_df = CORR_CACHE[cache_key] - title = f"{bus} Correlation" + title = f"{vehicle} {bus} Correlation" if target_id: title += f" ({target_id})" return plot_correlation_heatmap(corr_df, target_id=target_id, title=title) if __name__ == '__main__': - app.run(debug=True) + app.run(debug=False) diff --git a/stats/correlation.py b/stats/correlation.py index 3035f39..6cbca29 100644 --- a/stats/correlation.py +++ b/stats/correlation.py @@ -83,8 +83,12 @@ def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None = return c.max(axis=0) return np.zeros(n_cols, dtype=np.float64) - with ThreadPoolExecutor() as executor: - out = np.array(list(executor.map(_process_group, groups))) + out = np.zeros((len(unique_ids), n_cols), dtype=np.float64) + if len(groups) > 0: + with ThreadPoolExecutor() as executor: + results = list(executor.map(_process_group, groups)) + for i, res in enumerate(results): + out[i] = res result = pd.DataFrame(out, index=unique_ids, columns=available_cols) result.index.name = 'Identifier' diff --git a/stats/entropy.py b/stats/entropy.py index ac3295e..e39e5d2 100644 --- a/stats/entropy.py +++ b/stats/entropy.py @@ -70,8 +70,12 @@ def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame: res[ci] = _entropy_col(sub[:, ci]) return res - with ThreadPoolExecutor() as executor: - out = np.array(list(executor.map(_process_group, groups))) + out = np.zeros((len(unique_ids), n_cols), dtype=np.float64) + if len(groups) > 0: + with ThreadPoolExecutor() as executor: + results = list(executor.map(_process_group, groups)) + for i, res in enumerate(results): + out[i] = res result = pd.DataFrame(out, index=unique_ids, columns=available_cols) result.index.name = 'Identifier'