1use 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
18pub struct AsyncMavlinkReader<R> {
24 source: R,
25 decoder: FrameDecoder,
26}
27
28impl<R> AsyncMavlinkReader<R> {
29 #[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 #[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 #[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 pub const fn get_ref(&self) -> &R {
73 &self.source
74 }
75
76 pub const fn get_mut(&mut self) -> &mut R {
81 &mut self.source
82 }
83
84 pub fn into_inner(self) -> R {
88 self.source
89 }
90}
91
92#[cfg(feature = "tokio")]
93impl<R: AsyncRead + Unpin> AsyncMavlinkReader<R> {
94 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 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 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 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 #[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 #[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 #[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 #[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 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 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 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 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}