Skip to main content

android_abx/decode/
stream.rs

1//! Streaming parser over any [`Read`] source.
2
3use std::io::Read;
4
5use winnow::{
6    Partial,
7    error::{ErrMode, Needed},
8    stream::{Offset, StreamIsPartial},
9};
10
11use super::grammar;
12use crate::{Attribute, AttributeValue, Event, Result, render_event};
13
14use std::collections::HashMap;
15
16const INITIAL_BUF: usize = 4096;
17const READ_CHUNK: usize = 4096;
18
19/// A pull parser that reads an ABX document from any [`Read`] source.
20///
21/// The input is read in 4 KiB chunks into an internal buffer, which only grows
22/// to fit a single value larger than that. Same methods as
23/// [`AbxParser`](crate::AbxParser); it also implements [`Iterator`], yielding
24/// `Result<Event>`.
25///
26/// # Examples
27///
28/// ```
29/// use android_abx::{AbxStreamParser, Event};
30///
31/// # let data = include_bytes!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/simple_pkg.abx"));
32/// let parser = AbxStreamParser::new(&data[..])?;
33/// let tags = parser
34///     .filter(|ev| matches!(ev, Ok(Event::StartTag { .. })))
35///     .count();
36/// assert_eq!(tags, 1);
37/// # Ok::<(), android_abx::AbxError>(())
38/// ```
39#[derive(Debug)]
40pub struct AbxStreamParser<R: Read> {
41    reader: R,
42    buf: Vec<u8>,
43    pos: usize,
44    len: usize,
45    eof: bool,
46    pool: Vec<crate::InternedStr>,
47}
48
49impl<R: Read> AbxStreamParser<R> {
50    /// Creates a parser and reads the header from `reader`.
51    ///
52    /// # Errors
53    ///
54    /// Returns [`AbxError::Io`](crate::AbxError::Io) if reading fails, [`AbxError::UnexpectedEof`](crate::AbxError::UnexpectedEof) if
55    /// the input is shorter than 4 bytes, or [`AbxError::InvalidMagic`](crate::AbxError::InvalidMagic) if it does
56    /// not start with [`MAGIC`](crate::MAGIC).
57    pub fn new(reader: R) -> Result<Self> {
58        let mut p = AbxStreamParser {
59            reader,
60            buf: vec![0u8; INITIAL_BUF],
61            pos: 0,
62            len: 0,
63            eof: false,
64            pool: Vec::with_capacity(32),
65        };
66
67        p.ensure(4)?;
68        crate::decode::check_magic(&p.buf[p.pos..p.len])?;
69        p.pos += 4;
70        Ok(p)
71    }
72
73    #[inline]
74    fn available(&self) -> usize {
75        self.len - self.pos
76    }
77
78    /// Compact, then read until `needed` bytes are available or EOF.
79    fn ensure(&mut self, needed: usize) -> Result<()> {
80        // Hot path: skip the compaction memmove when no refill is needed.
81        if self.available() >= needed || self.eof {
82            return Ok(());
83        }
84
85        if self.pos > 0 {
86            self.buf.copy_within(self.pos..self.len, 0);
87            self.len -= self.pos;
88            self.pos = 0;
89        }
90
91        while self.available() < needed && !self.eof {
92            let spare = self.buf.len() - self.len;
93            if spare < READ_CHUNK {
94                self.buf
95                    .resize(self.len + READ_CHUNK.max(needed - self.available()), 0);
96            }
97
98            let n = self.reader.read(&mut self.buf[self.len..])?;
99            if n == 0 {
100                self.eof = true;
101            } else {
102                self.len += n;
103            }
104        }
105
106        Ok(())
107    }
108
109    /// Reads the next event, or returns `None` at the end of the input.
110    ///
111    /// # Errors
112    ///
113    /// Returns an error if the input is truncated or malformed (unknown token,
114    /// invalid interned-string index, invalid UTF-8).
115    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
116    /// On error, the parser is left at the start of the failing event.
117    pub fn next_event(&mut self) -> Result<Option<Event>> {
118        self.ensure(1)?;
119        if self.available() == 0 {
120            return Ok(None);
121        }
122        let pool_len = self.pool.len();
123        let result = self.parse_event();
124        if result.is_err() {
125            self.pool.truncate(pool_len);
126        }
127        result.map(Some)
128    }
129
130    /// Parses one event from the buffer, refilling and retrying from its start on `Incomplete`.
131    fn parse_event(&mut self) -> Result<Event> {
132        loop {
133            let window = &self.buf[self.pos..self.len];
134            let mut input = Partial::new(window);
135            if self.eof {
136                let _ = input.complete();
137            }
138            let pool_len = self.pool.len();
139            match grammar::event(&mut input, &mut self.pool) {
140                Ok(ev) => {
141                    self.pos += input.offset_from(&Partial::new(window));
142                    return Ok(ev);
143                }
144                Err(ErrMode::Incomplete(needed)) => {
145                    self.pool.truncate(pool_len);
146                    let more = match needed {
147                        Needed::Size(n) => n.get(),
148                        Needed::Unknown => 1,
149                    };
150                    self.ensure(self.available() + more)?;
151                }
152                Err(e) => return Err(grammar::into_abx_error(e)),
153            }
154        }
155    }
156
157    /// Reads all remaining events into a `Vec`.
158    ///
159    /// # Errors
160    ///
161    /// Returns an error if the input is truncated or malformed (unknown token,
162    /// invalid interned-string index, invalid UTF-8).
163    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
164    pub fn collect_events(&mut self) -> Result<Vec<Event>> {
165        let mut out = Vec::new();
166        while let Some(ev) = self.next_event()? {
167            out.push(ev);
168        }
169        Ok(out)
170    }
171
172    /// Returns the value of attribute `attr` on the next `<element>` tag that has it.
173    ///
174    /// Events are consumed up to and including the matching tag. Returns `Ok(None)`
175    /// if the document ends first.
176    ///
177    /// # Errors
178    ///
179    /// Returns an error if the input is truncated or malformed (unknown token,
180    /// invalid interned-string index, invalid UTF-8).
181    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
182    ///
183    /// # Examples
184    ///
185    /// ```
186    /// use android_abx::{AbxStreamParser, AttributeValue};
187    ///
188    /// # let data = include_bytes!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/simple_pkg.abx"));
189    /// let mut parser = AbxStreamParser::new(&data[..])?;
190    /// let name = parser.find_attribute("pkg", "name")?;
191    /// assert_eq!(name, Some(AttributeValue::String("com.example.chat".into())));
192    /// # Ok::<(), android_abx::AbxError>(())
193    /// ```
194    pub fn find_attribute(&mut self, element: &str, attr: &str) -> Result<Option<AttributeValue>> {
195        loop {
196            match self.next_event()? {
197                Some(Event::StartTag { name, attributes }) if name == element => {
198                    if let Some(a) = attributes.into_iter().find(|a| a.name == attr) {
199                        return Ok(Some(a.value));
200                    }
201                }
202                Some(Event::EndDocument) | None => return Ok(None),
203                _ => {}
204            }
205        }
206    }
207
208    /// Returns the value of attribute `attr` on every remaining `<element>` tag.
209    ///
210    /// Tags without `attr` are skipped. Consumes the rest of the document.
211    ///
212    /// # Errors
213    ///
214    /// Returns an error if the input is truncated or malformed (unknown token,
215    /// invalid interned-string index, invalid UTF-8).
216    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
217    pub fn find_all_attributes(
218        &mut self,
219        element: &str,
220        attr: &str,
221    ) -> Result<Vec<AttributeValue>> {
222        let mut out = Vec::new();
223        while let Some(ev) = self.next_event()? {
224            if let Event::StartTag { name, attributes } = ev
225                && name == element
226            {
227                out.extend(
228                    attributes
229                        .into_iter()
230                        .filter(|a| a.name == attr)
231                        .map(|a| a.value),
232                );
233            }
234        }
235        Ok(out)
236    }
237
238    /// Returns the attributes of the next `<element>` tag.
239    ///
240    /// Events are consumed up to and including the matching tag. Returns `Ok(None)`
241    /// if the document ends first.
242    ///
243    /// # Errors
244    ///
245    /// Returns an error if the input is truncated or malformed (unknown token,
246    /// invalid interned-string index, invalid UTF-8).
247    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
248    pub fn attributes_of(&mut self, element: &str) -> Result<Option<Vec<Attribute>>> {
249        loop {
250            match self.next_event()? {
251                Some(Event::StartTag { name, attributes }) if name == element => {
252                    return Ok(Some(attributes));
253                }
254                Some(Event::EndDocument) | None => return Ok(None),
255                _ => {}
256            }
257        }
258    }
259
260    /// Returns the attributes of every remaining `<element>` tag.
261    ///
262    /// Consumes the rest of the document.
263    ///
264    /// # Errors
265    ///
266    /// Returns an error if the input is truncated or malformed (unknown token,
267    /// invalid interned-string index, invalid UTF-8).
268    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
269    pub fn all_attributes_of(&mut self, element: &str) -> Result<Vec<Vec<Attribute>>> {
270        let mut out = Vec::new();
271        while let Some(ev) = self.next_event()? {
272            if let Event::StartTag { name, attributes } = ev
273                && name == element
274            {
275                out.push(attributes);
276            }
277        }
278        Ok(out)
279    }
280
281    /// Deserializes the next `<element>` into `T`.
282    ///
283    /// Events are consumed up to and including the element's closing tag. Returns
284    /// `Ok(None)` if the document ends first. See
285    /// [Deserializing with serde](crate#deserializing-with-serde) for how fields are
286    /// matched.
287    ///
288    /// # Errors
289    ///
290    /// Returns a parse error if the input is malformed, or
291    /// [`AbxError::Deserialization`](crate::AbxError::Deserialization)(crate::AbxError::Deserialization) if the element
292    /// does not match `T`.
293    ///
294    /// # Examples
295    ///
296    /// ```
297    /// use android_abx::AbxStreamParser;
298    /// use serde::Deserialize;
299    ///
300    /// #[derive(Deserialize)]
301    /// struct Permission {
302    ///     name: String,
303    /// }
304    ///
305    /// #[derive(Deserialize)]
306    /// struct Pkg {
307    ///     name: String,
308    ///     description: String,
309    ///     permission: Vec<Permission>,
310    /// }
311    ///
312    /// # let data = include_bytes!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/nested_permissions.abx"));
313    /// let mut parser = AbxStreamParser::new(&data[..])?;
314    /// let pkg: Pkg = parser.deserialize_next("pkg")?.unwrap();
315    /// assert_eq!(pkg.name, "com.example.chat");
316    /// assert_eq!(pkg.description, "A chat app");
317    /// assert_eq!(pkg.permission.len(), 2);
318    /// # Ok::<(), android_abx::AbxError>(())
319    /// ```
320    #[cfg(feature = "serde")]
321    pub fn deserialize_next<T: serde::de::DeserializeOwned>(
322        &mut self,
323        element: &str,
324    ) -> Result<Option<T>> {
325        crate::de::find_and_consume_element(self, element)
326    }
327
328    /// Deserializes every remaining `<element>` into a `Vec<T>`.
329    ///
330    /// # Errors
331    ///
332    /// Same as [`deserialize_next`](Self::deserialize_next).
333    #[cfg(feature = "serde")]
334    pub fn deserialize_all<T: serde::de::DeserializeOwned>(
335        &mut self,
336        element: &str,
337    ) -> Result<Vec<T>> {
338        let mut out = Vec::new();
339        while let Some(item) = self.deserialize_next(element)? {
340            out.push(item);
341        }
342        Ok(out)
343    }
344
345    /// Returns an iterator that deserializes each remaining `<element>` into `T`.
346    ///
347    /// Elements are read one at a time, so memory use does not grow with the
348    /// document.
349    ///
350    /// # Examples
351    ///
352    /// ```
353    /// use android_abx::AbxStreamParser;
354    /// use serde::Deserialize;
355    ///
356    /// #[derive(Deserialize)]
357    /// struct Permission {
358    ///     name: String,
359    /// }
360    ///
361    /// # let data = include_bytes!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/nested_permissions.abx"));
362    /// let mut parser = AbxStreamParser::new(&data[..])?;
363    /// let names = parser
364    ///     .deserialize_iter::<Permission>("permission")
365    ///     .map(|p| p.map(|p| p.name))
366    ///     .collect::<Result<Vec<_>, _>>()?;
367    /// assert_eq!(names, ["INTERNET", "CAMERA"]);
368    /// # Ok::<(), android_abx::AbxError>(())
369    /// ```
370    #[cfg(feature = "serde")]
371    pub fn deserialize_iter<'p, T: serde::de::DeserializeOwned>(
372        &'p mut self,
373        element: &'p str,
374    ) -> DeserializeIter<'p, R, T> {
375        DeserializeIter {
376            parser: self,
377            element,
378            _marker: std::marker::PhantomData,
379        }
380    }
381
382    /// Renders the remaining events as an XML string.
383    ///
384    /// The output starts with an `<?xml ...?>` declaration. Text and attribute values
385    /// are escaped; empty elements are written as an opening and a closing tag.
386    ///
387    /// # Errors
388    ///
389    /// Returns an error if the input is truncated or malformed (unknown token,
390    /// invalid interned-string index, invalid UTF-8).
391    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
392    ///
393    /// # Examples
394    ///
395    /// ```
396    /// use android_abx::AbxStreamParser;
397    ///
398    /// # let data = include_bytes!(concat!(env!("CARGO_MANIFEST_DIR"), "/tests/fixtures/simple_pkg.abx"));
399    /// let xml = AbxStreamParser::new(&data[..])?.to_xml()?;
400    /// assert!(xml.ends_with(r#"<pkg name="com.example.chat" version="3" flags="1"></pkg>"#));
401    /// # Ok::<(), android_abx::AbxError>(())
402    /// ```
403    pub fn to_xml(&mut self) -> Result<String> {
404        let mut buf = String::from(r#"<?xml version="1.0" encoding="UTF-8"?>"#);
405        while let Some(ev) = self.next_event()? {
406            if matches!(ev, Event::EndDocument) {
407                break;
408            }
409            render_event(&ev, &mut buf);
410        }
411        Ok(buf)
412    }
413
414    /// Writes the remaining events as XML to `writer`.
415    ///
416    /// Same output as [`to_xml`](Self::to_xml), without building the whole string
417    /// in memory.
418    ///
419    /// # Errors
420    ///
421    /// Returns an error if the input is truncated or malformed (unknown token,
422    /// invalid interned-string index, invalid UTF-8).
423    /// Also returns [`AbxError::Io`](crate::AbxError::Io)(crate::AbxError::Io) if reading or writing fails.
424    pub fn write_xml(&mut self, writer: &mut impl std::io::Write) -> Result<()> {
425        writer.write_all(b"<?xml version=\"1.0\" encoding=\"UTF-8\"?>")?;
426        let mut tmp = String::new();
427        while let Some(ev) = self.next_event()? {
428            if matches!(ev, Event::EndDocument) {
429                break;
430            }
431            tmp.clear();
432            render_event(&ev, &mut tmp);
433            writer.write_all(tmp.as_bytes())?;
434        }
435        Ok(())
436    }
437
438    /// Collects the attributes of every remaining tag, grouped by tag name.
439    ///
440    /// Values are rendered with [`AttributeValue::as_str`](crate::AttributeValue::as_str).
441    /// Text and nesting are discarded.
442    ///
443    /// # Errors
444    ///
445    /// Returns an error if the input is truncated or malformed (unknown token,
446    /// invalid interned-string index, invalid UTF-8).
447    /// Also returns [`AbxError::Io`](crate::AbxError::Io) if reading fails.
448    pub fn into_map(mut self) -> Result<HashMap<String, Vec<HashMap<String, String>>>> {
449        let mut map: HashMap<String, Vec<HashMap<String, String>>> = HashMap::new();
450        while let Some(ev) = self.next_event()? {
451            if let Event::StartTag { name, attributes } = ev {
452                let entry = map.entry(name.into()).or_default();
453                let mut attrs = HashMap::new();
454                for attr in attributes {
455                    attrs.insert(attr.name.into(), attr.value.as_str().into_owned());
456                }
457                entry.push(attrs);
458            }
459        }
460        Ok(map)
461    }
462
463    /// Returns the underlying reader. Bytes already buffered but not parsed are lost.
464    pub fn into_inner(self) -> R {
465        self.reader
466    }
467}
468
469impl<R: Read> Iterator for AbxStreamParser<R> {
470    type Item = Result<Event>;
471
472    fn next(&mut self) -> Option<Self::Item> {
473        match self.next_event() {
474            Ok(Some(ev)) => Some(Ok(ev)),
475            Ok(None) => None,
476            Err(e) => Some(Err(e)),
477        }
478    }
479}
480
481/// Iterator returned by [`AbxStreamParser::deserialize_iter`].
482///
483/// Yields `Result<T>` for each matching element.
484#[cfg(feature = "serde")]
485pub struct DeserializeIter<'p, R: Read, T> {
486    parser: &'p mut AbxStreamParser<R>,
487    element: &'p str,
488    _marker: std::marker::PhantomData<T>,
489}
490
491#[cfg(feature = "serde")]
492impl<'p, R: Read, T: serde::de::DeserializeOwned> Iterator for DeserializeIter<'p, R, T> {
493    type Item = Result<T>;
494
495    fn next(&mut self) -> Option<Self::Item> {
496        self.parser.deserialize_next(self.element).transpose()
497    }
498}