# File: decoder.py # Copyright (C) 2026 Erick Ahmed # SPDX-License-Identifier: AGPL-3.0-or-later import argparse import polars as pl def get_j1939_mask() -> pl.Expr: """ Returns a Polars expression representing the strict J1939 filtering rules. """ id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32) return id_int > 0x7FF def decode_j1939_metadata(lf: pl.LazyFrame) -> pl.LazyFrame: """ Decodes J1939 fields and bundles them into a Struct column. """ id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32) priority = ((id_int // 67108864) % 8).cast(pl.UInt8) pf = ((id_int // 65536) % 256).cast(pl.UInt8) ps = ((id_int // 256) % 256).cast(pl.UInt8) sa = (id_int % 256).cast(pl.UInt8) da = pl.when(pf < 240).then(ps).otherwise(pl.lit(255, dtype=pl.UInt8)).cast(pl.UInt8) pgn = pl.when(pf < 240).then( ((id_int // 256) & 0x3FF00) ).otherwise( ((id_int // 256) & 0x3FFFF) ).cast(pl.UInt32) return lf.with_columns( pl.struct([ priority.alias("Priority"), pf.alias("PF"), ps.alias("PS"), sa.alias("SA"), da.alias("DA"), pgn.alias("PGN") ]).alias("j1939_metadata") ) def decode_j1939_frames(df: pl.DataFrame) -> pl.DataFrame: id_int = pl.col("ID").str.to_integer(base=16).cast(pl.UInt32) is_j1939 = id_int > 0x7FF priority = ((id_int // 67108864) % 8).cast(pl.UInt8) pf = ((id_int // 65536) % 256).cast(pl.UInt8) ps = ((id_int // 256) % 256).cast(pl.UInt8) sa = (id_int % 256).cast(pl.UInt8) da = pl.when(pf < 240).then(ps).otherwise(pl.lit(255, dtype=pl.UInt8)).cast(pl.UInt8) pgn = pl.when(pf < 240).then( ((id_int // 256) & 0x3FF00) ).otherwise( ((id_int // 256) & 0x3FFFF) ).cast(pl.UInt32) j1939_meta = pl.when(is_j1939).then( pl.struct([ priority.alias("Priority"), pf.alias("PF"), ps.alias("PS"), sa.alias("SA"), da.alias("DA"), pgn.alias("PGN") ]) ).otherwise(None) return df.with_columns(j1939_meta.alias("j1939_metadata")) if __name__ == "__main__": parser = argparse.ArgumentParser(description="J1939 decoder") parser.add_argument("input_parquet", help="Path to the raw .parquet file") parser.add_argument("output_parquet", help="Path to save the decoded .parquet file") args = parser.parse_args() df = pl.scan_parquet(args.input_parquet).collect() decoded_df = decode_j1939_frames(df) decoded_df.write_parquet(args.output_parquet)