22 Commits

Author SHA1 Message Date
eeeck f5450da96d Merge pull request 'Implement Plotly Dash app with efficient dynamic resampler' (#6) from dev-dash into main
Reviewed-on: erickahmed/CANveyor#6
2026-07-22 20:18:19 +02:00
eeeck d8ca263c0d Bump version and update project dependencies 2026-07-22 20:16:52 +02:00
eeeck 1463fa12ff Parallelize data processing tasks with ThreadPoolExecutor 2026-07-22 20:13:19 +02:00
eeeck 6c198d83c5 Remove debug flag 2026-07-22 20:06:50 +02:00
eeeck 57505074cd Change title to project name 2026-07-22 20:01:58 +02:00
eeeck af8e916116 Create an Overview menu
- To use as a sort of main menu
2026-07-22 20:00:59 +02:00
eeeck d9262e365a Put all CAN bus statistics submenus under a Statistics menu 2026-07-22 20:00:13 +02:00
eeeck 24ce8dad60 Explicitly specify grouping column in dataframe iteration 2026-07-22 19:47:11 +02:00
eeeck 22d4af292c Refactor CSV to Parquet conversion logic 2026-07-22 19:47:05 +02:00
eeeck 5e01c3bb44 Ensure data directories exist before pipeline execution 2026-07-22 19:46:59 +02:00
eeeck d6baaaa1fa Remove redundant Formatted_ID column in frequency calculation 2026-07-22 19:46:51 +02:00
eeeck 9965bc761f Simplify byte column selection in correlation calculation 2026-07-22 19:46:45 +02:00
eeeck 0a4dc5801e Merge pull request 'Implement plotly resamper and precompute data' (#5) from dev-plotly-resampler into dev-dash
Reviewed-on: erickahmed/CANveyor#5
2026-07-22 18:40:44 +02:00
eeeck c06d813c26 Remove resampling information on legend 2026-07-22 18:39:03 +02:00
eeeck 2da646fa80 Precompute CAN data
- Slower startup
- Much faster visualization (from O(n) to O(1))
2026-07-22 18:20:18 +02:00
eeeck 93e0e3f648 Suppress callback exceptions 2026-07-22 18:16:43 +02:00
eeeck 02e46ddf0b Implement plotly-resampler 2026-07-22 18:11:35 +02:00
eeeck d274897cf3 Implement lttbc 2026-07-22 18:07:33 +02:00
eeeck 65591bbc6b Refactor main application to use Polars pipeline
- replaced the caching layer with a pre-processing pipeline that parses
  raw logs into decoded Parquet files
2026-07-15 00:48:51 +02:00
eeeck c98563f541 Fix schema check and update import paths
- Use `collect_schema` for accurate column validation in Polars and
  correct
  relative import paths for statistical modules.
2026-07-15 00:48:30 +02:00
eeeck 3df42fb497 Make subdirectories Python packages 2026-07-15 00:37:19 +02:00
eeeck 73e76d3adc Rename to avoid conflict with Python stat module 2026-07-15 00:31:27 +02:00
10 changed files with 251 additions and 63 deletions
+208 -37
View File
@@ -2,50 +2,221 @@
# Copyright (C) 2026 Erick Ahmed # Copyright (C) 2026 Erick Ahmed
# SPDX-License-Identifier: AGPL-3.0-or-later # SPDX-License-Identifier: AGPL-3.0-or-later
import diskcache import os
import flask_caching from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
import polars as pl
import dash import dash
from dash import Input, Output, State, dcc, html, no_update from dash import dcc, html, Input, Output
import dash_bootstrap_components as dbc import dash_bootstrap_components as dbc
import numpy as np
CACHE_DATA_DIR = ".cache_data" from parser import parse_log, parse_csv
background_callback_manager = dash.DiskcacheManager(cache_dir=CACHE_DATA_DIR) from decoder import decode_j1939_frames
data_cache = flask_caching.Cache(config={'CACHE_TYPE': 'FileSystemCache', 'CACHE_DIR': CACHE_DATA_DIR}) 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
app = dash.Dash( RAW_LOG = "data/logs/rawlog.txt"
__name__, BUS1_CSV = "data/csv/bus1.csv"
external_stylesheets=[dbc.themes.BOOTSTRAP], BUS2_CSV = "data/csv/bus2.csv"
background_callback_manager=background_callback_manager BUS1_PARQUET = "data/parquet/bus1.parquet"
) BUS2_PARQUET = "data/parquet/bus2.parquet"
data_cache.init_app(app.server) BUS1_DECODED = "data/parquet/bus1_decoded.parquet"
BUS2_DECODED = "data/parquet/bus2_decoded.parquet"
def _get_df(session_data): def run_pipeline():
if not session_data or "token" not in session_data: os.makedirs("data/logs", exist_ok=True)
return None os.makedirs("data/csv", exist_ok=True)
return data_cache.get(session_data["token"]) 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)
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)
}
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"),
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")
], fluid=True)
@app.callback( @app.callback(
Output("graph-correlation", "figure"), Output('tab-content', 'children'),
Input("corr-method", "value"), Input('tabs', 'active_tab'),
Input("corr-target", "value"), Input('bus-selector', 'value')
Input("session-store", "data"),
background=True,
prevent_initial_call=True,
) )
def update_correlation(method, target, session_data): def render_content(tab, bus):
df = _get_df(session_data) df = DATA[bus]
if df is None or not method:
return no_update
target_id = None if target == "all" else target if tab == 'freq':
try: return dcc.Graph(figure=PRECOMPUTED_FIGURES[f"{bus}_freq"], style={'height': '80vh'})
elif tab == 'id_viewer':
ids = sorted(DATA_BY_ID[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[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"{bus}_entropy"], style={'height': '80vh'})
return html.Div("Tab not found")
@app.callback(
Output('id-viewer-graph', 'figure'),
Input('id-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:
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]
return plot_bits(filtered_df, byte_cols, selected_id, title=f"{bus} Byte Visualization")
@app.callback(
Output('corr-graph', 'figure'),
Input('corr-method', 'value'),
Input('corr-target', 'value'),
Input('bus-selector', 'value'),
Input('tabs', 'active_tab'),
)
def update_corr(method, target, bus, tab):
if tab != 'corr':
return dash.no_update
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_df = calculate_correlation(df, method=method, target_id=target_id)
title = f"Inter-Byte Correlation ({method.capitalize()})" CORR_CACHE[cache_key] = corr_df
if target_id: else:
title += f" - {target_id}" corr_df = CORR_CACHE[cache_key]
return plot_correlation_heatmap(corr_df, target_id=target_id, title=title)
except Exception as exc: title = f"{bus} Correlation"
fig = dash.go.Figure() if target_id:
fig.update_layout(title=f"Error: {exc}") title += f" ({target_id})"
return fig
return plot_correlation_heatmap(corr_df, target_id=target_id, title=title)
if __name__ == '__main__':
app.run(debug=True)
+1 -1
View File
@@ -77,7 +77,7 @@ def parse_csv(csv_path: PathLike) -> pl.LazyFrame:
""" """
lf = pl.scan_csv(csv_path, schema_overrides={"ID": pl.String, "Data": pl.String}) lf = pl.scan_csv(csv_path, schema_overrides={"ID": pl.String, "Data": pl.String})
if "Timestamp" not in lf.columns: if "Timestamp" not in lf.collect_schema().names():
lf = lf.with_row_index("Timestamp") lf = lf.with_row_index("Timestamp")
byte_exprs = [] byte_exprs = []
+11 -3
View File
@@ -1,7 +1,15 @@
[project] [project]
name = "CANveyor" name = "CANveyor"
version = "0.0.4" version = "0.1.0"
description = "J1939 CAN bus parser that works in pair with CANdigger" description = "J1939 CAN bus parser that works in pair with CANdigger"
readme = "README.md" readme = "README.md"
requires-python = ">=3.14" requires-python = ">=3.10"
dependencies = ["polars", "pathlib", "typing"] dependencies = [
"polars",
"dash",
"dash-bootstrap-components",
"numpy",
"pandas",
"plotly",
"plotly-resampler"
]
View File
+10 -7
View File
@@ -4,13 +4,14 @@
import argparse import argparse
from pathlib import Path from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
import numpy as np import numpy as np
import pandas as pd import pandas as pd
import plotly.graph_objects as go import plotly.graph_objects as go
from utils.extractor import load_data from stats.utils.extractor import load_data
from utils.extractor import to_int from stats.utils.extractor import to_int
def _format_can_id_vec(s: pd.Series) -> pd.Series: def _format_can_id_vec(s: pd.Series) -> pd.Series:
s = s.astype('string').str.strip() s = s.astype('string').str.strip()
@@ -27,8 +28,7 @@ def _ensure_int_bytes(df: pd.DataFrame, cols: list) -> pd.DataFrame:
return df return df
def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None = None) -> pd.DataFrame: def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None = None) -> pd.DataFrame:
byte_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns] available_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: if not available_cols:
raise ValueError("No byte columns (b0-b7) found in the DataFrame") raise ValueError("No byte columns (b0-b7) found in the DataFrame")
@@ -70,8 +70,7 @@ def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None =
else: else:
groups = [] groups = []
out = np.zeros((len(unique_ids), n_cols), dtype=np.float64) def _process_group(sub):
for gi, sub in enumerate(groups):
mask = ~np.isnan(sub).any(axis=1) mask = ~np.isnan(sub).any(axis=1)
sub = sub[mask] sub = sub[mask]
if len(sub) > 1: if len(sub) > 1:
@@ -81,7 +80,11 @@ def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None =
c = np.abs(np.corrcoef(sub, rowvar=False)) c = np.abs(np.corrcoef(sub, rowvar=False))
np.nan_to_num(c, copy=False, nan=0.0) np.nan_to_num(c, copy=False, nan=0.0)
np.fill_diagonal(c, 0.0) np.fill_diagonal(c, 0.0)
out[gi] = c.max(axis=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)))
result = pd.DataFrame(out, index=unique_ids, columns=available_cols) result = pd.DataFrame(out, index=unique_ids, columns=available_cols)
result.index.name = 'Identifier' result.index.name = 'Identifier'
+11 -7
View File
@@ -4,13 +4,14 @@
import argparse import argparse
from pathlib import Path from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
import numpy as np import numpy as np
import pandas as pd import pandas as pd
import plotly.graph_objects as go import plotly.graph_objects as go
from utils.extractor import load_data from stats.utils.extractor import load_data
from utils.extractor import to_int from stats.utils.extractor import to_int
def _format_can_id_vec(s: pd.Series) -> pd.Series: def _format_can_id_vec(s: pd.Series) -> pd.Series:
s = s.astype('string').str.strip() s = s.astype('string').str.strip()
@@ -36,8 +37,7 @@ def _entropy_col(a: np.ndarray) -> float:
return float(-np.sum(p * np.log2(p))) return float(-np.sum(p * np.log2(p)))
def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame: def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame:
byte_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns] available_cols = [f"b{i}" for i in range(8) if f"b{i}" in df.columns]
available_cols = byte_cols
if not available_cols: if not available_cols:
raise ValueError("No byte columns (b0-b7) found in the DataFrame") raise ValueError("No byte columns (b0-b7) found in the DataFrame")
@@ -64,10 +64,14 @@ def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame:
else: else:
groups = [] groups = []
out = np.zeros((len(unique_ids), n_cols), dtype=np.float64) def _process_group(sub):
for gi, sub in enumerate(groups): res = np.zeros(n_cols, dtype=np.float64)
for ci in range(n_cols): for ci in range(n_cols):
out[gi, ci] = _entropy_col(sub[:, ci]) res[ci] = _entropy_col(sub[:, ci])
return res
with ThreadPoolExecutor() as executor:
out = np.array(list(executor.map(_process_group, groups)))
result = pd.DataFrame(out, index=unique_ids, columns=available_cols) result = pd.DataFrame(out, index=unique_ids, columns=available_cols)
result.index.name = 'Identifier' result.index.name = 'Identifier'
+1 -2
View File
@@ -7,7 +7,7 @@ from pathlib import Path
import numpy as np import numpy as np
import pandas as pd import pandas as pd
import plotly.graph_objects as go import plotly.graph_objects as go
from utils.extractor import load_data from stats.utils.extractor import load_data
def _format_can_id_vec(s: pd.Series) -> pd.Series: def _format_can_id_vec(s: pd.Series) -> pd.Series:
s = s.astype('string').str.strip() s = s.astype('string').str.strip()
@@ -18,7 +18,6 @@ def _format_can_id_vec(s: pd.Series) -> pd.Series:
def calculate_frequency(df: pd.DataFrame) -> pd.DataFrame: def calculate_frequency(df: pd.DataFrame) -> pd.DataFrame:
can_id_col = 'ID' if 'ID' in df.columns else 'Identifier' can_id_col = 'ID' if 'ID' in df.columns else 'Identifier'
formatted = _format_can_id_vec(df[can_id_col]) formatted = _format_can_id_vec(df[can_id_col])
df['Formatted_ID'] = formatted
counts = formatted.value_counts() counts = formatted.value_counts()
freq_df = pd.DataFrame({ freq_df = pd.DataFrame({
+9 -6
View File
@@ -7,7 +7,8 @@ from pathlib import Path
import numpy as np import numpy as np
import pandas as pd import pandas as pd
import plotly.graph_objects as go import plotly.graph_objects as go
from utils.extractor import load_data from plotly_resampler import FigureResampler
from stats.utils.extractor import load_data
def _format_can_id_vec(s: pd.Series) -> pd.Series: def _format_can_id_vec(s: pd.Series) -> pd.Series:
s = s.astype('string').str.strip() s = s.astype('string').str.strip()
@@ -40,22 +41,24 @@ def prepare_data(df, target_id):
return filtered, byte_cols return filtered, byte_cols
def plot_bits(df, byte_cols, can_id, title): def plot_bits(df, byte_cols, can_id, title):
fig = go.Figure() fig = FigureResampler(
resampled_trace_prefix_suffix=("", ""),
show_mean_aggregation_size=False
)
colors = ['#e41a1c', '#377eb8', '#4daf4a', '#984ea3', '#ff7f00', '#ffff33', '#a65628', '#f781bf'] colors = ['#e41a1c', '#377eb8', '#4daf4a', '#984ea3', '#ff7f00', '#ffff33', '#a65628', '#f781bf']
n = len(byte_cols) n = len(byte_cols)
x = df['Timestamp'].to_numpy() if not df.empty else np.array([]) x = df['Timestamp'].to_numpy() if not df.empty else np.array([])
for i, col in enumerate(byte_cols): for i, col in enumerate(byte_cols):
y = df[col].to_numpy(dtype=np.float32, copy=False) if not df.empty else np.array([]) y = df[col].to_numpy(dtype=np.float32, copy=False) if not df.empty else np.array([])
fig.add_trace(go.Scattergl(
x=x, fig.add_trace(go.Scatter(
y=y,
mode='lines', mode='lines',
line=dict(shape='hv', width=2, color=colors[i % len(colors)]), line=dict(shape='hv', width=2, color=colors[i % len(colors)]),
name=col.upper(), name=col.upper(),
legendgroup=col.upper(), legendgroup=col.upper(),
hovertemplate=f"<b>{col.upper()}</b><br>Time: %{{x}}<br>Value: %{{y}}<extra></extra>", 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}]) all_button = dict(label='ALL', method='restyle', args=[{'visible': [True] * n}])
none_button = dict(label='NONE', method='restyle', args=[{'visible': ['legendonly'] * n}]) none_button = dict(label='NONE', method='restyle', args=[{'visible': ['legendonly'] * n}])
View File