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}