1#![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; const CRC_FIELD_OFFSET: usize = 4 + 1 + 1 + 8 + 4 + 4;
18
19#[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 #[must_use]
31 pub fn new() -> Self {
32 Self::default()
33 }
34
35 pub fn push_bytes(&mut self, data: &[u8]) {
37 self.buf.extend_from_slice(data);
38 }
39
40 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 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 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;