396 lines
14 KiB
Python
396 lines
14 KiB
Python
# File: main.py
|
||
# Copyright (C) 2026 Erick Ahmed
|
||
# SPDX-License-Identifier: AGPL-3.0-or-later
|
||
|
||
import os
|
||
from pathlib import Path
|
||
from concurrent.futures import ThreadPoolExecutor
|
||
|
||
import polars as pl
|
||
import dash
|
||
from dash import dcc, html, Input, Output, State
|
||
import dash_bootstrap_components as dbc
|
||
import numpy as np
|
||
|
||
from parser import parse_log, parse_csv
|
||
from decoder import decode_j1939_frames
|
||
from stats.utils.extractor import load_data
|
||
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_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(RAW_LOG_DIR, exist_ok=True)
|
||
os.makedirs("data/csv", exist_ok=True)
|
||
os.makedirs("data/parquet", exist_ok=True)
|
||
|
||
for log_file in Path(RAW_LOG_DIR).glob("*.txt"):
|
||
vehicle, brand, model = parse_vehicle_from_filename(log_file.name)
|
||
|
||
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 = {}
|
||
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(vehicle, bus, df):
|
||
precomp = {}
|
||
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 = {}
|
||
for can_id, group in df.groupby(by='Formatted_ID'):
|
||
byte_cols = [f"b{i}" for i in range(8) if f"b{i}" in group.columns]
|
||
if not group.empty and len(byte_cols) > 0:
|
||
arr = group[byte_cols].to_numpy(dtype=np.float32, copy=False)
|
||
if len(arr) > 1:
|
||
changed = np.any(arr[1:] != arr[:-1], axis=1)
|
||
keep = np.concatenate(([True], changed))
|
||
group = group.iloc[keep]
|
||
grouped[can_id] = (group, byte_cols)
|
||
|
||
return vehicle, bus, precomp, grouped
|
||
|
||
with ThreadPoolExecutor() as executor:
|
||
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:
|
||
v, b, precomp, grouped = future.result()
|
||
PRECOMPUTED_FIGURES.update(precomp)
|
||
DATA_BY_ID[(v, b)] = grouped
|
||
|
||
app = dash.Dash(__name__, external_stylesheets=[dbc.themes.BOOTSTRAP])
|
||
app.config.suppress_callback_exceptions = True
|
||
|
||
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="Logs", tab_id="logs", children=[
|
||
dbc.Row([
|
||
dbc.Col(html.Label("Select Vehicle:", className="mt-2"), width="auto"),
|
||
dbc.Col(dcc.Dropdown(
|
||
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", 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"),
|
||
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, 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, 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'})
|
||
|
||
elif tab == 'id_viewer':
|
||
ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys())
|
||
return html.Div([
|
||
html.Label("Select 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'})
|
||
])
|
||
|
||
elif 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'})
|
||
])
|
||
|
||
elif 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_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"{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=False)
|