From a8f7a11df53c602d80118d434962c14114740b1b Mon Sep 17 00:00:00 2001 From: Drew Rife Date: Thu, 20 Aug 2026 14:40:12 -0400 Subject: [PATCH] fix(j1939-21): respect sequence number for fragmented frames on reassembly --- j1939/j1939_21.py | 28 +++++++++++++- test/test_ecu.py | 95 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 122 insertions(+), 1 deletion(-) diff --git a/j1939/j1939_21.py b/j1939/j1939_21.py index e3e55ff..e399d3f 100644 --- a/j1939/j1939_21.py +++ b/j1939/j1939_21.py @@ -314,6 +314,7 @@ def _process_tp_cm(self, mid, dest_address, data, timestamp): 'pgn': pgn, 'message_size': message_size, 'num_packages': num_packages, + 'next_expected_packet': 1, 'next_packet': min(self._max_cmdt_packets, max_num_packages), 'max_cmdt_packages': self._max_cmdt_packets, 'num_packages_max_rec': min(self._max_cmdt_packets, max_num_packages), @@ -379,6 +380,7 @@ def _process_tp_cm(self, mid, dest_address, data, timestamp): "pgn": pgn, "message_size": message_size, "num_packages": num_packages, + "next_expected_packet": 1, "next_packet": 1, "max_cmdt_packages": self._max_cmdt_packets, "data": [], @@ -410,8 +412,33 @@ def _process_tp_dt(self, mid, dest_address, data, timestamp): # TODO: LOG/TRACE/EXCEPTION? return + expected_sequence_number = self._rcv_buffer[buffer_hash]['next_expected_packet'] + if sequence_number != expected_sequence_number: + sequence_error = ( + "duplicate" + if 0 < sequence_number < expected_sequence_number + else "invalid" + ) + if dest_address != ParameterGroupNumber.Address.GLOBAL: + # Keep the wire value in the standardized reason range. The + # specific sequence error remains available in the exception. + self.__send_tp_abort( + dest_address, + src_address, + self.ConnectionAbortReason.RESOURCES, + self._rcv_buffer[buffer_hash]['pgn'], + ) + del self._rcv_buffer[buffer_hash] + self.__job_thread_wakeup() + raise ValueError( + "J1939-21 TP.DT packet out of sequence " + f"({sequence_error}): " + f"expected {expected_sequence_number}, received {sequence_number}" + ) + # get data self._rcv_buffer[buffer_hash]['data'].extend(data[1:]) + self._rcv_buffer[buffer_hash]['next_expected_packet'] += 1 # message is complete with sending an acknowledge if len(self._rcv_buffer[buffer_hash]['data']) >= self._rcv_buffer[buffer_hash]['message_size']: @@ -564,4 +591,3 @@ def notify(self, can_id, data, timestamp): # stack having to own the destination address. self.__notify_subscribers(mid.priority, pgn_value, mid.source_address, dest_address, timestamp, data) return - diff --git a/test/test_ecu.py b/test/test_ecu.py index 3d789b5..843acb9 100644 --- a/test/test_ecu.py +++ b/test/test_ecu.py @@ -1,8 +1,11 @@ import time import can +import pytest import j1939 +from j1939.j1939_21 import J1939_21 +from j1939.message_id import MessageId from test.helpers.feeder import Feeder # def test_connect(self): @@ -58,6 +61,98 @@ def test_broadcast_receive_long(feeder): feeder.receive() +def test_broadcast_receive_out_of_sequence_packet_raises(): + """Reject and terminate a BAM session with an invalid sequence number.""" + sent = [] + notified = [] + dll = J1939_21( + send_message=lambda *args: sent.append(args), + job_thread_wakeup=lambda: None, + notify_subscribers=lambda *args: notified.append(args), + max_cmdt_packets=1, + minimum_tp_rts_cts_dt_interval=None, + minimum_tp_bam_dt_interval=None, + ecu_is_message_acceptable=lambda dest: True, + ) + bam_mid = MessageId(can_id=0x00ECFF01) + dll._process_tp_cm(bam_mid, 0xFF, [32, 20, 0, 3, 255, 0xB0, 0xFE, 0], 0.0) + buffer_hash = dll._buffer_hash(0x01, 0xFF) + mid = MessageId(can_id=0x00EBFF01) + + with pytest.raises(ValueError, match='out of sequence'): + dll._process_tp_dt(mid, 0xFF, [2, 8, 9, 10, 11, 12, 13, 14], 0.0) + + assert sent == [] + assert notified == [] + assert buffer_hash not in dll._rcv_buffer + + +def test_peer_to_peer_receive_out_of_sequence_packet_aborts(): + """Reject an out-of-sequence CMDT packet and abort the session.""" + sent = [] + dll = J1939_21( + send_message=lambda *args: sent.append(args), + job_thread_wakeup=lambda: None, + notify_subscribers=lambda *args: None, + max_cmdt_packets=1, + minimum_tp_rts_cts_dt_interval=None, + minimum_tp_bam_dt_interval=None, + ecu_is_message_acceptable=lambda dest: True, + ) + rts_mid = MessageId(can_id=0x00EC0201) + dll._process_tp_cm(rts_mid, 0x02, [16, 20, 0, 3, 1, 0, 223, 0], 0.0) + sent.clear() + + dt_mid = MessageId(can_id=0x00EB0201) + with pytest.raises(ValueError, match='out of sequence'): + dll._process_tp_dt(dt_mid, 0x02, [2, 1, 2, 3, 4, 5, 6, 7], 0.0) + + assert sent == [ + ( + 0x1CEC0102, + True, + [255, 2, 255, 255, 255, 0, 223, 0], + ) + ] + assert not dll._rcv_buffer + + dll._process_tp_cm(rts_mid, 0x02, [16, 20, 0, 3, 1, 0, 223, 0], 0.0) + sent.clear() + + with pytest.raises(ValueError, match='out of sequence'): + dll._process_tp_dt(dt_mid, 0x02, [0, 1, 2, 3, 4, 5, 6, 7], 0.0) + + assert sent[0][2][1] == 2 + + +def test_peer_to_peer_sequence_gap_after_valid_packet_aborts(): + """Reject a sequence gap without delivering a partial RTS/CTS payload.""" + sent = [] + notified = [] + dll = J1939_21( + send_message=lambda *args: sent.append(args), + job_thread_wakeup=lambda: None, + notify_subscribers=lambda *args: notified.append(args), + max_cmdt_packets=2, + minimum_tp_rts_cts_dt_interval=None, + minimum_tp_bam_dt_interval=None, + ecu_is_message_acceptable=lambda dest: True, + ) + rts_mid = MessageId(can_id=0x00EC0201) + dll._process_tp_cm(rts_mid, 0x02, [16, 20, 0, 3, 2, 0, 223, 0], 0.0) + sent.clear() + + dt_mid = MessageId(can_id=0x00EB0201) + dll._process_tp_dt(dt_mid, 0x02, [1, 1, 2, 3, 4, 5, 6, 7], 0.0) + + with pytest.raises(ValueError, match='out of sequence'): + dll._process_tp_dt(dt_mid, 0x02, [3, 8, 9, 10, 11, 12, 13, 14], 0.0) + + assert sent[0][2][1] == 2 + assert notified == [] + assert not dll._rcv_buffer + + def test_peer_to_peer_receive_short(feeder): """Test the receivement of a normal peer-to-peer message