2 Commits

Author SHA1 Message Date
eeeck e35c56c67a Suppress callback references 2026-07-22 17:53:35 +02:00
eeeck 4bd0c60d2d Implement plotly-resampler 2026-07-22 17:41:22 +02:00
6 changed files with 60 additions and 122 deletions
+34 -87
View File
@@ -4,19 +4,18 @@
import os
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
import polars as pl
import dash
from dash import dcc, html, Input, Output
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.id_viewer import prepare_data, plot_bits
from stats.correlation import calculate_correlation, plot_correlation_heatmap
from stats.entropy import calculate_byte_entropy, plot_entropy_heatmap
@@ -30,15 +29,15 @@ BUS2_DECODED = "data/parquet/bus2_decoded.parquet"
def run_pipeline():
os.makedirs("data/logs", 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)
lf1 = parse_csv(BUS1_CSV)
lf1.sink_parquet(BUS1_PARQUET)
lf2 = parse_csv(BUS2_CSV)
lf2.sink_parquet(BUS2_PARQUET)
print("Decoding J1939...")
df1 = pl.read_parquet(BUS1_PARQUET)
@@ -56,70 +55,27 @@ DATA = {
"Bus 2": load_data(BUS2_DECODED)
}
PRECOMPUTED_FIGURES = {}
DATA_BY_ID = {}
CORR_CACHE = {}
def process_bus_data(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")
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 precomp, grouped
with ThreadPoolExecutor() as executor:
futures = {executor.submit(process_bus_data, bus, df): bus for bus, df in DATA.items()}
for future in futures:
bus = futures[future]
precomp, grouped = future.result()
PRECOMPUTED_FIGURES.update(precomp)
DATA_BY_ID[bus] = 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"),
html.H1("CAN Bus Analyzer", className="my-4"),
dbc.Row([
dbc.Col(html.Label("Select Bus:"), width=1, className="mt-2"),
dbc.Col(dcc.Dropdown(
id='bus-selector',
options=[{'label': k, 'value': k} for k in DATA.keys()],
value='Bus 1',
clearable=False
), width=2),
], className="mb-3"),
dbc.Tabs([
dbc.Tab(label="Overview", tab_id="overview", children=[
html.Div(id="overview-content")
]),
dbc.Tab(label="Statistics", tab_id="statistics", children=[
dbc.Row([
dbc.Col(html.Label("Select Bus:"), width=1, className="mt-2"),
dbc.Col(dcc.Dropdown(
id='bus-selector',
options=[{'label': k, 'value': k} for k in DATA.keys()],
value='Bus 1',
clearable=False
), width=2),
], className="mb-3 mt-3"),
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")
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")
], fluid=True)
@app.callback(
@@ -131,10 +87,12 @@ def render_content(tab, bus):
df = DATA[bus]
if tab == 'freq':
return dcc.Graph(figure=PRECOMPUTED_FIGURES[f"{bus}_freq"], style={'height': '80vh'})
stats = calculate_frequency(df)
fig = plot_frequency(stats, title=f"{bus} Frequency")
return dcc.Graph(figure=fig, style={'height': '80vh'})
elif tab == 'id_viewer':
ids = sorted(DATA_BY_ID[bus].keys())
ids = sorted(df['ID'].unique().tolist())
return html.Div([
html.Label("Select CAN ID:"),
dcc.Dropdown(
@@ -148,7 +106,7 @@ def render_content(tab, bus):
])
elif tab == 'corr':
ids = sorted(DATA_BY_ID[bus].keys())
ids = sorted(df['ID'].unique().tolist())
return html.Div([
dbc.Row([
dbc.Col(html.Label("Method:"), width=1, className="mt-2"),
@@ -170,7 +128,9 @@ def render_content(tab, bus):
])
elif tab == 'entropy':
return dcc.Graph(figure=PRECOMPUTED_FIGURES[f"{bus}_entropy"], style={'height': '80vh'})
entropy_df = calculate_byte_entropy(df)
fig = plot_entropy_heatmap(entropy_df, title=f"{bus} Byte-Level Entropy")
return dcc.Graph(figure=fig, style={'height': '80vh'})
return html.Div("Tab not found")
@@ -183,12 +143,8 @@ def render_content(tab, bus):
def update_id_viewer(selected_id, bus, tab):
if tab != 'id_viewer' or not selected_id:
return dash.no_update
grouped_data = DATA_BY_ID.get(bus, {})
if selected_id not in grouped_data:
return dash.no_update
filtered_df, byte_cols = grouped_data[selected_id]
df = DATA[bus]
filtered_df, byte_cols = prepare_data(df, selected_id)
return plot_bits(filtered_df, byte_cols, selected_id, title=f"{bus} Byte Visualization")
@app.callback(
@@ -201,21 +157,12 @@ def update_id_viewer(selected_id, bus, tab):
def update_corr(method, target, bus, tab):
if tab != 'corr':
return dash.no_update
df = DATA[bus]
target_id = None if target == 'all' or not target else target
cache_key = (bus, method, target_id)
if cache_key not in CORR_CACHE:
df = DATA[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]
corr_df = calculate_correlation(df, method=method, target_id=target_id)
title = f"{bus} Correlation"
if target_id:
title += f" ({target_id})"
return plot_correlation_heatmap(corr_df, target_id=target_id, title=title)
if __name__ == '__main__':
+3 -11
View File
@@ -1,15 +1,7 @@
[project]
name = "CANveyor"
version = "0.1.0"
version = "0.0.4"
description = "J1939 CAN bus parser that works in pair with CANdigger"
readme = "README.md"
requires-python = ">=3.10"
dependencies = [
"polars",
"dash",
"dash-bootstrap-components",
"numpy",
"pandas",
"plotly",
"plotly-resampler"
]
requires-python = ">=3.14"
dependencies = ["polars", "pathlib", "typing"]
+5 -8
View File
@@ -4,7 +4,6 @@
import argparse
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
import numpy as np
import pandas as pd
@@ -28,7 +27,8 @@ def _ensure_int_bytes(df: pd.DataFrame, cols: list) -> pd.DataFrame:
return df
def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None = None) -> pd.DataFrame:
available_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns]
byte_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns]
available_cols = [col for col in byte_cols if col in df.columns]
if not available_cols:
raise ValueError("No byte columns (b0-b7) found in the DataFrame")
@@ -70,7 +70,8 @@ def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None =
else:
groups = []
def _process_group(sub):
out = np.zeros((len(unique_ids), n_cols), dtype=np.float64)
for gi, sub in enumerate(groups):
mask = ~np.isnan(sub).any(axis=1)
sub = sub[mask]
if len(sub) > 1:
@@ -80,11 +81,7 @@ def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None =
c = np.abs(np.corrcoef(sub, rowvar=False))
np.nan_to_num(c, copy=False, nan=0.0)
np.fill_diagonal(c, 0.0)
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[gi] = c.max(axis=0)
result = pd.DataFrame(out, index=unique_ids, columns=available_cols)
result.index.name = 'Identifier'
+5 -9
View File
@@ -4,7 +4,6 @@
import argparse
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
import numpy as np
import pandas as pd
@@ -37,7 +36,8 @@ def _entropy_col(a: np.ndarray) -> float:
return float(-np.sum(p * np.log2(p)))
def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame:
available_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns]
byte_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns]
available_cols = byte_cols
if not available_cols:
raise ValueError("No byte columns (b0-b7) found in the DataFrame")
@@ -64,14 +64,10 @@ def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame:
else:
groups = []
def _process_group(sub):
res = np.zeros(n_cols, dtype=np.float64)
out = np.zeros((len(unique_ids), n_cols), dtype=np.float64)
for gi, sub in enumerate(groups):
for ci in range(n_cols):
res[ci] = _entropy_col(sub[:, ci])
return res
with ThreadPoolExecutor() as executor:
out = np.array(list(executor.map(_process_group, groups)))
out[gi, ci] = _entropy_col(sub[:, ci])
result = pd.DataFrame(out, index=unique_ids, columns=available_cols)
result.index.name = 'Identifier'
+1
View File
@@ -18,6 +18,7 @@ def _format_can_id_vec(s: pd.Series) -> pd.Series:
def calculate_frequency(df: pd.DataFrame) -> pd.DataFrame:
can_id_col = 'ID' if 'ID' in df.columns else 'Identifier'
formatted = _format_can_id_vec(df[can_id_col])
df['Formatted_ID'] = formatted
counts = formatted.value_counts()
freq_df = pd.DataFrame({
+12 -7
View File
@@ -7,7 +7,7 @@ from pathlib import Path
import numpy as np
import pandas as pd
import plotly.graph_objects as go
from plotly_resampler import FigureResampler
import lttbc
from stats.utils.extractor import load_data
def _format_can_id_vec(s: pd.Series) -> pd.Series:
@@ -41,24 +41,29 @@ def prepare_data(df, target_id):
return filtered, byte_cols
def plot_bits(df, byte_cols, can_id, title):
fig = FigureResampler(
resampled_trace_prefix_suffix=("", ""),
show_mean_aggregation_size=False
)
fig = go.Figure()
colors = ['#e41a1c', '#377eb8', '#4daf4a', '#984ea3', '#ff7f00', '#ffff33', '#a65628', '#f781bf']
n = len(byte_cols)
max_points = 2000
x = df['Timestamp'].to_numpy() if not df.empty else np.array([])
for i, col in enumerate(byte_cols):
y = df[col].to_numpy(dtype=np.float32, copy=False) if not df.empty else np.array([])
fig.add_trace(go.Scatter(
if len(x) > max_points and len(x) == len(y):
x_plot, y_plot = lttbc.downsample(x, y, max_points)
else:
x_plot, y_plot = x, y
fig.add_trace(go.Scattergl(
x=x_plot,
y=y_plot,
mode='lines',
line=dict(shape='hv', width=2, color=colors[i % len(colors)]),
name=col.upper(),
legendgroup=col.upper(),
hovertemplate=f"<b>{col.upper()}</b><br>Time: %{{x}}<br>Value: %{{y}}<extra></extra>",
), hf_x=x, hf_y=y)
))
all_button = dict(label='ALL', method='restyle', args=[{'visible': [True] * n}])
none_button = dict(label='NONE', method='restyle', args=[{'visible': ['legendonly'] * n}])