Skip to main content

pamoja_mavlink/
parser.rs

1//! A streaming frame parser that turns a byte stream into whole frames.
2//!
3//! A link delivers bytes, not frames: a serial port hands over whatever arrived since the
4//! last read, and even a UDP datagram can carry several frames back to back. This parser
5//! bridges that gap. Bytes are fed in as they arrive, and a complete, checksum-verified
6//! [`Frame`] is returned as soon as one is recognized. A stray byte or a frame mangled in
7//! transit makes the parser resynchronize on the next start marker rather than wedge, so a
8//! noisy link recovers on its own. It holds a single fixed buffer, so it runs unchanged on
9//! a microcontroller.
10
11use crate::frame::{Frame, MAGIC_V1, MAGIC_V2, MAX_FRAME};
12
13// What to do given the bytes buffered so far.
14enum Decision {
15    // Not enough bytes yet to decide; keep accumulating.
16    NeedMore,
17    // The buffer does not start a valid frame; drop the leading byte and re-examine.
18    Resync,
19    // A complete, valid frame occupies the first `usize` bytes of the buffer.
20    Frame(usize),
21}
22
23/// Accumulates bytes from a link and emits complete frames.
24///
25/// Feed bytes with [`push_byte`](Parser::push_byte); each call returns a [`Frame`] on the
26/// byte that completes one. Resolving a message id to its `CRC_EXTRA` is left to the
27/// caller, so the parser works against any dialect.
28#[derive(Clone)]
29pub struct Parser {
30    buf: [u8; MAX_FRAME],
31    len: usize,
32}
33
34impl Default for Parser {
35    fn default() -> Self {
36        Self::new()
37    }
38}
39
40impl Parser {
41    /// Creates an empty parser.
42    ///
43    /// # Returns
44    ///
45    /// A parser with nothing buffered.
46    pub fn new() -> Self {
47        Parser {
48            buf: [0u8; MAX_FRAME],
49            len: 0,
50        }
51    }
52
53    /// Feeds one byte, returning a frame if this byte completes a valid one.
54    ///
55    /// A frame whose checksum fails, or whose message id `crc_extra_for` does not
56    /// recognize, is discarded and the parser resynchronizes on the next start marker.
57    ///
58    /// # Arguments
59    ///
60    /// * `byte` - the next byte from the link.
61    /// * `crc_extra_for` - resolves a message id to its `CRC_EXTRA`, or `None` if unknown.
62    ///
63    /// # Returns
64    ///
65    /// The completed frame, or [`None`] if more bytes are needed.
66    pub fn push_byte<F>(&mut self, byte: u8, crc_extra_for: &F) -> Option<Frame>
67    where
68        F: Fn(u32) -> Option<u8>,
69    {
70        if self.len >= MAX_FRAME {
71            // A frame can never exceed the buffer; an overlong run of bytes that never
72            // resolved is noise, so start fresh.
73            self.len = 0;
74        }
75        self.buf[self.len] = byte;
76        self.len += 1;
77
78        loop {
79            match self.decide(crc_extra_for) {
80                Decision::NeedMore => return None,
81                Decision::Resync => {
82                    self.drop_front(1);
83                    if self.len == 0 {
84                        return None;
85                    }
86                }
87                Decision::Frame(total) => {
88                    let frame = Frame::parse_with(&self.buf[..total], crc_extra_for)
89                        .expect("decide only reports a frame the checksum accepted");
90                    self.drop_front(total);
91                    return Some(frame);
92                }
93            }
94        }
95    }
96
97    fn decide<F>(&self, crc_extra_for: &F) -> Decision
98    where
99        F: Fn(u32) -> Option<u8>,
100    {
101        if self.len == 0 {
102            return Decision::NeedMore;
103        }
104        let header_len = match self.buf[0] {
105            MAGIC_V1 => 6,
106            MAGIC_V2 => 10,
107            _ => return Decision::Resync,
108        };
109        // The length byte, and for v2 the incompat flags, decide the total frame size.
110        let need_for_size = if self.buf[0] == MAGIC_V2 { 3 } else { 2 };
111        if self.len < need_for_size {
112            return Decision::NeedMore;
113        }
114        let plen = self.buf[1] as usize;
115        let signed = self.buf[0] == MAGIC_V2 && self.buf[2] & 0x01 != 0;
116        let total = header_len + plen + 2 + if signed { 13 } else { 0 };
117        if self.len < total {
118            return Decision::NeedMore;
119        }
120        match Frame::parse_with(&self.buf[..total], crc_extra_for) {
121            Ok(_) => Decision::Frame(total),
122            Err(_) => Decision::Resync,
123        }
124    }
125
126    fn drop_front(&mut self, n: usize) {
127        let n = n.min(self.len);
128        self.buf.copy_within(n..self.len, 0);
129        self.len -= n;
130    }
131}
132
133#[cfg(test)]
134mod tests {
135    use super::*;
136    use crate::frame::{Frame, Header};
137
138    // The common dialect's CRC_EXTRA for HEARTBEAT.
139    fn heartbeat_crc(id: u32) -> Option<u8> {
140        (id == 0).then_some(50)
141    }
142
143    fn heartbeat_frame(seq: u8) -> Frame {
144        Frame::encode_v2(Header::new(1, 1, seq), 0, &[0, 0, 0, 0, 6, 8, 0, 3, 3], 50).unwrap()
145    }
146
147    fn feed(parser: &mut Parser, bytes: &[u8]) -> Vec<Frame> {
148        let mut out = Vec::new();
149        for &byte in bytes {
150            if let Some(frame) = parser.push_byte(byte, &heartbeat_crc) {
151                out.push(frame);
152            }
153        }
154        out
155    }
156
157    #[test]
158    fn a_whole_frame_is_parsed() {
159        let mut parser = Parser::new();
160        let frame = heartbeat_frame(1);
161        let parsed = feed(&mut parser, frame.as_bytes());
162        assert_eq!(parsed.len(), 1);
163        assert_eq!(parsed[0].sequence(), 1);
164    }
165
166    #[test]
167    fn a_frame_split_across_feeds_is_parsed() {
168        let mut parser = Parser::new();
169        let frame = heartbeat_frame(2);
170        let bytes = frame.as_bytes();
171        let (head, tail) = bytes.split_at(4);
172        assert!(feed(&mut parser, head).is_empty());
173        let parsed = feed(&mut parser, tail);
174        assert_eq!(parsed.len(), 1);
175        assert_eq!(parsed[0].sequence(), 2);
176    }
177
178    #[test]
179    fn garbage_before_a_frame_is_skipped() {
180        let mut parser = Parser::new();
181        let frame = heartbeat_frame(3);
182        let mut stream = vec![0x00, 0xFF, 0x12, 0x34]; // leading noise, no start marker
183        stream.extend_from_slice(frame.as_bytes());
184        let parsed = feed(&mut parser, &stream);
185        assert_eq!(parsed.len(), 1);
186        assert_eq!(parsed[0].sequence(), 3);
187    }
188
189    #[test]
190    fn two_back_to_back_frames_both_emit() {
191        let mut parser = Parser::new();
192        let mut stream = Vec::new();
193        stream.extend_from_slice(heartbeat_frame(10).as_bytes());
194        stream.extend_from_slice(heartbeat_frame(11).as_bytes());
195        let parsed = feed(&mut parser, &stream);
196        assert_eq!(parsed.len(), 2);
197        assert_eq!(parsed[0].sequence(), 10);
198        assert_eq!(parsed[1].sequence(), 11);
199    }
200
201    #[test]
202    fn a_corrupt_frame_is_dropped_and_the_next_recovers() {
203        let mut parser = Parser::new();
204        let mut corrupt = heartbeat_frame(20).as_bytes().to_vec();
205        let last = corrupt.len() - 1;
206        corrupt[last] ^= 0xFF; // break the checksum
207        let mut stream = corrupt;
208        stream.extend_from_slice(heartbeat_frame(21).as_bytes());
209        let parsed = feed(&mut parser, &stream);
210        assert_eq!(parsed.len(), 1);
211        assert_eq!(parsed[0].sequence(), 21);
212    }
213}