Add multi-vehicle support to dashboard and pipeline
This commit is contained in:
@@ -20,55 +20,80 @@ from stats.frequency import calculate_frequency, plot_frequency
|
|||||||
from stats.correlation import calculate_correlation, plot_correlation_heatmap
|
from stats.correlation import calculate_correlation, plot_correlation_heatmap
|
||||||
from stats.entropy import calculate_byte_entropy, plot_entropy_heatmap
|
from stats.entropy import calculate_byte_entropy, plot_entropy_heatmap
|
||||||
|
|
||||||
RAW_LOG = "data/logs/rawlog.txt"
|
RAW_LOG_DIR = "data/logs"
|
||||||
BUS1_CSV = "data/csv/bus1.csv"
|
|
||||||
BUS2_CSV = "data/csv/bus2.csv"
|
def parse_vehicle_from_filename(filename: str):
|
||||||
BUS1_PARQUET = "data/parquet/bus1.parquet"
|
stem = Path(filename).stem
|
||||||
BUS2_PARQUET = "data/parquet/bus2.parquet"
|
if '-' in stem:
|
||||||
BUS1_DECODED = "data/parquet/bus1_decoded.parquet"
|
brand, model_part = stem.split('-', 1)
|
||||||
BUS2_DECODED = "data/parquet/bus2_decoded.parquet"
|
else:
|
||||||
|
brand, model_part = stem, "Unknown"
|
||||||
|
model = model_part.replace('_', ' ')
|
||||||
|
vehicle = f"{brand} {model}".strip()
|
||||||
|
return vehicle, brand, model
|
||||||
|
|
||||||
def run_pipeline():
|
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/csv", exist_ok=True)
|
||||||
os.makedirs("data/parquet", 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...")
|
for log_file in Path(RAW_LOG_DIR).glob("*.txt"):
|
||||||
parse_log(RAW_LOG, BUS1_CSV, BUS2_CSV)
|
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...")
|
print("Converting to parquet...")
|
||||||
parse_csv(BUS1_CSV).sink_parquet(BUS1_PARQUET)
|
parse_csv(bus1_csv).sink_parquet(bus1_parquet)
|
||||||
parse_csv(BUS2_CSV).sink_parquet(BUS2_PARQUET)
|
parse_csv(bus2_csv).sink_parquet(bus2_parquet)
|
||||||
|
|
||||||
print("Decoding J1939...")
|
print("Decoding J1939...")
|
||||||
df1 = pl.read_parquet(BUS1_PARQUET)
|
df1 = pl.read_parquet(bus1_parquet)
|
||||||
df2 = pl.read_parquet(BUS2_PARQUET)
|
df2 = pl.read_parquet(bus2_parquet)
|
||||||
dec1 = decode_j1939_frames(df1)
|
dec1 = decode_j1939_frames(df1)
|
||||||
dec2 = decode_j1939_frames(df2)
|
dec2 = decode_j1939_frames(df2)
|
||||||
dec1.write_parquet(BUS1_DECODED)
|
dec1.write_parquet(bus1_decoded)
|
||||||
dec2.write_parquet(BUS2_DECODED)
|
dec2.write_parquet(bus2_decoded)
|
||||||
|
|
||||||
run_pipeline()
|
run_pipeline()
|
||||||
|
|
||||||
print("Loading data into memory...")
|
print("Loading data into memory...")
|
||||||
DATA = {
|
DATA = {}
|
||||||
"Bus 1": load_data(BUS1_DECODED),
|
VEHICLE_META = {}
|
||||||
"Bus 2": load_data(BUS2_DECODED)
|
|
||||||
|
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 = {}
|
PRECOMPUTED_FIGURES = {}
|
||||||
DATA_BY_ID = {}
|
DATA_BY_ID = {}
|
||||||
CORR_CACHE = {}
|
CORR_CACHE = {}
|
||||||
|
|
||||||
def process_bus_data(bus, df):
|
def process_bus_data(vehicle, bus, df):
|
||||||
precomp = {}
|
precomp = {}
|
||||||
precomp[f"{bus}_freq"] = plot_frequency(calculate_frequency(df), title=f"{bus} Frequency")
|
precomp[f"{vehicle}_{bus}_freq"] = plot_frequency(calculate_frequency(df), title=f"{vehicle} {bus} Frequency")
|
||||||
precomp[f"{bus}_entropy"] = plot_entropy_heatmap(calculate_byte_entropy(df), title=f"{bus} Byte-Level Entropy")
|
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'
|
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 = df.assign(Formatted_ID=formatted)
|
df = df.assign(Formatted_ID=formatted)
|
||||||
|
|
||||||
df = df.sort_values(['Formatted_ID', 'Timestamp'], kind='stable')
|
df = df.sort_values(['Formatted_ID', 'Timestamp'], kind='stable')
|
||||||
|
|
||||||
grouped = {}
|
grouped = {}
|
||||||
@@ -82,15 +107,17 @@ def process_bus_data(bus, df):
|
|||||||
group = group.iloc[keep]
|
group = group.iloc[keep]
|
||||||
grouped[can_id] = (group, byte_cols)
|
grouped[can_id] = (group, byte_cols)
|
||||||
|
|
||||||
return precomp, grouped
|
return vehicle, bus, precomp, grouped
|
||||||
|
|
||||||
with ThreadPoolExecutor() as executor:
|
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:
|
for future in futures:
|
||||||
bus = futures[future]
|
v, b, precomp, grouped = future.result()
|
||||||
precomp, grouped = future.result()
|
|
||||||
PRECOMPUTED_FIGURES.update(precomp)
|
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 = dash.Dash(__name__, external_stylesheets=[dbc.themes.BOOTSTRAP])
|
||||||
app.config.suppress_callback_exceptions = True
|
app.config.suppress_callback_exceptions = True
|
||||||
@@ -103,14 +130,21 @@ app.layout = dbc.Container([
|
|||||||
]),
|
]),
|
||||||
dbc.Tab(label="Statistics", tab_id="statistics", children=[
|
dbc.Tab(label="Statistics", tab_id="statistics", children=[
|
||||||
dbc.Row([
|
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='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(
|
dbc.Col(dcc.Dropdown(
|
||||||
id='bus-selector',
|
id='bus-selector',
|
||||||
options=[{'label': k, 'value': k} for k in DATA.keys()],
|
options=[{'label': 'Bus 1', 'value': 'Bus 1'}, {'label': 'Bus 2', 'value': 'Bus 2'}],
|
||||||
value='Bus 1',
|
value='Bus 1',
|
||||||
clearable=False
|
clearable=False
|
||||||
), width=2),
|
), width=2),
|
||||||
], className="mb-3 mt-3"),
|
], className="mb-3 mt-3", align="end"),
|
||||||
dbc.Tabs([
|
dbc.Tabs([
|
||||||
dbc.Tab(label="Frequency", tab_id="freq"),
|
dbc.Tab(label="Frequency", tab_id="freq"),
|
||||||
dbc.Tab(label="ID Viewer", tab_id="id_viewer"),
|
dbc.Tab(label="ID Viewer", tab_id="id_viewer"),
|
||||||
@@ -125,16 +159,20 @@ app.layout = dbc.Container([
|
|||||||
@app.callback(
|
@app.callback(
|
||||||
Output('tab-content', 'children'),
|
Output('tab-content', 'children'),
|
||||||
Input('tabs', 'active_tab'),
|
Input('tabs', 'active_tab'),
|
||||||
|
Input('vehicle-selector', 'value'),
|
||||||
Input('bus-selector', 'value')
|
Input('bus-selector', 'value')
|
||||||
)
|
)
|
||||||
def render_content(tab, bus):
|
def render_content(tab, vehicle, bus):
|
||||||
df = DATA[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':
|
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':
|
elif tab == 'id_viewer':
|
||||||
ids = sorted(DATA_BY_ID[bus].keys())
|
ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys())
|
||||||
return html.Div([
|
return html.Div([
|
||||||
html.Label("Select CAN ID:"),
|
html.Label("Select CAN ID:"),
|
||||||
dcc.Dropdown(
|
dcc.Dropdown(
|
||||||
@@ -148,7 +186,7 @@ def render_content(tab, bus):
|
|||||||
])
|
])
|
||||||
|
|
||||||
elif tab == 'corr':
|
elif tab == 'corr':
|
||||||
ids = sorted(DATA_BY_ID[bus].keys())
|
ids = sorted(DATA_BY_ID.get((vehicle, bus), {}).keys())
|
||||||
return html.Div([
|
return html.Div([
|
||||||
dbc.Row([
|
dbc.Row([
|
||||||
dbc.Col(html.Label("Method:"), width=1, className="mt-2"),
|
dbc.Col(html.Label("Method:"), width=1, className="mt-2"),
|
||||||
@@ -170,53 +208,55 @@ def render_content(tab, bus):
|
|||||||
])
|
])
|
||||||
|
|
||||||
elif tab == 'entropy':
|
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")
|
return html.Div("Tab not found")
|
||||||
|
|
||||||
@app.callback(
|
@app.callback(
|
||||||
Output('id-viewer-graph', 'figure'),
|
Output('id-viewer-graph', 'figure'),
|
||||||
Input('id-selector', 'value'),
|
Input('id-selector', 'value'),
|
||||||
|
Input('vehicle-selector', 'value'),
|
||||||
Input('bus-selector', 'value'),
|
Input('bus-selector', 'value'),
|
||||||
Input('tabs', 'active_tab'),
|
Input('tabs', 'active_tab'),
|
||||||
)
|
)
|
||||||
def update_id_viewer(selected_id, bus, tab):
|
def update_id_viewer(selected_id, vehicle, bus, tab):
|
||||||
if tab != 'id_viewer' or not selected_id:
|
if tab != 'id_viewer' or not selected_id or not vehicle or not bus:
|
||||||
return dash.no_update
|
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:
|
if selected_id not in grouped_data:
|
||||||
return dash.no_update
|
return dash.no_update
|
||||||
|
|
||||||
filtered_df, byte_cols = grouped_data[selected_id]
|
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(
|
@app.callback(
|
||||||
Output('corr-graph', 'figure'),
|
Output('corr-graph', 'figure'),
|
||||||
Input('corr-method', 'value'),
|
Input('corr-method', 'value'),
|
||||||
Input('corr-target', 'value'),
|
Input('corr-target', 'value'),
|
||||||
|
Input('vehicle-selector', 'value'),
|
||||||
Input('bus-selector', 'value'),
|
Input('bus-selector', 'value'),
|
||||||
Input('tabs', 'active_tab'),
|
Input('tabs', 'active_tab'),
|
||||||
)
|
)
|
||||||
def update_corr(method, target, bus, tab):
|
def update_corr(method, target, vehicle, bus, tab):
|
||||||
if tab != 'corr':
|
if tab != 'corr' or not vehicle or not bus:
|
||||||
return dash.no_update
|
return dash.no_update
|
||||||
|
|
||||||
target_id = None if target == 'all' or not target else target
|
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:
|
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_df = calculate_correlation(df, method=method, target_id=target_id)
|
||||||
CORR_CACHE[cache_key] = corr_df
|
CORR_CACHE[cache_key] = corr_df
|
||||||
else:
|
else:
|
||||||
corr_df = CORR_CACHE[cache_key]
|
corr_df = CORR_CACHE[cache_key]
|
||||||
|
|
||||||
title = f"{bus} Correlation"
|
title = f"{vehicle} {bus} Correlation"
|
||||||
if target_id:
|
if target_id:
|
||||||
title += f" ({target_id})"
|
title += f" ({target_id})"
|
||||||
|
|
||||||
return plot_correlation_heatmap(corr_df, target_id=target_id, title=title)
|
return plot_correlation_heatmap(corr_df, target_id=target_id, title=title)
|
||||||
|
|
||||||
if __name__ == '__main__':
|
if __name__ == '__main__':
|
||||||
app.run(debug=True)
|
app.run(debug=False)
|
||||||
|
|||||||
@@ -83,8 +83,12 @@ def calculate_correlation(df: pd.DataFrame, method: str, target_id: str | None =
|
|||||||
return c.max(axis=0)
|
return c.max(axis=0)
|
||||||
return np.zeros(n_cols, dtype=np.float64)
|
return np.zeros(n_cols, dtype=np.float64)
|
||||||
|
|
||||||
|
out = np.zeros((len(unique_ids), n_cols), dtype=np.float64)
|
||||||
|
if len(groups) > 0:
|
||||||
with ThreadPoolExecutor() as executor:
|
with ThreadPoolExecutor() as executor:
|
||||||
out = np.array(list(executor.map(_process_group, groups)))
|
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 = pd.DataFrame(out, index=unique_ids, columns=available_cols)
|
||||||
result.index.name = 'Identifier'
|
result.index.name = 'Identifier'
|
||||||
|
|||||||
+5
-1
@@ -70,8 +70,12 @@ def calculate_byte_entropy(df: pd.DataFrame) -> pd.DataFrame:
|
|||||||
res[ci] = _entropy_col(sub[:, ci])
|
res[ci] = _entropy_col(sub[:, ci])
|
||||||
return res
|
return res
|
||||||
|
|
||||||
|
out = np.zeros((len(unique_ids), n_cols), dtype=np.float64)
|
||||||
|
if len(groups) > 0:
|
||||||
with ThreadPoolExecutor() as executor:
|
with ThreadPoolExecutor() as executor:
|
||||||
out = np.array(list(executor.map(_process_group, groups)))
|
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 = pd.DataFrame(out, index=unique_ids, columns=available_cols)
|
||||||
result.index.name = 'Identifier'
|
result.index.name = 'Identifier'
|
||||||
|
|||||||
Reference in New Issue
Block a user