1use alloc::vec::Vec;
27use dvb_bbframe::header::{BBHEADER_LEN, Bbheader, Mode};
28use dvb_bbframe::packet::{CarryOverExtractor, NM_UP_SIZE};
29
30use crate::payload::AnyPayload;
31use crate::pump::{Stats, T2miPump};
32
33pub struct InnerTsRecovery {
44 pump: T2miPump,
45 extractor: CarryOverExtractor,
46 out: Vec<[u8; NM_UP_SIZE]>,
47 up_buf: Vec<[u8; NM_UP_SIZE]>,
48 target_plp: Option<u8>,
49 filtered_out: u64,
50}
51
52impl InnerTsRecovery {
53 #[must_use]
56 pub fn new(t2mi_pid: u16) -> Self {
57 Self::build(t2mi_pid, None)
58 }
59
60 #[must_use]
64 pub fn new_for_plp(t2mi_pid: u16, plp_id: u8) -> Self {
65 Self::build(t2mi_pid, Some(plp_id))
66 }
67
68 fn build(t2mi_pid: u16, target_plp: Option<u8>) -> Self {
69 Self {
70 pump: T2miPump::new(t2mi_pid),
71 extractor: CarryOverExtractor::new(),
72 out: Vec::new(),
73 up_buf: Vec::new(),
74 target_plp,
75 filtered_out: 0,
76 }
77 }
78
79 pub fn feed(&mut self, ts_packet: &[u8]) -> &[[u8; NM_UP_SIZE]] {
83 self.out.clear();
84 let events: Vec<_> = self.pump.feed_ts(ts_packet).collect();
88 for event in events {
89 let Ok(AnyPayload::Bbframe(bb)) = event.payload() else {
90 continue;
91 };
92 if self.target_plp.is_some_and(|t| bb.plp_id != t) {
93 self.filtered_out += 1;
94 continue;
95 }
96 if bb.bbframe.len() < BBHEADER_LEN {
97 continue;
98 }
99 let Ok(hdr) = Bbheader::parse(bb.bbframe) else {
100 continue;
101 };
102 let header_bytes: [u8; BBHEADER_LEN] = match bb.bbframe[..BBHEADER_LEN].try_into() {
103 Ok(b) => b,
104 Err(_) => continue,
105 };
106 let data_field = &bb.bbframe[BBHEADER_LEN..];
107 match hdr.mode {
108 Mode::Normal => {
109 self.extractor
110 .feed_nm_into(&header_bytes, data_field, &mut self.up_buf);
111 }
112 Mode::HighEfficiency if !hdr.matype.npd => {
113 self.extractor.feed_hem_into(
114 &header_bytes,
115 data_field,
116 false,
117 &mut self.up_buf,
118 );
119 }
120 _ => continue,
122 }
123 self.out.append(&mut self.up_buf);
124 }
125 &self.out
126 }
127
128 #[must_use]
131 pub fn stats(&self) -> Stats {
132 self.pump.stats()
133 }
134
135 #[must_use]
139 pub fn filtered_bbframes(&self) -> u64 {
140 self.filtered_out
141 }
142}
143
144#[cfg(test)]
145mod tests {
146 use super::*;
147 use broadcast_common::crc32_mpeg2;
148 use dvb_bbframe::crc::crc8;
149 use dvb_bbframe::header::{Matype, TsGs};
150
151 const TS_SYNC: u8 = 0x47;
152 const TS_LEN: usize = 188;
153
154 fn inner_packet() -> [u8; TS_LEN] {
156 let mut p = [0xAAu8; TS_LEN];
157 p[0] = TS_SYNC;
158 p[1] = 0x41; p[2] = 0x00;
160 p[3] = 0x10; p
162 }
163
164 fn nm_bbframe(inner: &[u8; TS_LEN]) -> Vec<u8> {
166 let hdr = Bbheader {
167 matype: Matype {
168 ts_gs: TsGs::Ts,
169 sis: true,
170 ccm: true,
171 issyi: false,
172 npd: false,
173 ext: 0,
174 isi: 0,
175 },
176 upl: 1504,
177 sync: TS_SYNC,
178 dfl: 1504,
179 syncd: 0,
180 mode: Mode::Normal,
181 issy_in_header: None,
182 };
183 let mut frame = hdr.serialize().to_vec();
184 let mut data = [0u8; TS_LEN];
185 data[0] = crc8(&[0u8; TS_LEN]); data[1..].copy_from_slice(&inner[1..]);
187 frame.extend_from_slice(&data);
188 frame
189 }
190
191 fn t2mi_packet(bbframe: &[u8]) -> Vec<u8> {
193 let mut payload = vec![0x00, 0x05, 0x80]; payload.extend_from_slice(bbframe);
195 let mut pkt = vec![0x00u8, 0x01, 0x00, 0x00];
196 pkt.extend_from_slice(&((payload.len() * 8) as u16).to_be_bytes());
197 pkt.extend_from_slice(&payload);
198 let crc = crc32_mpeg2::compute(&pkt);
199 pkt.extend_from_slice(&crc.to_be_bytes());
200 pkt
201 }
202
203 fn outer_ts(pid: u16, data: &[u8]) -> Vec<[u8; TS_LEN]> {
205 let mut out = Vec::new();
206 let first_cap = TS_LEN - 5;
207 let cont_cap = TS_LEN - 4;
208 let mut off = 0;
209 let mut first = true;
210 while off < data.len() {
211 let mut pkt = [0xFFu8; TS_LEN];
212 pkt[0] = TS_SYNC;
213 let cap = if first { first_cap } else { cont_cap };
214 pkt[1] = (if first { 0x40 } else { 0x00 }) | (((pid >> 8) as u8) & 0x1F);
215 pkt[2] = (pid & 0xFF) as u8;
216 pkt[3] = 0x10;
217 let hdr_len = if first {
218 pkt[4] = 0x00; 5
220 } else {
221 4
222 };
223 let n = (data.len() - off).min(cap);
224 pkt[hdr_len..hdr_len + n].copy_from_slice(&data[off..off + n]);
225 out.push(pkt);
226 off += n;
227 first = false;
228 }
229 out
230 }
231
232 #[test]
233 fn recovers_inner_ts_from_nm_bbframe_chain() {
234 let pid = 0x1000;
235 let inner = inner_packet();
236 let outer = outer_ts(pid, &t2mi_packet(&nm_bbframe(&inner)));
237
238 let mut rec = InnerTsRecovery::new(pid);
239 let mut recovered: Vec<[u8; TS_LEN]> = Vec::new();
240 for pkt in &outer {
241 recovered.extend_from_slice(rec.feed(pkt));
242 }
243
244 assert_eq!(recovered.len(), 1, "exactly one inner TS packet expected");
245 assert_eq!(recovered[0][0], TS_SYNC, "sync byte restored");
246 assert_eq!(&recovered[0][1..], &inner[1..]);
248 }
249
250 #[test]
251 fn wrong_pid_yields_nothing() {
252 let inner = inner_packet();
253 let outer = outer_ts(0x1000, &t2mi_packet(&nm_bbframe(&inner)));
254 let mut rec = InnerTsRecovery::new(0x0064); let mut n = 0;
256 for pkt in &outer {
257 n += rec.feed(pkt).len();
258 }
259 assert_eq!(n, 0);
260 }
261
262 #[test]
263 fn garbage_packet_no_panic_no_output() {
264 let mut rec = InnerTsRecovery::new(0x1000);
265 let junk = [0u8; TS_LEN];
266 assert!(rec.feed(&junk).is_empty());
267 }
268
269 fn t2mi_packet_for_plp(bbframe: &[u8], plp_id: u8) -> Vec<u8> {
271 let mut payload = vec![0x00, plp_id, 0x80]; payload.extend_from_slice(bbframe);
273 let mut pkt = vec![0x00u8, 0x01, 0x00, 0x00];
274 pkt.extend_from_slice(&((payload.len() * 8) as u16).to_be_bytes());
275 pkt.extend_from_slice(&payload);
276 let crc = crc32_mpeg2::compute(&pkt);
277 pkt.extend_from_slice(&crc.to_be_bytes());
278 pkt
279 }
280
281 fn tagged_inner_packet(marker: u8) -> [u8; TS_LEN] {
283 let mut p = [0xAAu8; TS_LEN];
284 p[0] = TS_SYNC;
285 p[1] = 0x41;
286 p[2] = 0x00;
287 p[3] = 0x10;
288 p[4] = marker;
289 p
290 }
291
292 #[test]
293 fn plp_filter_keeps_only_target_plp() {
294 let pid = 0x1000;
295 let inner_plp0 = tagged_inner_packet(0xA0);
296 let inner_plp1 = tagged_inner_packet(0xB0);
297
298 let bb_plp0 = nm_bbframe(&inner_plp0);
300 let bb_plp1 = nm_bbframe(&inner_plp1);
301
302 let t2mi_plp0 = t2mi_packet_for_plp(&bb_plp0, 0);
304 let t2mi_plp1 = t2mi_packet_for_plp(&bb_plp1, 1);
305
306 let mut combined = t2mi_plp0;
308 combined.extend_from_slice(&t2mi_plp1);
309 let outer = outer_ts(pid, &combined);
310
311 let mut rec0 = InnerTsRecovery::new_for_plp(pid, 0);
313 let mut recovered_0: Vec<[u8; TS_LEN]> = Vec::new();
314 for pkt in &outer {
315 recovered_0.extend_from_slice(rec0.feed(pkt));
316 }
317 assert_eq!(
318 recovered_0.len(),
319 1,
320 "plp 0 filter should recover exactly one inner packet"
321 );
322 assert_eq!(recovered_0[0][4], 0xA0, "should be the plp 0 packet");
323 assert_eq!(
324 rec0.filtered_bbframes(),
325 1,
326 "one BBFRAME (plp 1) filtered out"
327 );
328
329 let mut rec1 = InnerTsRecovery::new_for_plp(pid, 1);
331 let mut recovered_1: Vec<[u8; TS_LEN]> = Vec::new();
332 for pkt in &outer {
333 recovered_1.extend_from_slice(rec1.feed(pkt));
334 }
335 assert_eq!(
336 recovered_1.len(),
337 1,
338 "plp 1 filter should recover exactly one inner packet"
339 );
340 assert_eq!(recovered_1[0][4], 0xB0, "should be the plp 1 packet");
341 assert_eq!(
342 rec1.filtered_bbframes(),
343 1,
344 "one BBFRAME (plp 0) filtered out"
345 );
346
347 let mut all = InnerTsRecovery::new(pid);
349 let mut recovered_all: Vec<[u8; TS_LEN]> = Vec::new();
350 for pkt in &outer {
351 recovered_all.extend_from_slice(all.feed(pkt));
352 }
353 assert_eq!(
354 recovered_all.len(),
355 2,
356 "unfiltered should recover both inner packets"
357 );
358 assert_eq!(
359 all.filtered_bbframes(),
360 0,
361 "no filtering when target is None"
362 );
363 }
364}