Skip to main content

ogg_core/
demux.rs

1//! Ogg page reader — incremental, byte-chunk push/poll. Reassembles packets
2//! that span continuation pages and pages carrying multiple packets (the
3//! general case any real Ogg encoder produces, even though this crate's own
4//! [`crate::Muxer`] only ever emits the simpler one-packet-per-page form).
5
6#![forbid(unsafe_code)]
7
8use std::collections::VecDeque;
9
10use bytes::Bytes;
11
12use crate::crc::crc32_ogg;
13use crate::error::Error;
14use crate::types::Packet;
15
16const HEADER_LEN: usize = 27; // capture(4) + version(1) + flags(1) + granule(8) + serial(4) + sequence(4) + crc(4) + page_segments(1)
17const CRC_FIELD_OFFSET: usize = 4 + 1 + 1 + 8 + 4 + 4;
18
19/// Reads Ogg pages from pushed byte chunks and yields fully reassembled packets.
20#[derive(Debug, Default)]
21pub struct Demuxer {
22    buf: Vec<u8>,
23    partial: Vec<u8>,
24    has_partial: bool,
25    ready: VecDeque<Packet>,
26}
27
28impl Demuxer {
29    /// New, empty demux session.
30    #[must_use]
31    pub fn new() -> Self {
32        Self::default()
33    }
34
35    /// Append incoming bytes.
36    pub fn push_bytes(&mut self, data: &[u8]) {
37        self.buf.extend_from_slice(data);
38    }
39
40    /// Pop the next fully reassembled packet, or `Ok(None)` if no more pages are
41    /// buffered yet — call again after more `push_bytes`.
42    pub fn poll_packet(&mut self) -> Result<Option<Packet>, Error> {
43        loop {
44            if let Some(packet) = self.ready.pop_front() {
45                return Ok(Some(packet));
46            }
47            if !self.parse_one_page()? {
48                return Ok(None);
49            }
50        }
51    }
52
53    /// Try to parse and consume one page from `self.buf`. Returns `Ok(true)` if a
54    /// page was parsed (queuing 0+ packets into `self.ready`), `Ok(false)` if not
55    /// enough bytes are buffered yet.
56    fn parse_one_page(&mut self) -> Result<bool, Error> {
57        if self.buf.len() < HEADER_LEN {
58            return Ok(false);
59        }
60        if &self.buf[0..4] != b"OggS" {
61            return Err(Error::BadCapturePattern);
62        }
63        let version = self.buf[4];
64        if version != 0 {
65            return Err(Error::UnsupportedVersion(version));
66        }
67        let flags = self.buf[5];
68        let continued = flags & 0x01 != 0;
69        let bos = flags & 0x02 != 0;
70        let eos = flags & 0x04 != 0;
71        let granule_position = i64::from_le_bytes(self.buf[6..14].try_into().unwrap_or_default());
72        let serial = u32::from_le_bytes(self.buf[14..18].try_into().unwrap_or_default());
73        let crc_declared = u32::from_le_bytes(self.buf[22..26].try_into().unwrap_or_default());
74        let page_segments = usize::from(self.buf[26]);
75        let with_seg_table_len = HEADER_LEN + page_segments;
76        if self.buf.len() < with_seg_table_len {
77            return Ok(false);
78        }
79        let segment_table = self.buf[HEADER_LEN..with_seg_table_len].to_vec();
80        let payload_len: usize = segment_table.iter().map(|&s| usize::from(s)).sum();
81        let total_page_len = with_seg_table_len + payload_len;
82        if self.buf.len() < total_page_len {
83            return Ok(false);
84        }
85
86        if continued != self.has_partial {
87            return Err(Error::ContinuationFlagMismatch { flag: continued });
88        }
89
90        let mut page_for_crc = self.buf[0..total_page_len].to_vec();
91        page_for_crc[CRC_FIELD_OFFSET..CRC_FIELD_OFFSET + 4].fill(0);
92        let computed = crc32_ogg(&page_for_crc);
93        if computed != crc_declared {
94            return Err(Error::CrcMismatch {
95                expected: crc_declared,
96                computed,
97            });
98        }
99
100        let payload_start = with_seg_table_len;
101        let mut seg_start = payload_start;
102        let mut offset = payload_start;
103        // Packets completed on this page (each non-255 segment terminator
104        // completes exactly one packet). `page_count` is known upfront so the
105        // demuxer can hand each packet its position among the page's completed
106        // packets — codec-aware consumers back-compute per-packet positions
107        // from the page granule using these.
108        let page_count =
109            u32::try_from(segment_table.iter().filter(|&&s| s < 255).count()).unwrap_or(u32::MAX);
110        let mut page_index = 0u32;
111        for &seg in &segment_table {
112            offset += usize::from(seg);
113            if seg < 255 {
114                let chunk = &self.buf[seg_start..offset];
115                seg_start = offset;
116                if self.has_partial {
117                    self.partial.extend_from_slice(chunk);
118                    self.ready.push_back(Packet {
119                        data: Bytes::copy_from_slice(&self.partial),
120                        granule_position,
121                        serial,
122                        bos,
123                        eos,
124                        page_index,
125                        page_count,
126                    });
127                    self.partial.clear();
128                    self.has_partial = false;
129                } else {
130                    self.ready.push_back(Packet {
131                        data: Bytes::copy_from_slice(chunk),
132                        granule_position,
133                        serial,
134                        bos,
135                        eos,
136                        page_index,
137                        page_count,
138                    });
139                }
140                page_index += 1;
141            }
142        }
143        if seg_start < total_page_len {
144            self.partial
145                .extend_from_slice(&self.buf[seg_start..total_page_len]);
146            self.has_partial = true;
147        }
148
149        self.buf.drain(0..total_page_len);
150        Ok(true)
151    }
152}
153
154#[cfg(test)]
155#[path = "demux_tests.rs"]
156mod tests;