From 4dd2139ef4efb7594394df16e9214d9561c89cee Mon Sep 17 00:00:00 2001 From: projectmoon Date: Tue, 4 Aug 2026 23:30:14 +0200 Subject: [PATCH 01/11] Handle very rare deadlock situation when loading nodeinfo packets from radio. --- src/connections/stream_buffer.rs | 66 ++++++++++++++++++++++++++++++++ 1 file changed, 66 insertions(+) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 78e228c..d97b9ff 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -37,12 +37,20 @@ pub enum StreamBufferError { MissingLSB { lsb_index: usize }, #[error("Detected malformed packet, packet buffer contains a framing byte at index {next_packet_start_idx}")] MalformedPacket { next_packet_start_idx: usize }, + #[error("Declared packet size of {declared_size} bytes exceeds max packet size of {max_size} bytes")] + ImplausiblePacketSize { + declared_size: usize, + max_size: usize, + }, #[error(transparent)] DecodeFailure(#[from] prost::DecodeError), } const PACKET_HEADER_SIZE: usize = 4; +// Matches firmware's MAX_TO_FROM_RADIO_SIZE (src/mesh/PhoneAPI.h). +const MAX_TO_FROM_RADIO_SIZE: usize = 512; + impl StreamBuffer { /// Creates a new StreamBuffer instance that will send decoded FromRadio packets /// to the given broadcast channel. @@ -123,6 +131,17 @@ impl StreamBuffer { continue; // Don't need more data to continue, purge from buffer } + StreamBufferError::ImplausiblePacketSize { + declared_size, + max_size, + } => { + error!( + "Header declares implausible packet size {declared_size} (max {max_size}), discarding false header" + ); + + self.buffer.drain(0..PACKET_HEADER_SIZE.min(self.buffer.len())); + continue; // Don't need more data to continue, purge false header + } StreamBufferError::DecodeFailure { .. } => { error!("Failed to decode chunk from packet, this does not affect the next iteration"); @@ -286,6 +305,13 @@ impl StreamBuffer { // Recall that packet size doesn't include the first four magic bytes let incoming_packet_data_size: usize = usize::from(u16::from_le_bytes([*lsb, *msb])); + if incoming_packet_data_size > MAX_TO_FROM_RADIO_SIZE { + return Err(StreamBufferError::ImplausiblePacketSize { + declared_size: incoming_packet_data_size, + max_size: MAX_TO_FROM_RADIO_SIZE, + }); + } + Ok(incoming_packet_data_size) } @@ -718,6 +744,46 @@ mod tests { #[tokio::test] async fn detect_malformed_packets_with_internal_header_sequence() {} + /// Test for a [0x94, 0xc3] sequence occurring by chance inside another packet's + /// payload. Expected behavior is that the bogus size is rejected and the buffer + /// resyncs on the next real packet. + #[tokio::test] + async fn resync_after_false_header_match_with_implausible_size() { + // Arrange + + let payload_variant_1 = + protobufs::from_radio::PayloadVariant::MyInfo(protobufs::MyNodeInfo::default()); + let payload_variant_2 = + protobufs::from_radio::PayloadVariant::MyInfo(protobufs::MyNodeInfo::default()); + + let (packet_1, packet_data_1) = mock_encoded_from_radio_packet(payload_variant_1, None); + let (packet_2, packet_data_2) = mock_encoded_from_radio_packet(payload_variant_2, None); + + let encoded_packet_1 = format_data_packet(packet_data_1.into()).unwrap(); + let encoded_packet_2 = format_data_packet(packet_data_2.into()).unwrap(); + + // Header with a declared size (big-endian MSB/LSB) larger than MAX_TO_FROM_RADIO_SIZE. + let false_header = vec![0x94, 0xc3, 254, 6]; + + let (mock_tx, mut mock_rx) = unbounded_channel::(); + + // Act + + let mut buffer = StreamBuffer::new(mock_tx); + + let mut data = Vec::new(); + data.extend_from_slice(encoded_packet_1.data()); + data.extend_from_slice(&false_header); + data.extend_from_slice(encoded_packet_2.data()); + buffer.process_incoming_bytes(data.into()); + + // Assert + + assert_eq!(timeout_test(mock_rx.recv(), None).await, Some(packet_1)); + assert_eq!(timeout_test(mock_rx.recv(), None).await, Some(packet_2)); + assert_eq!(buffer.buffer.len(), 0); + } + #[tokio::test] async fn process_log_lines() { let (mock_tx, mut _mock_rx) = unbounded_channel::(); From 7d68ebdfd460d8b83a04320e342c7af82f9a363a Mon Sep 17 00:00:00 2001 From: projectmoon Date: Tue, 4 Aug 2026 23:40:01 +0200 Subject: [PATCH 02/11] cargo format --- src/connections/stream_buffer.rs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index d97b9ff..79e4921 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -37,7 +37,9 @@ pub enum StreamBufferError { MissingLSB { lsb_index: usize }, #[error("Detected malformed packet, packet buffer contains a framing byte at index {next_packet_start_idx}")] MalformedPacket { next_packet_start_idx: usize }, - #[error("Declared packet size of {declared_size} bytes exceeds max packet size of {max_size} bytes")] + #[error( + "Declared packet size of {declared_size} bytes exceeds max packet size of {max_size} bytes" + )] ImplausiblePacketSize { declared_size: usize, max_size: usize, @@ -139,7 +141,8 @@ impl StreamBuffer { "Header declares implausible packet size {declared_size} (max {max_size}), discarding false header" ); - self.buffer.drain(0..PACKET_HEADER_SIZE.min(self.buffer.len())); + self.buffer + .drain(0..PACKET_HEADER_SIZE.min(self.buffer.len())); continue; // Don't need more data to continue, purge false header } StreamBufferError::DecodeFailure { .. } => { From f3f54bcbfba18d8d59d4586a682a6fa68057a8b8 Mon Sep 17 00:00:00 2001 From: projectmoon Date: Wed, 5 Aug 2026 00:26:58 +0200 Subject: [PATCH 03/11] Better solution --- src/connections/stream_buffer.rs | 84 ++++++++++++++++++++++++++++++++ 1 file changed, 84 insertions(+) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 79e4921..df9b86b 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -352,6 +352,13 @@ impl StreamBuffer { .map(|idx| idx + packet_data_start_index); if let Some(next_packet_start_idx) = next_packet_start_index { + // Only malformed if it doesn't decode at the declared length. + let declared_packet = + &self.buffer[packet_data_start_index..packet_data_start_index + packet_data_size]; + if protobufs::FromRadio::decode(declared_packet).is_ok() { + return Ok(()); + } + // Remove malformed packet from buffer self.buffer.drain(..next_packet_start_idx); @@ -787,6 +794,83 @@ mod tests { assert_eq!(buffer.buffer.len(), 0); } + /// `outer_packet`'s serialized bytes coincidentally contain [0x94, 0xc3], + /// the same collision that caused a real hang before MAX_TO_FROM_RADIO_SIZE + /// was added. Identity fields are placeholders; device_metrics is unchanged + /// since that's what produces the collision. Expected behavior is that both + /// packets decode, since the coincidental match doesn't mean outer_packet + /// is malformed. + #[tokio::test] + async fn resync_after_coincidental_header_match_inside_packet_payload() { + let outer_packet = protobufs::FromRadio { + id: 0, + payload_variant: Some(protobufs::from_radio::PayloadVariant::NodeInfo( + protobufs::NodeInfo { + num: 1234567890, + user: Some(protobufs::User { + id: "!deadbeef".to_string(), + long_name: "Test Node".to_string(), + short_name: "TEST".to_string(), + hw_model: protobufs::HardwareModel::HeltecV3 as i32, + public_key: vec![0xab; 32], + ..Default::default() + }), + device_metrics: Some(protobufs::DeviceMetrics { + battery_level: Some(100), + voltage: Some(4.219), + channel_utilization: Some(22.93), + air_util_tx: Some(2.791), + uptime_seconds: Some(14655892), + }), + channel: 1, + via_mqtt: true, + hops_away: Some(1), + ..Default::default() + }, + )), + }; + let encoded_outer = format_data_packet(outer_packet.encode_to_vec().into()).unwrap(); + assert!( + encoded_outer.data_vec().windows(2).any(|w| w == [0x94, 0xc3]), + "outer_packet no longer reproduces the coincidental header collision" + ); + + let recovery_packet = protobufs::FromRadio { + id: 0, + payload_variant: Some(protobufs::from_radio::PayloadVariant::NodeInfo( + protobufs::NodeInfo { + num: 987654321, + user: Some(protobufs::User { + id: "!cafef00d".to_string(), + long_name: "Recovery Node".to_string(), + short_name: "RCVR".to_string(), + hw_model: protobufs::HardwareModel::HeltecV3 as i32, + public_key: vec![0xcd; 32], + ..Default::default() + }), + channel: 1, + via_mqtt: false, + hops_away: Some(2), + ..Default::default() + }, + )), + }; + let encoded_recovery = format_data_packet(recovery_packet.encode_to_vec().into()).unwrap(); + + let (mock_tx, mut mock_rx) = unbounded_channel::(); + let mut buffer = StreamBuffer::new(mock_tx); + + let mut data = encoded_outer.data_vec(); + data.extend_from_slice(encoded_recovery.data()); + buffer.process_incoming_bytes(data.into()); + + assert_eq!(timeout_test(mock_rx.recv(), None).await, Some(outer_packet)); + assert_eq!( + timeout_test(mock_rx.recv(), None).await, + Some(recovery_packet) + ); + } + #[tokio::test] async fn process_log_lines() { let (mock_tx, mut _mock_rx) = unbounded_channel::(); From 89a63d28472559b2d3099f8279fec502c2f599f5 Mon Sep 17 00:00:00 2001 From: projectmoon Date: Wed, 5 Aug 2026 00:33:48 +0200 Subject: [PATCH 04/11] minimal diff --- src/connections/stream_buffer.rs | 73 +------------------------------- 1 file changed, 2 insertions(+), 71 deletions(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index df9b86b..babe12c 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -37,22 +37,12 @@ pub enum StreamBufferError { MissingLSB { lsb_index: usize }, #[error("Detected malformed packet, packet buffer contains a framing byte at index {next_packet_start_idx}")] MalformedPacket { next_packet_start_idx: usize }, - #[error( - "Declared packet size of {declared_size} bytes exceeds max packet size of {max_size} bytes" - )] - ImplausiblePacketSize { - declared_size: usize, - max_size: usize, - }, #[error(transparent)] DecodeFailure(#[from] prost::DecodeError), } const PACKET_HEADER_SIZE: usize = 4; -// Matches firmware's MAX_TO_FROM_RADIO_SIZE (src/mesh/PhoneAPI.h). -const MAX_TO_FROM_RADIO_SIZE: usize = 512; - impl StreamBuffer { /// Creates a new StreamBuffer instance that will send decoded FromRadio packets /// to the given broadcast channel. @@ -133,18 +123,6 @@ impl StreamBuffer { continue; // Don't need more data to continue, purge from buffer } - StreamBufferError::ImplausiblePacketSize { - declared_size, - max_size, - } => { - error!( - "Header declares implausible packet size {declared_size} (max {max_size}), discarding false header" - ); - - self.buffer - .drain(0..PACKET_HEADER_SIZE.min(self.buffer.len())); - continue; // Don't need more data to continue, purge false header - } StreamBufferError::DecodeFailure { .. } => { error!("Failed to decode chunk from packet, this does not affect the next iteration"); @@ -308,13 +286,6 @@ impl StreamBuffer { // Recall that packet size doesn't include the first four magic bytes let incoming_packet_data_size: usize = usize::from(u16::from_le_bytes([*lsb, *msb])); - if incoming_packet_data_size > MAX_TO_FROM_RADIO_SIZE { - return Err(StreamBufferError::ImplausiblePacketSize { - declared_size: incoming_packet_data_size, - max_size: MAX_TO_FROM_RADIO_SIZE, - }); - } - Ok(incoming_packet_data_size) } @@ -754,49 +725,9 @@ mod tests { #[tokio::test] async fn detect_malformed_packets_with_internal_header_sequence() {} - /// Test for a [0x94, 0xc3] sequence occurring by chance inside another packet's - /// payload. Expected behavior is that the bogus size is rejected and the buffer - /// resyncs on the next real packet. - #[tokio::test] - async fn resync_after_false_header_match_with_implausible_size() { - // Arrange - - let payload_variant_1 = - protobufs::from_radio::PayloadVariant::MyInfo(protobufs::MyNodeInfo::default()); - let payload_variant_2 = - protobufs::from_radio::PayloadVariant::MyInfo(protobufs::MyNodeInfo::default()); - - let (packet_1, packet_data_1) = mock_encoded_from_radio_packet(payload_variant_1, None); - let (packet_2, packet_data_2) = mock_encoded_from_radio_packet(payload_variant_2, None); - - let encoded_packet_1 = format_data_packet(packet_data_1.into()).unwrap(); - let encoded_packet_2 = format_data_packet(packet_data_2.into()).unwrap(); - - // Header with a declared size (big-endian MSB/LSB) larger than MAX_TO_FROM_RADIO_SIZE. - let false_header = vec![0x94, 0xc3, 254, 6]; - - let (mock_tx, mut mock_rx) = unbounded_channel::(); - - // Act - - let mut buffer = StreamBuffer::new(mock_tx); - - let mut data = Vec::new(); - data.extend_from_slice(encoded_packet_1.data()); - data.extend_from_slice(&false_header); - data.extend_from_slice(encoded_packet_2.data()); - buffer.process_incoming_bytes(data.into()); - - // Assert - - assert_eq!(timeout_test(mock_rx.recv(), None).await, Some(packet_1)); - assert_eq!(timeout_test(mock_rx.recv(), None).await, Some(packet_2)); - assert_eq!(buffer.buffer.len(), 0); - } - /// `outer_packet`'s serialized bytes coincidentally contain [0x94, 0xc3], - /// the same collision that caused a real hang before MAX_TO_FROM_RADIO_SIZE - /// was added. Identity fields are placeholders; device_metrics is unchanged + /// the same collision that caused a real hang. Identity fields are + /// placeholders; device_metrics is unchanged /// since that's what produces the collision. Expected behavior is that both /// packets decode, since the coincidental match doesn't mean outer_packet /// is malformed. From 794e1e16bff703c0f2e40027adcd44730550f9f5 Mon Sep 17 00:00:00 2001 From: projectmoon Date: Wed, 5 Aug 2026 00:38:26 +0200 Subject: [PATCH 05/11] cargo format again --- src/connections/stream_buffer.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index babe12c..8b94730 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -762,7 +762,10 @@ mod tests { }; let encoded_outer = format_data_packet(outer_packet.encode_to_vec().into()).unwrap(); assert!( - encoded_outer.data_vec().windows(2).any(|w| w == [0x94, 0xc3]), + encoded_outer + .data_vec() + .windows(2) + .any(|w| w == [0x94, 0xc3]), "outer_packet no longer reproduces the coincidental header collision" ); From 6e7e236173d80f8385f93f4ab42025ea88453d3e Mon Sep 17 00:00:00 2001 From: ProjectMoon Date: Wed, 5 Aug 2026 19:29:01 +0200 Subject: [PATCH 06/11] check outer packet magic bytes only after the legit magic bytes. Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- src/connections/stream_buffer.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 8b94730..ef5586a 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -764,6 +764,7 @@ mod tests { assert!( encoded_outer .data_vec() + [PACKET_HEADER_SIZE..] .windows(2) .any(|w| w == [0x94, 0xc3]), "outer_packet no longer reproduces the coincidental header collision" From 6db5f7cb7cdf0da8a5f9fcb0fcde77c85b2e05ff Mon Sep 17 00:00:00 2001 From: projectmoon Date: Wed, 5 Aug 2026 21:15:47 +0200 Subject: [PATCH 07/11] format rabbit changes --- src/connections/stream_buffer.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index ef5586a..3753575 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -762,9 +762,7 @@ mod tests { }; let encoded_outer = format_data_packet(outer_packet.encode_to_vec().into()).unwrap(); assert!( - encoded_outer - .data_vec() - [PACKET_HEADER_SIZE..] + encoded_outer.data_vec()[PACKET_HEADER_SIZE..] .windows(2) .any(|w| w == [0x94, 0xc3]), "outer_packet no longer reproduces the coincidental header collision" From a427bcec3893cc9b0f73388db1bf1dfec7e0070e Mon Sep 17 00:00:00 2001 From: ProjectMoon Date: Mon, 17 Aug 2026 10:11:36 +0200 Subject: [PATCH 08/11] More succinct reproduction test packet. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Lukáš Poláček --- src/connections/stream_buffer.rs | 19 ++----------------- 1 file changed, 2 insertions(+), 17 deletions(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 3753575..6e5c041 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -734,31 +734,16 @@ mod tests { #[tokio::test] async fn resync_after_coincidental_header_match_inside_packet_payload() { let outer_packet = protobufs::FromRadio { - id: 0, payload_variant: Some(protobufs::from_radio::PayloadVariant::NodeInfo( protobufs::NodeInfo { - num: 1234567890, - user: Some(protobufs::User { - id: "!deadbeef".to_string(), - long_name: "Test Node".to_string(), - short_name: "TEST".to_string(), - hw_model: protobufs::HardwareModel::HeltecV3 as i32, - public_key: vec![0xab; 32], - ..Default::default() - }), device_metrics: Some(protobufs::DeviceMetrics { - battery_level: Some(100), - voltage: Some(4.219), - channel_utilization: Some(22.93), - air_util_tx: Some(2.791), uptime_seconds: Some(14655892), + ..Default::default() }), - channel: 1, - via_mqtt: true, - hops_away: Some(1), ..Default::default() }, )), + ..Default::default() }; let encoded_outer = format_data_packet(outer_packet.encode_to_vec().into()).unwrap(); assert!( From f431467856976accd8eeabccbc30347af9f1be7b Mon Sep 17 00:00:00 2001 From: ProjectMoon Date: Mon, 17 Aug 2026 10:11:56 +0200 Subject: [PATCH 09/11] Shorter-er protos. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Lukáš Poláček --- src/connections/stream_buffer.rs | 11 ----------- 1 file changed, 11 deletions(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 6e5c041..2f33093 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -758,17 +758,6 @@ mod tests { payload_variant: Some(protobufs::from_radio::PayloadVariant::NodeInfo( protobufs::NodeInfo { num: 987654321, - user: Some(protobufs::User { - id: "!cafef00d".to_string(), - long_name: "Recovery Node".to_string(), - short_name: "RCVR".to_string(), - hw_model: protobufs::HardwareModel::HeltecV3 as i32, - public_key: vec![0xcd; 32], - ..Default::default() - }), - channel: 1, - via_mqtt: false, - hops_away: Some(2), ..Default::default() }, )), From 30caf872c23d3b1f2f60e22b92f51f678d968bb7 Mon Sep 17 00:00:00 2001 From: ProjectMoon Date: Mon, 17 Aug 2026 10:12:12 +0200 Subject: [PATCH 10/11] Check if whole packet parsed. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Lukáš Poláček --- src/connections/stream_buffer.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 2f33093..1ff2019 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -324,9 +324,10 @@ impl StreamBuffer { if let Some(next_packet_start_idx) = next_packet_start_index { // Only malformed if it doesn't decode at the declared length. - let declared_packet = + let mut declared_packet = &self.buffer[packet_data_start_index..packet_data_start_index + packet_data_size]; - if protobufs::FromRadio::decode(declared_packet).is_ok() { + if protobufs::FromRadio::decode(&mut declared_packet).is_ok() + && declared_packet.is_empty() return Ok(()); } From fe39cfd114de7a7419518d1cda27efff1ad8ea64 Mon Sep 17 00:00:00 2001 From: projectmoon Date: Mon, 17 Aug 2026 10:15:11 +0200 Subject: [PATCH 11/11] syntax fix (missing } ). --- src/connections/stream_buffer.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/connections/stream_buffer.rs b/src/connections/stream_buffer.rs index 1ff2019..eaabc4d 100644 --- a/src/connections/stream_buffer.rs +++ b/src/connections/stream_buffer.rs @@ -328,6 +328,7 @@ impl StreamBuffer { &self.buffer[packet_data_start_index..packet_data_start_index + packet_data_size]; if protobufs::FromRadio::decode(&mut declared_packet).is_ok() && declared_packet.is_empty() + { return Ok(()); }