Skip to main content

mavlink_core/
async_reader.rs

1//! Incremental asynchronous MAVLink reader.
2
3use crate::{
4    MAVLinkMessageRaw, MavHeader, MavlinkVersion, Message,
5    error::MessageReadError,
6    frame_decoder::{FrameDecoder, VersionFilter},
7    reader::{try_decode_message, try_decode_raw_message},
8};
9
10#[cfg(feature = "tokio")]
11use crate::SigningData;
12
13#[cfg(all(feature = "embedded", not(feature = "std")))]
14use embedded_io_async::Read;
15#[cfg(feature = "tokio")]
16use tokio::io::{AsyncRead, AsyncReadExt};
17
18/// Incrementally reads complete MAVLink frames from an asynchronous byte stream.
19///
20/// Partial input remains buffered if a read is cancelled. Use the same MAVLink
21/// dialect `M` for the lifetime of a stream, because CRC validation is
22/// dialect-specific.
23pub struct AsyncMavlinkReader<R> {
24    source: R,
25    decoder: FrameDecoder,
26}
27
28impl<R> AsyncMavlinkReader<R> {
29    /// Creates a reader around `source` with the default read-ahead capacity.
30    ///
31    /// The buffer is allocated once and reused for the reader's lifetime.
32    #[cfg(feature = "std")]
33    pub fn new(source: R) -> Self
34    where
35        R: AsyncRead + Unpin,
36    {
37        Self {
38            source,
39            decoder: FrameDecoder::new(),
40        }
41    }
42
43    /// Creates an allocation-free reader around `source`.
44    #[cfg(not(feature = "std"))]
45    pub const fn new(source: R) -> Self
46    where
47        R: Read,
48    {
49        Self {
50            source,
51            decoder: FrameDecoder::new(),
52        }
53    }
54
55    /// Creates a reader with at least `capacity` bytes of read-ahead space.
56    ///
57    /// Capacities smaller than one maximum-size MAVLink frame are raised to
58    /// [`consts::MAX_FRAME_SIZE`](crate::consts::MAX_FRAME_SIZE). The buffer is
59    /// allocated once and reused for the lifetime of the reader.
60    #[cfg(feature = "std")]
61    pub fn with_capacity(capacity: usize, source: R) -> Self
62    where
63        R: AsyncRead + Unpin,
64    {
65        Self {
66            source,
67            decoder: FrameDecoder::with_capacity(capacity),
68        }
69    }
70
71    /// Returns a shared reference to the underlying source.
72    pub const fn get_ref(&self) -> &R {
73        &self.source
74    }
75
76    /// Returns a mutable reference to the underlying source.
77    ///
78    /// Reading directly from the source can skip bytes already buffered by
79    /// this reader and should therefore be avoided.
80    pub const fn get_mut(&mut self) -> &mut R {
81        &mut self.source
82    }
83
84    /// Returns the underlying source.
85    ///
86    /// Any bytes read ahead by this reader are discarded.
87    pub fn into_inner(self) -> R {
88        self.source
89    }
90}
91
92#[cfg(feature = "tokio")]
93impl<R: AsyncRead + Unpin> AsyncMavlinkReader<R> {
94    /// Reads and parses the next CRC-valid message accepted by `version`.
95    pub async fn read_message<M: Message>(
96        &mut self,
97        version: MavlinkVersion,
98    ) -> Result<(MavHeader, M), MessageReadError> {
99        self.read_message_inner(VersionFilter::Exact(version), None)
100            .await
101    }
102
103    /// Reads and parses the next CRC-valid MAVLink 1 or MAVLink 2 message.
104    pub async fn read_any_message<M: Message>(
105        &mut self,
106    ) -> Result<(MavHeader, M), MessageReadError> {
107        self.read_message_inner(VersionFilter::Any, None).await
108    }
109
110    /// Reads the next CRC-valid raw message accepted by `version`.
111    pub async fn read_raw_message<M: Message>(
112        &mut self,
113        version: MavlinkVersion,
114    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
115        self.read_raw_message_inner::<M>(VersionFilter::Exact(version), None)
116            .await
117    }
118
119    /// Reads the next CRC-valid MAVLink 1 or MAVLink 2 raw message.
120    pub async fn read_any_raw_message<M: Message>(
121        &mut self,
122    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
123        self.read_raw_message_inner::<M>(VersionFilter::Any, None)
124            .await
125    }
126
127    /// Reads, verifies, and parses the next message accepted by `version`.
128    /// With signing data, MAVLink 1 and unsigned MAVLink 2 frames require
129    /// `SigningConfig::allow_unsigned`; signed MAVLink 2 frames must have a valid
130    /// signature. With `None`, signature verification is skipped.
131    #[cfg(feature = "mav2-message-signing")]
132    pub async fn read_message_signed<M: Message>(
133        &mut self,
134        version: MavlinkVersion,
135        signing_data: Option<&SigningData>,
136    ) -> Result<(MavHeader, M), MessageReadError> {
137        self.read_message_inner(VersionFilter::Exact(version), signing_data)
138            .await
139    }
140
141    /// Reads, verifies, and parses the next MAVLink 1 or MAVLink 2 message.
142    /// With signing data, MAVLink 1 frames follow `SigningConfig::allow_unsigned`.
143    /// With `None`, signature verification is skipped.
144    #[cfg(feature = "mav2-message-signing")]
145    pub async fn read_any_message_signed<M: Message>(
146        &mut self,
147        signing_data: Option<&SigningData>,
148    ) -> Result<(MavHeader, M), MessageReadError> {
149        self.read_message_inner(VersionFilter::Any, signing_data)
150            .await
151    }
152
153    /// Reads and verifies the next raw message accepted by `version`.
154    /// With signing data, MAVLink 1 and unsigned MAVLink 2 frames require
155    /// `SigningConfig::allow_unsigned`; signed MAVLink 2 frames must have a valid
156    /// signature. With `None`, signature verification is skipped.
157    #[cfg(feature = "mav2-message-signing")]
158    pub async fn read_raw_message_signed<M: Message>(
159        &mut self,
160        version: MavlinkVersion,
161        signing_data: Option<&SigningData>,
162    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
163        self.read_raw_message_inner::<M>(VersionFilter::Exact(version), signing_data)
164            .await
165    }
166
167    /// Reads and verifies the next MAVLink 1 or MAVLink 2 raw message.
168    /// With signing data, MAVLink 1 frames follow `SigningConfig::allow_unsigned`.
169    /// With `None`, signature verification is skipped.
170    #[cfg(feature = "mav2-message-signing")]
171    pub async fn read_any_raw_message_signed<M: Message>(
172        &mut self,
173        signing_data: Option<&SigningData>,
174    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
175        self.read_raw_message_inner::<M>(VersionFilter::Any, signing_data)
176            .await
177    }
178
179    #[inline]
180    pub(crate) async fn read_message_inner<M: Message>(
181        &mut self,
182        filter: VersionFilter,
183        signing_data: Option<&SigningData>,
184    ) -> Result<(MavHeader, M), MessageReadError> {
185        loop {
186            if let Some(result) = try_decode_message::<M>(&mut self.decoder, filter, signing_data) {
187                return result;
188            }
189            self.read_more().await?;
190        }
191    }
192
193    #[inline]
194    pub(crate) async fn read_raw_message_inner<M: Message>(
195        &mut self,
196        filter: VersionFilter,
197        signing_data: Option<&SigningData>,
198    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
199        loop {
200            if let Some(message) =
201                try_decode_raw_message::<M>(&mut self.decoder, filter, signing_data)
202            {
203                return Ok(message);
204            }
205            self.read_more().await?;
206        }
207    }
208
209    async fn read_more(&mut self) -> Result<(), MessageReadError> {
210        let destination = self.decoder.spare_capacity_mut();
211        assert!(
212            !destination.is_empty(),
213            "a pending MAVLink frame must fit in the decoder buffer"
214        );
215        let count = loop {
216            match self.source.read(destination).await {
217                Err(error) if error.kind() == std::io::ErrorKind::Interrupted => continue,
218                result => break result?,
219            }
220        };
221        if count == 0 {
222            return Err(MessageReadError::eof());
223        }
224        self.decoder.commit(count);
225        Ok(())
226    }
227}
228
229#[cfg(all(feature = "embedded", not(feature = "std")))]
230impl<R: Read> AsyncMavlinkReader<R> {
231    /// Reads and parses the next CRC-valid message accepted by `version`.
232    pub async fn read_message<M: Message>(
233        &mut self,
234        version: MavlinkVersion,
235    ) -> Result<(MavHeader, M), MessageReadError> {
236        let filter = VersionFilter::Exact(version);
237        loop {
238            if let Some(result) = try_decode_message::<M>(&mut self.decoder, filter, None) {
239                return result;
240            }
241            self.read_more().await?;
242        }
243    }
244
245    /// Reads and parses the next CRC-valid MAVLink 1 or MAVLink 2 message.
246    pub async fn read_any_message<M: Message>(
247        &mut self,
248    ) -> Result<(MavHeader, M), MessageReadError> {
249        loop {
250            if let Some(result) =
251                try_decode_message::<M>(&mut self.decoder, VersionFilter::Any, None)
252            {
253                return result;
254            }
255            self.read_more().await?;
256        }
257    }
258
259    /// Reads the next CRC-valid raw message accepted by `version`.
260    pub async fn read_raw_message<M: Message>(
261        &mut self,
262        version: MavlinkVersion,
263    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
264        let filter = VersionFilter::Exact(version);
265        loop {
266            if let Some(message) = try_decode_raw_message::<M>(&mut self.decoder, filter, None) {
267                return Ok(message);
268            }
269            self.read_more().await?;
270        }
271    }
272
273    /// Reads the next CRC-valid MAVLink 1 or MAVLink 2 raw message.
274    pub async fn read_any_raw_message<M: Message>(
275        &mut self,
276    ) -> Result<MAVLinkMessageRaw, MessageReadError> {
277        loop {
278            if let Some(message) =
279                try_decode_raw_message::<M>(&mut self.decoder, VersionFilter::Any, None)
280            {
281                return Ok(message);
282            }
283            self.read_more().await?;
284        }
285    }
286
287    async fn read_more(&mut self) -> Result<(), MessageReadError> {
288        let destination = self.decoder.spare_capacity_mut();
289        assert!(
290            !destination.is_empty(),
291            "a pending MAVLink frame must fit in the decoder buffer"
292        );
293        let count = self
294            .source
295            .read(destination)
296            .await
297            .map_err(|_| MessageReadError::Io)?;
298        if count == 0 {
299            return Err(MessageReadError::eof());
300        }
301        self.decoder.commit(count);
302        Ok(())
303    }
304}