From 948e3df0d95f2b050eb060a2a8882713166f5023 Mon Sep 17 00:00:00 2001 From: dinacaran <75803638+dinacaran@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:16:21 +0300 Subject: [PATCH 1/3] fix(j1939): decode SA-specific PGNs only from their own source address When a database defines a J1939 PGN at more than one source address (e.g. DM1_239 at SA 0xEF and DM1_243 at SA 0xF3), a frame of that PGN now only decodes into the message for its own SA. Frames from an SA the database does not define stay undecoded instead of borrowing the first node's message, which merged DM1 from every ECU into one signal. A PGN the database defines at a single SA keeps the any-SA fallback, since that SA is usually a placeholder. The candidate lookup moves to DBCDecoder.candidates_for(), which the vectorized decoder and the debug inspector now call instead of keeping their own copies of the rule. Refs #13 Co-Authored-By: Claude Opus 5.5 --- core/dbc_decoder.py | 51 +++++++++++++++++---- core/debug_inspector.py | 18 +------- core/vectorized_decoder.py | 22 +-------- tests/test_dbc_decoder.py | 94 ++++++++++++++++++++++++++++++++++++++ 4 files changed, 138 insertions(+), 47 deletions(-) diff --git a/core/dbc_decoder.py b/core/dbc_decoder.py index 43619fe..92bc2f4 100644 --- a/core/dbc_decoder.py +++ b/core/dbc_decoder.py @@ -243,6 +243,13 @@ def __init__(self, dbc_path: str | Path) -> None: # Primary lookup: arbitration_id → [message, ...] self._messages_exact: dict[int, list[Any]] = {} self._messages_pgn: dict[int, list[Any]] = {} + # (PGN, source address) → [message, ...]. Used for the PGNs in + # _pgn_sa_specific, which the database defines at more than one SA: + # their frames only decode into the message for their own SA, and an + # SA the database does not define stays undecoded instead of + # borrowing another node's message. + self._messages_pgn_sa: dict[tuple[int, int], list[Any]] = {} + self._pgn_sa_specific: set[int] = set() # Perf: per-frame candidate cache (same ID seen repeatedly → reuse result) self._candidate_cache: dict[int, list[Any]] = {} @@ -291,12 +298,20 @@ def _build_indexes(self) -> None: pgn = self._extract_j1939_pgn(frame_id) if pgn is not None: self._messages_pgn.setdefault(pgn, []).append(message) + self._messages_pgn_sa.setdefault( + (pgn, frame_id & 0xFF), [] + ).append(message) # Perf: pre-cache signal choices so decode_frame avoids getattr per sample for signal in getattr(message, "signals", []): choices = getattr(signal, "choices", None) or {} self._choices_cache[(message.name, signal.name)] = dict(choices) + self._pgn_sa_specific = { + pgn for pgn, messages in self._messages_pgn.items() + if len({int(m.frame_id) & 0xFF for m in messages}) > 1 + } + # ── Decode kwargs — built once, reused every frame ──────────────────── def _get_decode_kwargs(self) -> dict[str, Any]: @@ -329,16 +344,15 @@ def _extract_j1939_pgn(frame_id: int) -> int | None: ps = (can_id >> 8) & 0xFF return (pf << 8) if pf < 240 else ((pf << 8) | ps) - def _get_candidates(self, frame: RawFrame) -> list[Any]: - """ - Return message candidates for this frame's arbitration_id. - Result is cached after first lookup — same ID seen in every periodic frame. + def candidates_for(self, arb_id: int, is_extended: bool) -> list[Any]: """ - arb_id = frame.arbitration_id - cached = self._candidate_cache.get(arb_id) - if cached is not None: - return cached + Return the messages that may decode ``arb_id``, best match first. + Exact and masked ID matches come first. Extended frames then fall back + to J1939 PGN matching: a PGN the database defines at a single source + address matches any SA (the DBC's SA is a placeholder), while a PGN + defined at several SAs only matches the message for the frame's SA. + """ seen: set[tuple[str, int]] = set() candidates: list[Any] = [] @@ -354,12 +368,29 @@ def add(msg: Any) -> None: add(msg) # J1939 PGN fallback - if frame.is_extended_id or arb_id > 0x7FF: + if is_extended or arb_id > 0x7FF: pgn = self._extract_j1939_pgn(arb_id) if pgn is not None: - for msg in self._messages_pgn.get(pgn, []): + if pgn in self._pgn_sa_specific: + matches = self._messages_pgn_sa.get((pgn, arb_id & 0xFF), []) + else: + matches = self._messages_pgn.get(pgn, []) + for msg in matches: add(msg) + return candidates + + def _get_candidates(self, frame: RawFrame) -> list[Any]: + """ + Return message candidates for this frame's arbitration_id. + Result is cached after first lookup — same ID seen in every periodic frame. + """ + arb_id = frame.arbitration_id + cached = self._candidate_cache.get(arb_id) + if cached is not None: + return cached + + candidates = self.candidates_for(arb_id, frame.is_extended_id) self._candidate_cache[arb_id] = candidates return candidates diff --git a/core/debug_inspector.py b/core/debug_inspector.py index a6db27d..53c5ee4 100644 --- a/core/debug_inspector.py +++ b/core/debug_inspector.py @@ -1408,23 +1408,7 @@ def inspect_measurement( def _database_candidates(decoder, frame_id: int) -> list[object]: - candidates: list[object] = [] - seen: set[int] = set() - for lookup in (frame_id, frame_id & 0x1FFFFFFF, frame_id & 0x7FF): - for message in decoder._messages_exact.get(lookup, []): - marker = id(message) - if marker not in seen: - seen.add(marker) - candidates.append(message) - if frame_id > 0x7FF: - pgn = decoder._extract_j1939_pgn(frame_id) - if pgn is not None: - for message in decoder._messages_pgn.get(pgn, []): - marker = id(message) - if marker not in seen: - seen.add(marker) - candidates.append(message) - return candidates + return decoder.candidates_for(frame_id, frame_id > 0x7FF) def _inspect_ldf( diff --git a/core/vectorized_decoder.py b/core/vectorized_decoder.py index 0545554..2877d10 100644 --- a/core/vectorized_decoder.py +++ b/core/vectorized_decoder.py @@ -311,23 +311,5 @@ def get_message_decoder(self, message: Any) -> MessageVectorDecoder: return dec def get_candidates(self, arb_id: int, is_extended: bool) -> list[Any]: - """Mirror DBCDecoder._get_candidates without needing a RawFrame.""" - seen: set[tuple[str, int]] = set() - candidates: list[Any] = [] - - def add(msg: Any) -> None: - key = (getattr(msg, 'name', ''), int(getattr(msg, 'frame_id', -1))) - if key not in seen: - seen.add(key) - candidates.append(msg) - - for lookup_id in (arb_id, arb_id & 0x1FFFFFFF, arb_id & 0x7FF): - for msg in self.dbc._messages_exact.get(lookup_id, []): - add(msg) - - if is_extended or arb_id > 0x7FF: - pgn = self.dbc._extract_j1939_pgn(arb_id) - if pgn is not None: - for msg in self.dbc._messages_pgn.get(pgn, []): - add(msg) - return candidates + """Same candidates as DBCDecoder, without needing a RawFrame.""" + return self.dbc.candidates_for(arb_id, is_extended) diff --git a/tests/test_dbc_decoder.py b/tests/test_dbc_decoder.py index e36271e..8f5b77d 100644 --- a/tests/test_dbc_decoder.py +++ b/tests/test_dbc_decoder.py @@ -307,3 +307,97 @@ def test_diagnostics_text_contains_no_signals_counter(decoder, frame_diag): text = decoder.diagnostics_text() assert "Matched, no signals:" in text assert "1" in text + + +# ── J1939 source-address matching (issue #13) ───────────────────────────── + +_J1939_DBC = """\ +VERSION "" + +NS_ : + +BS_: + +BU_: MCU BMS ENG + +BO_ 2566834927 DM1_239: 8 MCU + SG_ DM1_239_DTC1 : 16|32@1+ (1,0) [0|4294967295] "" Vector__XXX + +BO_ 2566834931 DM1_243: 8 BMS + SG_ DM1_243_DTC1 : 16|32@1+ (1,0) [0|4294967295] "" Vector__XXX + +BO_ 2566840320 VehicleDist: 8 ENG + SG_ TotalDist : 0|32@1+ (0.125,0) [0|526385151.9] "km" Vector__XXX +""" + + +@pytest.fixture +def j1939_decoder(tmp_path): + """DM1 (PGN 0xFECA) at SAs 0xEF and 0xF3; VD (PGN 0xFEE0) at SA 0x00 only.""" + from core.dbc_decoder import DBCDecoder + + path = tmp_path / "j1939.dbc" + path.write_text(_J1939_DBC, encoding="utf-8") + return DBCDecoder(str(path)) + + +def _j1939_frame(arb_id: int, dtc: int = 0) -> RawFrame: + data = bytes([0x00, 0xFF]) + dtc.to_bytes(4, "little") + bytes([0xFF, 0xFF]) + return RawFrame( + timestamp=0.0, channel=1, arbitration_id=arb_id, + is_extended_id=True, is_fd=False, dlc=8, + data=data, direction="Rx", + ) + + +def _names(messages) -> list[str]: + return [m.name for m in messages] + + +@pytest.mark.parametrize("arb_id", [0x18FECA13, 0x18FECAE6]) +def test_j1939_undefined_sa_of_multi_sa_pgn_stays_undecoded(j1939_decoder, arb_id): + assert j1939_decoder.candidates_for(arb_id, True) == [] + assert j1939_decoder.decode_frame(_j1939_frame(arb_id)) == [] + + +@pytest.mark.parametrize("arb_id, name", [ + (0x18FECAEF, "DM1_239"), + (0x18FECAF3, "DM1_243"), +]) +def test_j1939_defined_sa_matches_only_its_own_message(j1939_decoder, arb_id, name): + assert _names(j1939_decoder.candidates_for(arb_id, True)) == [name] + + +def test_j1939_sa_match_ignores_priority(j1939_decoder): + # Priority 7 instead of the DBC's 6: same PGN and SA, still DM1_239. + samples = j1939_decoder.decode_frame(_j1939_frame(0x1CFECAEF, dtc=61765)) + assert [(s.signal_name, s.value) for s in samples] == [("DM1_239_DTC1", 61765)] + + +def test_j1939_single_sa_pgn_keeps_placeholder_fallback(j1939_decoder): + # VehicleDist is defined once, at SA 0x00; any sender's frame decodes into it. + assert _names(j1939_decoder.candidates_for(0x18FEE0FE, True)) == ["VehicleDist"] + + +def test_vectorized_candidates_follow_decoder(j1939_decoder): + from core.vectorized_decoder import VectorizedDBC + + vec = VectorizedDBC(j1939_decoder) + assert vec.get_candidates(0x18FECA13, is_extended=True) == [] + assert _names(vec.get_candidates(0x18FECAEF, is_extended=True)) == ["DM1_239"] + assert _names(vec.get_candidates(0x18FEE0FE, is_extended=True)) == ["VehicleDist"] + + +def test_j1939_dm1_trace_holds_only_its_own_node(j1939_decoder): + # Four nodes send DM1 in the same millisecond; only SA 0xEF has a DTC. + frames = [ + _j1939_frame(0x18FECAEF, dtc=61765), + _j1939_frame(0x18FECA13), + _j1939_frame(0x18FECAE6), + _j1939_frame(0x18FECAF3), + ] + by_signal: dict[str, list] = {} + for frame in frames: + for s in j1939_decoder.decode_frame(frame): + by_signal.setdefault(s.signal_name, []).append(s.value) + assert by_signal == {"DM1_239_DTC1": [61765], "DM1_243_DTC1": [0]} From f482afac8ca13b8c42d72a2dd1fe4485be443728 Mon Sep 17 00:00:00 2001 From: dinacaran <75803638+dinacaran@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:37:46 +0300 Subject: [PATCH 2/3] fix(j1939): split shared messages per sender and keep series in time order A database often defines a J1939 message once at a placeholder source address while several ECUs send it. In a truck log, CCVS came from four nodes: two with real brake/clutch switch states and two sending 3 (not available). All four decoded into one series, and because the bulk loaders insert one frame ID at a time, the series held four back-to-back sweeps of the recording instead of one time-ordered trace. Zoomed out, the plot drew a solid box between 0 and 3; zoomed in, pyqtgraph's binary-search clipping picked one sweep and showed a flat 3. - A message that frames from more than one source address decode into now gets one series per sender, named "CCVS [SA 0x17]". A message with a single sender keeps its database name. - SignalStore flags any series whose bulk insert starts before its last sample, and sort_merged_series() reorders only those series once loading finishes. All three bulk load paths call it, so one sender on two priorities still loads as one time-ordered series. - The MF4 reader applies the same rules to asammdf's extraction. asammdf does its own J1939 matching and gives each frame ID its own group. The reader now reads the frame ID from each group's comment and drops groups our decoder would not match, which brings the #13 rule to MF4 files. It also splits multi-sender messages the same way instead of merging the groups back together. Refs #13 Co-Authored-By: Claude Opus 5.5 --- core/dbc_decoder.py | 24 +++- core/load_worker.py | 45 ++++-- core/readers/mdf_can_reader.py | 99 +++++++++++++- core/signal_store.py | 36 +++++ tests/test_j1939_senders.py | 243 +++++++++++++++++++++++++++++++++ 5 files changed, 432 insertions(+), 15 deletions(-) create mode 100644 tests/test_j1939_senders.py diff --git a/core/dbc_decoder.py b/core/dbc_decoder.py index 92bc2f4..aa180ee 100644 --- a/core/dbc_decoder.py +++ b/core/dbc_decoder.py @@ -8,7 +8,7 @@ from dataclasses import dataclass from pathlib import Path -from typing import Any +from typing import Any, Hashable, Iterable import inspect import xml.etree.ElementTree as ET @@ -233,6 +233,28 @@ def load_database_file( return db, load_messages +def source_address_message_name(message_name: str, source_address: int) -> str: + """Name the series of one J1939 sender of a message that several send.""" + return f"{message_name} [SA 0x{source_address:02X}]" + + +def multi_sender_messages(matches: Iterable[tuple[Hashable, int]]) -> set[Hashable]: + """ + Return the keys matched by extended frames from more than one source address. + + *matches* pairs a key identifying a database message with the ID of a + frame that decoded into it. A J1939 message matched through its PGN + placeholder can be sent by several ECUs; their series are kept apart + instead of interleaving, e.g. one ECU reporting a switch and another + reporting it as not available. + """ + sources: dict[Hashable, set[int]] = {} + for key, frame_id in matches: + if frame_id > 0x7FF: + sources.setdefault(key, set()).add(frame_id & 0xFF) + return {key for key, addresses in sources.items() if len(addresses) > 1} + + class DBCDecoder: def __init__(self, dbc_path: str | Path) -> None: self.dbc_path = Path(dbc_path) diff --git a/core/load_worker.py b/core/load_worker.py index 76ebdc7..784582d 100644 --- a/core/load_worker.py +++ b/core/load_worker.py @@ -22,6 +22,7 @@ store_key_prefix, ) from core.channel_config import ChannelConfig +from core.dbc_decoder import multi_sender_messages, source_address_message_name from core.signal_store import as_channel_key # ── Streaming constants ─────────────────────────────────────────────────── @@ -434,11 +435,12 @@ def flush_batch() -> None: f"Bulk decoding {n:,} frames across {total_groups:,} CAN ID groups..." ) + # Resolve every group's message before decoding any: whether a + # message gets one series per J1939 source address depends on all of + # the frame IDs that matched it. + group_matches: list[tuple | None] = [] for g in range(total_groups): - start, end = int(boundaries[g]), int(boundaries[g + 1]) - group_idx = sort_idx[start:end] - - first = group_idx[0] + first = sort_idx[int(boundaries[g])] ch_byte = int(channels_np[first]) arb_id = int(arb_ids_np[first]) bus = BusType.LIN if is_lin_np[first] else BusType.CAN @@ -450,6 +452,7 @@ def flush_batch() -> None: or _decoder_map.get((bus, ALL_CHANNELS_NUMBER)) ) if decoder is None: + group_matches.append(None) continue vec = vec_dbcs.get(id(decoder)) @@ -458,13 +461,31 @@ def flush_batch() -> None: vec_dbcs[id(decoder)] = vec candidates = vec.get_candidates(arb_id, is_extended=(arb_id > 0x7FF)) - if not candidates: + # Match existing single-decoder behaviour: pick the first candidate. + group_matches.append( + (ch_key, arb_id, decoder, vec, candidates[0]) if candidates else None + ) + + split_messages = multi_sender_messages( + ((match[0], id(match[4])), match[1]) + for match in group_matches if match is not None + ) + + for g in range(total_groups): + start, end = int(boundaries[g]), int(boundaries[g + 1]) + group_idx = sort_idx[start:end] + + match = group_matches[g] + if match is None: continue + ch_key, arb_id, decoder, vec, message = match - # Match existing single-decoder behaviour: pick the first candidate. - message = candidates[0] - msg_name = message.name - msg_id = int(getattr(message, 'frame_id', arb_id)) + if (ch_key, id(message)) in split_messages: + msg_name = source_address_message_name(message.name, arb_id & 0xFF) + msg_id = arb_id + else: + msg_name = message.name + msg_id = int(getattr(message, 'frame_id', arb_id)) msg_dec = vec.get_message_decoder(message) # Frames shorter than the message they matched decode from the @@ -601,6 +622,9 @@ def flush_batch() -> None: # unmatched because decoding has not happened yet). store.decoded_frames = decoded_total store.unmatched_frames = n - decoded_total + # Groups are decoded one frame ID at a time, so a signal fed by + # several IDs holds one sweep per ID until it is reordered. + store.sort_merged_series() decode_elapsed = time.perf_counter() - decode_started n_sigs = len(store._series_by_key) hint = ( @@ -1158,6 +1182,8 @@ def on_metadata_ready(metadata_rows): store.decoded_frames = trace_decoded_frames store.unmatched_frames = trace_frames - trace_decoded_frames + # asammdf yields one group per frame ID; see sort_merged_series(). + store.sort_merged_series() import_elapsed = time.perf_counter() - import_start self.progress.emit( f"Bulk import complete: {total:,} signals | samples: " @@ -1254,6 +1280,7 @@ def on_metadata_ready(metadata_rows): f"Loaded {ch_count:,} channels | samples: {store.total_samples:,}" ) + store.sort_merged_series() self.tree_update.emit(store.build_tree_payload()) self.partial_ready.emit() if metadata_first: diff --git a/core/readers/mdf_can_reader.py b/core/readers/mdf_can_reader.py index c63d617..41a0f6f 100644 --- a/core/readers/mdf_can_reader.py +++ b/core/readers/mdf_can_reader.py @@ -38,7 +38,11 @@ from core.bus_types import BusType from core.models import RawFrame, DecodedSignalSample -from core.dbc_decoder import DBCDecoder +from core.dbc_decoder import ( + DBCDecoder, + multi_sender_messages, + source_address_message_name, +) from core.raw_frame_store import FLAG_LIN from core.readers.mdf_reader import MDFReader, _channel_failure_text from core.readers.mdf_recovery import ( @@ -315,9 +319,15 @@ def iter_decoded_channel_arrays( message_name for _channel, message_name, _message_id, _signal_name, _unit in native_metadata_rows } + j1939_groups = ( + self._j1939_group_identities(extracted, channel_config) + if extracted is not None else {} + ) for group_idx, group in enumerate( extracted.groups if extracted is not None else () ): + if group_idx in j1939_groups and j1939_groups[group_idx] is None: + continue for ch_idx, decoded_channel in enumerate(group.channels): signal_name = ( getattr(decoded_channel, "name", None) or f"Ch{ch_idx}" @@ -327,9 +337,12 @@ def iter_decoded_channel_arrays( if signal_name.lower() in ("time", "t", "timestamps"): continue unit = str(getattr(decoded_channel, "unit", "") or "") - channel, message_name, message_id = self._decoded_group_metadata( - extracted, group_idx, signal_name, channel_config - ) + if group_idx in j1939_groups: + channel, message_name, message_id = j1939_groups[group_idx] + else: + channel, message_name, message_id = self._decoded_group_metadata( + extracted, group_idx, signal_name, channel_config + ) # Normally DBC signals have a concrete CAN channel and # therefore cannot collide with recorder-decoded CH? # signals. Keep both visible even when malformed metadata @@ -385,10 +398,14 @@ def iter_decoded_channel_arrays( batch_all_groups=True, channel_error=self._record_channel_error, ): + if group_idx in j1939_groups and j1939_groups[group_idx] is None: + continue signal_name = old_meta[1] unit = old_meta[2] meta = metadata_by_key.get((group_idx, signal_name)) - if meta is None: + if meta is None and group_idx in j1939_groups: + channel, message_name, message_id = j1939_groups[group_idx] + elif meta is None: channel, message_name, message_id = self._decoded_group_metadata( extracted, group_idx, signal_name, channel_config ) @@ -604,6 +621,78 @@ def field(suffix, required=True): channel_error(group_name, "LIN_Frame", exc) return total + @classmethod + def _j1939_group_identities(cls, extracted, channel_config): + """ + Identify the J1939 groups asammdf extracted, one per sender. + + asammdf extracts J1939 frames into one group per frame ID and records + the sender in the group comment (``CAN1 ID=0x18FEF100 CCVS PGN=0xFEF1 + SA=0x0``). Its acquisition source only carries the database ID, and + its ``acq_name`` prints the SA in decimal behind a ``0x`` prefix, so + the comment is the reliable record of which frame ID fed the group. + + Returns ``{group_idx: (channel, message_name, frame_id) | None}`` for + J1939 groups only. ``None`` marks a group the analyzer's own decoder + would not match: asammdf assigns an undefined sender of a PGN the + database defines at several source addresses to one of those + messages, which would put another node's data in its signals. When + several senders share one message, each gets its own message name. + """ + identities = {} + for group_idx, group in enumerate(extracted.groups): + channel_group = getattr(group, "channel_group", None) + comment = str(getattr(channel_group, "comment", "") or "") + match = re.search( + r"\bID=0x([0-9A-F]+)\b.*\bSA=0x([0-9A-F]+)\b", comment, + re.IGNORECASE, + ) + if not match: + continue + signal_name = next( + ( + decoded_channel.name + for decoded_channel in group.channels + if getattr(decoded_channel, "channel_type", -1) != 1 + and decoded_channel.name.lower() not in ("time", "t", "timestamps") + ), + None, + ) + if signal_name is None: + continue + channel, message_name, _database_id = cls._decoded_group_metadata( + extracted, group_idx, signal_name, channel_config + ) + frame_id = int(match.group(1), 16) + identity = (channel, message_name, frame_id) + if channel_config is not None and channel is not None: + try: + decoder = channel_config.decoder_for(*channel) + except Exception: + decoder = None + if decoder is not None and hasattr(decoder, "candidates_for"): + names = { + message.name + for message in decoder.candidates_for(frame_id, True) + } + if message_name not in names: + identity = None + identities[group_idx] = identity + + split = multi_sender_messages( + ((identity[0], identity[1]), identity[2]) + for identity in identities.values() if identity is not None + ) + for group_idx, identity in identities.items(): + if identity is not None and (identity[0], identity[1]) in split: + channel, message_name, frame_id = identity + identities[group_idx] = ( + channel, + source_address_message_name(message_name, frame_id & 0xFF), + frame_id, + ) + return identities + @staticmethod def _decoded_group_metadata(extracted, group_idx, signal_name, channel_config): """ diff --git a/core/signal_store.py b/core/signal_store.py index 61289d4..ac907cd 100644 --- a/core/signal_store.py +++ b/core/signal_store.py @@ -100,6 +100,9 @@ def __init__(self) -> None: # and raw_values is never appended to. self._choices_lookup: dict[tuple[BusChannel | None, str, str], dict] = {} self._tree_dirty: bool = True + # Series whose bulk inserts arrived out of time order; see + # sort_merged_series(). + self._unsorted_keys: set[str] = set() self.total_frames = 0 self.decoded_frames = 0 self.total_samples = 0 @@ -329,6 +332,8 @@ def add_series_bulk( self._tree_dirty = True else: series = self._series_by_key[key] + if series.timestamps and timestamps[0] < series.timestamps[-1]: + self._unsorted_keys.add(key) # C-level memcopy — no Python loop, no object allocation per sample ts_bytes = np.asarray(timestamps, dtype=np.float64).tobytes() @@ -351,6 +356,37 @@ def add_series_bulk( self.channels.add(channel) self.message_hits[(channel, message_name)] += n + def sort_merged_series(self) -> int: + """ + Put series assembled from several bulk inserts back in time order. + + A signal fed by more than one frame ID (J1939 senders sharing one + database message, or priority variants of one ID) is inserted one ID + at a time, so its samples arrive as back-to-back sweeps of the whole + recording. Plot clipping and cursor lookups binary-search the + timestamps, which only works when they increase. Returns the number + of series reordered. + """ + reordered = 0 + for key in self._unsorted_keys: + series = self._series_by_key.get(key) + if series is None: + continue + ts = series.numpy_timestamps() + order = np.argsort(ts, kind="stable") + timestamps = _array.array("d") + timestamps.frombytes(ts[order].tobytes()) + values = _array.array("d") + values.frombytes(series.numpy_values()[order].tobytes()) + series.timestamps = timestamps + series.values = values + if series.raw_values and len(series.raw_values) == len(order): + raw_values = series.raw_values + series.raw_values = [raw_values[i] for i in order.tolist()] + reordered += 1 + self._unsorted_keys.clear() + return reordered + def normalize_timestamps(self, already_normalized: bool = False) -> None: """ Shift all timestamps so that t=0 is the start of the recording. diff --git a/tests/test_j1939_senders.py b/tests/test_j1939_senders.py new file mode 100644 index 0000000..5284140 --- /dev/null +++ b/tests/test_j1939_senders.py @@ -0,0 +1,243 @@ +# This Source Code Form is subject to the terms of the Mozilla Public +# License, v. 2.0. If a copy of the MPL was not distributed with this +# file, You can obtain one at https://mozilla.org/MPL/2.0/. +# +# Copyright (c) 2025-2026 Dinakaran Ganesan + +""" +J1939 messages sent by several ECUs, through the real load paths. + +A database often defines a J1939 message once, at a placeholder source +address, while several ECUs send it: in a truck log CCVS came from the engine +and the body controller with real switch states and from two other nodes with +"not available". Those senders must end up in separate series, each in time +order, for both the ASC/BLF loader (analyzer decoder) and the MF4 loader +(asammdf extraction). +""" +from __future__ import annotations + +import numpy as np +import pytest + +from core.channel_config import ChannelConfig +from core.load_worker import LoadWorker +from core.signal_store import SignalStore + + +_DM1_SIGNAL = ' SG_ {prefix}DTC1 : 16|32@1+ (1,0) [0|4294967295] "" Vector__XXX' + +_DBC_HEADER = """\ +VERSION "" + +NS_ : + +BS_: + +BU_: ECU + +""" + +_DBC_J1939 = """ +BA_DEF_ BO_ "VFrameFormat" ENUM "StandardCAN","ExtendedCAN","reserved","J1939PG"; +BA_DEF_ "ProtocolType" STRING ; +BA_DEF_DEF_ "VFrameFormat" "J1939PG"; +BA_DEF_DEF_ "ProtocolType" "J1939"; +BA_ "ProtocolType" "J1939"; +""" + + +def _dbc(tmp_path, messages: list[tuple[str, int]]): + """Write a J1939 DBC defining DM1 (PGN 0xFECA) as *messages* [(name, id)].""" + body = "".join( + f"BO_ {frame_id | 0x80000000} {name}: 8 ECU\n" + f"{_DM1_SIGNAL.format(prefix=name + '_')}\n\n" + for name, frame_id in messages + ) + path = tmp_path / "dm1.dbc" + path.write_text(_DBC_HEADER + body + _DBC_J1939, encoding="utf-8") + return path + + +def _dm1(dtc: int) -> bytes: + return bytes([0x00, 0xFF]) + dtc.to_bytes(4, "little") + bytes([0xFF, 0xFF]) + + +# SA 0x13 reports DTC 34243 in cycles 3-4, SA 0xEF reports 61765 in cycles +# 1-6, SA 0xE6 never reports one. All three send every 10 ms, SA 0xEF first. +_SENDERS = (0xEF, 0x13, 0xE6) +_EVENTS = {0xEF: (1, 7, 61765), 0x13: (3, 5, 34243)} +_CYCLES = 10 + + +def _frames(): + for cycle in range(_CYCLES): + for order, sa in enumerate(_SENDERS): + start, end, value = _EVENTS.get(sa, (0, 0, 0)) + dtc = value if start <= cycle < end else 0 + yield cycle * 0.010 + order * 0.0002, 0x18FECA00 | sa, _dm1(dtc) + + +def _write_asc(path): + lines = [ + "date Thu Sep 24 10:00:00.000 am 2026", + "base hex timestamps absolute", + "internal events logged", + "Begin Triggerblock Thu Sep 24 10:00:00.000 am 2026", + ] + for timestamp, frame_id, data in _frames(): + payload = " ".join(f"{byte:02X}" for byte in data) + lines.append( + f"{timestamp + 0.001:11.6f} 1 {frame_id:08X}x Rx d 8 {payload}" + ) + lines.append("End TriggerBlock") + path.write_text("\n".join(lines) + "\n", encoding="utf-8") + return path + + +def _write_mf4(path): + can = pytest.importorskip("can") + pytest.importorskip("asammdf") + with can.Logger(str(path)) as logger: + for timestamp, frame_id, data in _frames(): + logger.on_message_received(can.Message( + timestamp=1000.0 + timestamp, arbitration_id=frame_id, + is_extended_id=True, data=data, channel=0, + )) + return path + + +@pytest.fixture(params=["asc", "mf4"]) +def measurement(request, tmp_path): + """The same frames as ASC (analyzer decoder) and MF4 (asammdf).""" + writer = _write_asc if request.param == "asc" else _write_mf4 + return writer(tmp_path / f"dm1.{request.param}") + + +def _load(measurement, dbc_path) -> dict[str, tuple[np.ndarray, np.ndarray]]: + """Run the real load path; return {message::signal: (timestamps, values)}.""" + worker = LoadWorker(str(measurement), ChannelConfig.from_single_dbc(str(dbc_path))) + result: dict = {} + worker.finished.connect(lambda store: result.setdefault("store", store)) + worker.failed.connect(lambda error: result.setdefault("error", error)) + worker.run() + if "error" in result: + raise AssertionError(result["error"]) + store = result["store"] + try: + return { + key.split("::", 1)[1]: ( + series.numpy_timestamps().copy(), series.numpy_values().copy() + ) + for key, series in store._series_by_key.items() + } + finally: + store.raw_frame_store.close() + + +def _nonzero(series) -> list[float]: + return sorted(set(series[1].tolist()) - {0.0}) + + +def test_placeholder_message_gets_one_series_per_sender(measurement, tmp_path): + dbc = _dbc(tmp_path, [("DM1", 0x18FECAFE)]) + + series = _load(measurement, dbc) + + assert set(series) == { + "DM1 [SA 0xEF]::DM1_DTC1", + "DM1 [SA 0x13]::DM1_DTC1", + "DM1 [SA 0xE6]::DM1_DTC1", + } + assert _nonzero(series["DM1 [SA 0xEF]::DM1_DTC1"]) == [61765.0] + assert _nonzero(series["DM1 [SA 0x13]::DM1_DTC1"]) == [34243.0] + assert _nonzero(series["DM1 [SA 0xE6]::DM1_DTC1"]) == [] + for timestamps, values in series.values(): + assert len(values) == _CYCLES + assert np.all(np.diff(timestamps) > 0) + + +def test_single_sender_keeps_the_database_name(tmp_path): + # Only SA 0xEF sends; there is nothing to tell apart. + asc = tmp_path / "one.asc" + _write_asc(asc) + asc.write_text( + "\n".join( + line for line in asc.read_text(encoding="utf-8").splitlines() + if "18FECA13x" not in line and "18FECAE6x" not in line + ) + "\n", + encoding="utf-8", + ) + dbc = _dbc(tmp_path, [("DM1", 0x18FECAFE)]) + + assert set(_load(asc, dbc)) == {"DM1::DM1_DTC1"} + + +def test_undefined_sender_of_sa_specific_pgn_stays_undecoded(measurement, tmp_path): + # Issue #13: DM1 defined for SA 0xEF and 0xF3; 0x13 and 0xE6 must not + # borrow either message, in the asammdf path as well as the decoder. + dbc = _dbc(tmp_path, [("DM1_239", 0x18FECAEF), ("DM1_243", 0x18FECAF3)]) + + series = _load(measurement, dbc) + + assert set(series) == {"DM1_239::DM1_239_DTC1"} + assert len(series["DM1_239::DM1_239_DTC1"][1]) == _CYCLES + assert _nonzero(series["DM1_239::DM1_239_DTC1"]) == [61765.0] + + +def test_one_sender_on_two_priorities_loads_in_time_order(tmp_path): + # SA 0xEF alternates priority 6 and 7: one sender, one series, but the + # loader decodes it one frame ID at a time. + asc = tmp_path / "priorities.asc" + _write_asc(asc) + lines = [ + line for line in asc.read_text(encoding="utf-8").splitlines() + if "18FECA13x" not in line and "18FECAE6x" not in line + ] + seen = 0 + for index, line in enumerate(lines): + if "18FECAEFx" in line: + if seen % 2: + lines[index] = line.replace("18FECAEFx", "1CFECAEFx") + seen += 1 + asc.write_text("\n".join(lines) + "\n", encoding="utf-8") + dbc = _dbc(tmp_path, [("DM1", 0x18FECAFE)]) + + series = _load(asc, dbc) + + assert set(series) == {"DM1::DM1_DTC1"} + timestamps, values = series["DM1::DM1_DTC1"] + assert len(values) == _CYCLES + assert np.all(np.diff(timestamps) > 0) + + +def test_store_reorders_series_fed_out_of_time_order(): + store = SignalStore() + for timestamps, values, labels in ( + ([0.0, 0.2, 0.4], [1.0, 1.0, 1.0], ["On", "On", "On"]), + ([0.1, 0.3], [3.0, 3.0], ["N/A", "N/A"]), + ): + store.add_series_bulk( + channel=1, message_name="CCVS", message_id=0x18FEF100, + signal_name="BrakeSwitch", unit="", + timestamps=np.array(timestamps), values=np.array(values), + raw_values=labels, has_labels=True, + ) + + assert store.sort_merged_series() == 1 + series = next(iter(store._series_by_key.values())) + assert series.numpy_timestamps().tolist() == [0.0, 0.1, 0.2, 0.3, 0.4] + assert series.numpy_values().tolist() == [1.0, 3.0, 1.0, 3.0, 1.0] + assert list(series.raw_values) == ["On", "N/A", "On", "N/A", "On"] + + +def test_store_leaves_in_order_series_alone(): + store = SignalStore() + for timestamps in ([0.0, 0.1], [0.2, 0.3]): + store.add_series_bulk( + channel=1, message_name="EEC1", message_id=0x0CF00400, + signal_name="EngSpeed", unit="rpm", + timestamps=np.array(timestamps), values=np.array([1.0, 2.0]), + raw_values=[], has_labels=False, + ) + + assert store.sort_merged_series() == 0 From ddd34db5a2e6c7ea7acd58884e3d30585bcadd48 Mon Sep 17 00:00:00 2001 From: dinacaran <75803638+dinacaran@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:40:22 +0300 Subject: [PATCH 3/3] fix(dbc): stop matching standard and extended IDs on their low 11 bits The decoder indexed every message under its ID masked to 11 bits and looked every frame up the same way, so the two ID types cross-matched: extended J1939 frame 0x18FECA13 became a candidate for a standard message 0x213 (and, listed first, won over its own DM1 message), and a standard frame 0x2EF decoded as the extended message 0x18FECAEF. Messages are now indexed by ID type. A standard frame only matches standard messages and an extended frame only extended ones; an ID that still carries the DBC's extended flag bit (0x80000000) keeps matching. The per-frame candidate cache is keyed on the ID type too, so a standard and an extended frame with the same ID no longer share an entry. Co-Authored-By: Claude Opus 5.5 --- core/dbc_decoder.py | 43 +++++++++++++++------------ tests/test_dbc_decoder.py | 62 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 86 insertions(+), 19 deletions(-) diff --git a/core/dbc_decoder.py b/core/dbc_decoder.py index aa180ee..bb56493 100644 --- a/core/dbc_decoder.py +++ b/core/dbc_decoder.py @@ -262,8 +262,9 @@ def __init__(self, dbc_path: str | Path) -> None: self._decode_signature = None self._decode_kwargs_cache: dict[str, Any] | None = None # perf: built once - # Primary lookup: arbitration_id → [message, ...] - self._messages_exact: dict[int, list[Any]] = {} + # Primary lookup: (is_extended, arbitration_id) → [message, ...]. A + # standard and an extended frame never match each other's messages. + self._messages_exact: dict[tuple[bool, int], list[Any]] = {} self._messages_pgn: dict[int, list[Any]] = {} # (PGN, source address) → [message, ...]. Used for the PGNs in # _pgn_sa_specific, which the database defines at more than one SA: @@ -274,7 +275,7 @@ def __init__(self, dbc_path: str | Path) -> None: self._pgn_sa_specific: set[int] = set() # Perf: per-frame candidate cache (same ID seen repeatedly → reuse result) - self._candidate_cache: dict[int, list[Any]] = {} + self._candidate_cache: dict[tuple[int, bool], list[Any]] = {} # Perf: per-signal choices dict cached at build time (avoids getattr per sample) # key = (message_name, signal_name) → {int_key: label_str} @@ -310,12 +311,12 @@ def _build_indexes(self) -> None: self._dbc_message_ids_preview.append( f"{message.name} | {frame_id_text} | len={getattr(message, 'length', '?')}" ) - # Register under all masked variants (exact, 29-bit, 11-bit) - for fid in (frame_id, frame_id & 0x1FFFFFFF, frame_id & 0x7FF): + is_extended = bool(getattr(message, "is_extended_frame", False)) or frame_id > 0x7FF + # Register under the exact ID and the ID without the extended flag bit + for fid in (frame_id, frame_id & 0x1FFFFFFF): if fid >= 0: - self._messages_exact.setdefault(fid, []).append(message) + self._messages_exact.setdefault((is_extended, fid), []).append(message) # J1939 PGN index - is_extended = bool(getattr(message, "is_extended_frame", False)) or frame_id > 0x7FF if is_extended: pgn = self._extract_j1939_pgn(frame_id) if pgn is not None: @@ -370,10 +371,12 @@ def candidates_for(self, arb_id: int, is_extended: bool) -> list[Any]: """ Return the messages that may decode ``arb_id``, best match first. - Exact and masked ID matches come first. Extended frames then fall back - to J1939 PGN matching: a PGN the database defines at a single source - address matches any SA (the DBC's SA is a placeholder), while a PGN - defined at several SAs only matches the message for the frame's SA. + Exact ID matches come first, among messages of the frame's own type: + a standard frame only matches standard messages and an extended frame + only extended ones. Extended frames then fall back to J1939 PGN + matching: a PGN the database defines at a single source address + matches any SA (the DBC's SA is a placeholder), while a PGN defined at + several SAs only matches the message for the frame's SA. """ seen: set[tuple[str, int]] = set() candidates: list[Any] = [] @@ -384,13 +387,15 @@ def add(msg: Any) -> None: seen.add(key) candidates.append(msg) - # Exact + masked lookups - for lookup_id in (arb_id, arb_id & 0x1FFFFFFF, arb_id & 0x7FF): - for msg in self._messages_exact.get(lookup_id, []): + extended = is_extended or arb_id > 0x7FF + + # Exact lookups, with and without the extended flag bit + for lookup_id in (arb_id, arb_id & 0x1FFFFFFF): + for msg in self._messages_exact.get((extended, lookup_id), []): add(msg) # J1939 PGN fallback - if is_extended or arb_id > 0x7FF: + if extended: pgn = self._extract_j1939_pgn(arb_id) if pgn is not None: if pgn in self._pgn_sa_specific: @@ -407,13 +412,13 @@ def _get_candidates(self, frame: RawFrame) -> list[Any]: Return message candidates for this frame's arbitration_id. Result is cached after first lookup — same ID seen in every periodic frame. """ - arb_id = frame.arbitration_id - cached = self._candidate_cache.get(arb_id) + cache_key = (frame.arbitration_id, bool(frame.is_extended_id)) + cached = self._candidate_cache.get(cache_key) if cached is not None: return cached - candidates = self.candidates_for(arb_id, frame.is_extended_id) - self._candidate_cache[arb_id] = candidates + candidates = self.candidates_for(*cache_key) + self._candidate_cache[cache_key] = candidates return candidates # ── Frame decode ────────────────────────────────────────────────────── diff --git a/tests/test_dbc_decoder.py b/tests/test_dbc_decoder.py index 8f5b77d..9e4bdb3 100644 --- a/tests/test_dbc_decoder.py +++ b/tests/test_dbc_decoder.py @@ -401,3 +401,65 @@ def test_j1939_dm1_trace_holds_only_its_own_node(j1939_decoder): for s in j1939_decoder.decode_frame(frame): by_signal.setdefault(s.signal_name, []).append(s.value) assert by_signal == {"DM1_239_DTC1": [61765], "DM1_243_DTC1": [0]} + + +# ── Standard vs extended IDs ────────────────────────────────────────────── + +_MIXED_ID_DBC = """\ +VERSION "" + +NS_ : + +BS_: + +BU_: GW MCU + +BO_ 531 Std213: 8 GW + SG_ Counter : 0|8@1+ (1,0) [0|255] "" Vector__XXX + +BO_ 2566834927 DM1: 8 MCU + SG_ DTC1 : 16|32@1+ (1,0) [0|4294967295] "" Vector__XXX +""" + + +@pytest.fixture +def mixed_id_decoder(tmp_path): + """Standard 0x213 and extended 0x18FECAEF (DM1, defined at SA 0xEF only).""" + from core.dbc_decoder import DBCDecoder + + path = tmp_path / "mixed.dbc" + path.write_text(_MIXED_ID_DBC, encoding="utf-8") + return DBCDecoder(str(path)) + + +def _frame(arb_id: int, is_extended: bool) -> RawFrame: + return RawFrame( + timestamp=0.0, channel=1, arbitration_id=arb_id, + is_extended_id=is_extended, is_fd=False, dlc=8, + data=bytes(8), direction="Rx", + ) + + +def test_extended_frame_does_not_match_standard_message_on_its_low_bits(mixed_id_decoder): + # 0x18FECA13 & 0x7FF == 0x213. The frame is DM1 from SA 0x13. + assert _names(mixed_id_decoder.candidates_for(0x18FECA13, True)) == ["DM1"] + + +def test_standard_frame_does_not_match_extended_message_on_its_low_bits(mixed_id_decoder): + # 0x18FECAEF & 0x7FF == 0x2EF. + assert mixed_id_decoder.candidates_for(0x2EF, False) == [] + + +def test_each_frame_type_still_matches_its_own_messages(mixed_id_decoder): + assert _names(mixed_id_decoder.candidates_for(0x213, False)) == ["Std213"] + assert _names(mixed_id_decoder.candidates_for(0x18FECAEF, True)) == ["DM1"] + # An ID that still carries the DBC's extended flag bit matches too. + assert _names(mixed_id_decoder.candidates_for(0x98FECAEF, True)) == ["DM1"] + + +def test_standard_and_extended_frame_with_the_same_id_decode_apart(mixed_id_decoder): + standard = mixed_id_decoder.decode_frame(_frame(0x213, is_extended=False)) + extended = mixed_id_decoder.decode_frame(_frame(0x213, is_extended=True)) + + assert [s.message_name for s in standard] == ["Std213"] + assert extended == []