diff --git a/core/dbc_decoder.py b/core/dbc_decoder.py index 43619fe..bb56493 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) @@ -240,12 +262,20 @@ 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: + # 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]] = {} + 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} @@ -281,22 +311,30 @@ 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: 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 +367,17 @@ 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]: + def candidates_for(self, arb_id: int, is_extended: bool) -> list[Any]: """ - Return message candidates for this frame's arbitration_id. - Result is cached after first lookup — same ID seen in every periodic frame. + Return the messages that may decode ``arb_id``, best match first. + + 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. """ - arb_id = frame.arbitration_id - cached = self._candidate_cache.get(arb_id) - if cached is not None: - return cached - seen: set[tuple[str, int]] = set() candidates: list[Any] = [] @@ -348,19 +387,38 @@ 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 frame.is_extended_id or arb_id > 0x7FF: + if extended: 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) - self._candidate_cache[arb_id] = candidates + 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. + """ + 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(*cache_key) + self._candidate_cache[cache_key] = candidates return candidates # ── Frame decode ────────────────────────────────────────────────────── 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/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/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..9e4bdb3 100644 --- a/tests/test_dbc_decoder.py +++ b/tests/test_dbc_decoder.py @@ -307,3 +307,159 @@ 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]} + + +# ── 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 == [] 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