Skip to main content

ffmpeg_pipeline/
decode.rs

1use super::*;
2use ffmpeg_next::{codec::context::Context, error::EAGAIN, Error as FFmpegOrigError, Packet};
3use std::collections::VecDeque;
4
5#[derive(Clone, Copy, Debug, Eq, PartialEq)]
6pub enum FrameProcess {
7    Passthrough,
8    Decode,
9}
10
11pub enum Frame {
12    Packet(Packet),
13    Frame(StreamFrame),
14}
15
16#[derive(Clone, Copy, Debug, Eq, PartialEq)]
17enum DecoderState {
18    Reading,
19    SendingEof,
20    Draining,
21    Done,
22}
23
24enum ReceiveStatus<T> {
25    Frame(T),
26    Again,
27    Eof,
28}
29
30#[cfg(test)]
31#[derive(Default)]
32struct DecoderStats {
33    receive_again: usize,
34    eof_sends: usize,
35}
36
37#[derive(Clone, Copy, Debug, Eq, PartialEq)]
38enum ReceiveEnd {
39    Again,
40    Eof,
41}
42
43fn drain_available<T, E>(
44    mut receive: impl FnMut() -> Result<ReceiveStatus<T>, E>,
45) -> Result<(VecDeque<T>, ReceiveEnd), E> {
46    let mut frames = VecDeque::new();
47    loop {
48        match receive()? {
49            ReceiveStatus::Frame(frame) => frames.push_back(frame),
50            ReceiveStatus::Again => return Ok((frames, ReceiveEnd::Again)),
51            ReceiveStatus::Eof => return Ok((frames, ReceiveEnd::Eof)),
52        }
53    }
54}
55
56/// Iterates packets or decoded frames for one stream.
57///
58/// Passthrough yields every packet from the selected stream and then `None` at
59/// demux EOF. Decode mode drains all frames made available by a packet before
60/// reading another packet, sends EOF once, drains delayed frames, and returns
61/// `None` only after decoder EOF. Every other demux or codec failure is yielded.
62pub struct Decoder<'i> {
63    index: usize,
64    decoder: StreamDecoder,
65    input: &'i mut Input,
66    process: FrameProcess,
67    time_base: Rational,
68    state: DecoderState,
69    pending_packet: Option<Packet>,
70    drain_again: bool,
71    ready_frames: VecDeque<StreamFrame>,
72    receive_end: Option<ReceiveEnd>,
73    #[cfg(test)]
74    stats: DecoderStats,
75}
76
77impl<'i> Decoder<'i> {
78    pub fn new_with_video(
79        input: &'i mut Input,
80        index: usize,
81        process: FrameProcess,
82    ) -> FFmpegResult<Self> {
83        let (decoder, time_base) = {
84            let stream = input
85                .stream(index)
86                .ok_or(FFmpegError::StreamNotFound(index))?;
87            let codec = Context::from_parameters(stream.parameters())?;
88            (
89                StreamDecoder::Video(codec.decoder().video()?),
90                stream.time_base(),
91            )
92        };
93        Ok(Self::new(input, index, process, decoder, time_base))
94    }
95
96    pub fn new_with_audio(
97        input: &'i mut Input,
98        index: usize,
99        process: FrameProcess,
100    ) -> FFmpegResult<Self> {
101        let (decoder, time_base) = {
102            let stream = input
103                .stream(index)
104                .ok_or(FFmpegError::StreamNotFound(index))?;
105            let codec = Context::from_parameters(stream.parameters())?;
106            (
107                StreamDecoder::Audio(codec.decoder().audio()?),
108                stream.time_base(),
109            )
110        };
111        Ok(Self::new(input, index, process, decoder, time_base))
112    }
113
114    fn new(
115        input: &'i mut Input,
116        index: usize,
117        process: FrameProcess,
118        decoder: StreamDecoder,
119        time_base: Rational,
120    ) -> Self {
121        Self {
122            index,
123            decoder,
124            input,
125            process,
126            time_base,
127            state: DecoderState::Reading,
128            pending_packet: None,
129            drain_again: false,
130            ready_frames: VecDeque::new(),
131            receive_end: None,
132            #[cfg(test)]
133            stats: DecoderStats::default(),
134        }
135    }
136
137    pub fn get_decoder(&self) -> &StreamDecoder {
138        &self.decoder
139    }
140
141    pub fn time_base(&self) -> Rational {
142        self.time_base
143    }
144
145    fn receive_frame(&mut self) -> FFmpegResult<ReceiveStatus<StreamFrame>> {
146        let result = match self.decoder {
147            StreamDecoder::Audio(ref mut decoder) => {
148                let mut frame = AudioFrame::empty();
149                decoder
150                    .receive_frame(&mut frame)
151                    .map(|_| ReceiveStatus::Frame(StreamFrame::Audio(frame)))
152            }
153            StreamDecoder::Video(ref mut decoder) => {
154                let mut frame = VideoFrame::empty();
155                decoder
156                    .receive_frame(&mut frame)
157                    .map(|_| ReceiveStatus::Frame(StreamFrame::Video(frame)))
158            }
159        };
160        match result {
161            Ok(frame) => Ok(frame),
162            Err(error) if is_again(error) => Ok(ReceiveStatus::Again),
163            Err(FFmpegOrigError::Eof) => Ok(ReceiveStatus::Eof),
164            Err(source) => Err(FFmpegError::decoder(
165                "receive frame",
166                self.index,
167                None,
168                source,
169            )),
170        }
171    }
172
173    fn send_packet(&mut self, packet: &Packet) -> Result<(), FFmpegOrigError> {
174        match self.decoder {
175            StreamDecoder::Audio(ref mut decoder) => decoder.send_packet(packet),
176            StreamDecoder::Video(ref mut decoder) => decoder.send_packet(packet),
177        }
178    }
179
180    fn send_eof(&mut self) -> Result<(), FFmpegOrigError> {
181        match self.decoder {
182            StreamDecoder::Audio(ref mut decoder) => decoder.send_eof(),
183            StreamDecoder::Video(ref mut decoder) => decoder.send_eof(),
184        }
185    }
186
187    fn read_selected_packet(&mut self) -> FFmpegResult<Option<Packet>> {
188        loop {
189            let mut packet = Packet::empty();
190            match packet.read(self.input) {
191                Ok(()) if packet.stream() == self.index => return Ok(Some(packet)),
192                Ok(()) => continue,
193                Err(FFmpegOrigError::Eof) => return Ok(None),
194                Err(source) => return Err(FFmpegError::decoder("demux", self.index, None, source)),
195            }
196        }
197    }
198
199    fn next_packet(&mut self) -> Option<FFmpegResult<Frame>> {
200        match self.read_selected_packet() {
201            Ok(Some(packet)) => Some(Ok(Frame::Packet(packet))),
202            Ok(None) => {
203                self.state = DecoderState::Done;
204                None
205            }
206            Err(error) => {
207                self.state = DecoderState::Done;
208                Some(Err(error))
209            }
210        }
211    }
212
213    fn fail(&mut self, error: FFmpegError) -> Option<FFmpegResult<Frame>> {
214        self.state = DecoderState::Done;
215        Some(Err(error))
216    }
217}
218
219impl Iterator for Decoder<'_> {
220    type Item = FFmpegResult<Frame>;
221
222    fn next(&mut self) -> Option<Self::Item> {
223        if self.state == DecoderState::Done {
224            return None;
225        }
226        if self.process == FrameProcess::Passthrough {
227            return self.next_packet();
228        }
229
230        loop {
231            if let Some(frame) = self.ready_frames.pop_front() {
232                return Some(Ok(Frame::Frame(frame)));
233            }
234            if let Some(end) = self.receive_end.take() {
235                if end == ReceiveEnd::Eof {
236                    self.state = DecoderState::Done;
237                    return None;
238                }
239            } else {
240                match drain_available(|| self.receive_frame()) {
241                    Ok((frames, end)) => {
242                        self.ready_frames = frames;
243                        self.receive_end = Some(end);
244                        self.drain_again = false;
245                        #[cfg(test)]
246                        if end == ReceiveEnd::Again {
247                            self.stats.receive_again += 1;
248                        }
249                        continue;
250                    }
251                    Err(error) => return self.fail(error),
252                }
253            }
254
255            match self.state {
256                DecoderState::Reading => {
257                    if self.pending_packet.is_none() {
258                        match self.read_selected_packet() {
259                            Ok(Some(packet)) => self.pending_packet = Some(packet),
260                            Ok(None) => {
261                                self.state = DecoderState::SendingEof;
262                                continue;
263                            }
264                            Err(error) => return self.fail(error),
265                        }
266                    }
267                    let Some(packet) = self.pending_packet.take() else {
268                        continue;
269                    };
270                    let position = packet.position();
271                    match self.send_packet(&packet) {
272                        Ok(()) => {}
273                        Err(error) if is_again(error) => self.pending_packet = Some(packet),
274                        Err(source) => {
275                            return self.fail(FFmpegError::decoder(
276                                "send packet",
277                                self.index,
278                                Some(position),
279                                source,
280                            ));
281                        }
282                    }
283                }
284                DecoderState::SendingEof => {
285                    #[cfg(test)]
286                    {
287                        self.stats.eof_sends += 1;
288                    }
289                    match self.send_eof() {
290                        Ok(()) => self.state = DecoderState::Draining,
291                        Err(error) if is_again(error) => continue,
292                        Err(FFmpegOrigError::Eof) => {
293                            self.state = DecoderState::Done;
294                            return None;
295                        }
296                        Err(source) => {
297                            return self
298                                .fail(FFmpegError::decoder("flush", self.index, None, source));
299                        }
300                    }
301                }
302                DecoderState::Draining => {
303                    if self.drain_again {
304                        return self.fail(FFmpegError::InvalidFormat(format!(
305                            "decoder stream {} returned EAGAIN while draining",
306                            self.index
307                        )));
308                    }
309                    self.drain_again = true;
310                }
311                DecoderState::Done => return None,
312            }
313        }
314    }
315}
316
317fn is_again(error: FFmpegOrigError) -> bool {
318    matches!(error, FFmpegOrigError::Other { errno } if errno == EAGAIN)
319}
320
321#[cfg(test)]
322mod tests {
323    use super::*;
324
325    #[test]
326    fn decoder_drains_to_stable_eof_and_preserves_time_base() {
327        initialize(log::Level::Error).unwrap();
328        let mut input = input_buffer_with_format(crate::tests::encoded_ivf(4), "ivf").unwrap();
329        let expected_time_base = input.as_ref().stream(0).unwrap().time_base();
330        let mut decoder = Decoder::new_with_video(input.as_mut(), 0, FrameProcess::Decode).unwrap();
331        assert_eq!(decoder.time_base(), expected_time_base);
332
333        let mut frames = 0;
334        let mut timestamps = Vec::new();
335        for frame in decoder.by_ref() {
336            let Frame::Frame(StreamFrame::Video(frame)) = frame.unwrap() else {
337                panic!("unexpected frame type");
338            };
339            frames += 1;
340            timestamps.push(frame.timestamp());
341        }
342        assert_eq!(frames, 4);
343        assert!(timestamps.iter().all(Option::is_some));
344        assert!(timestamps.windows(2).all(|pair| pair[0] < pair[1]));
345        assert!(decoder.stats.receive_again > 0);
346        assert_eq!(decoder.stats.eof_sends, 1);
347        assert!(decoder.next().is_none());
348    }
349
350    #[test]
351    fn receive_drain_handles_zero_multiple_and_delayed_frames() {
352        let mut zero = VecDeque::from([ReceiveStatus::<u8>::Again]);
353        let (frames, end) = drain_available(|| Ok::<_, ()>(zero.pop_front().unwrap())).unwrap();
354        assert!(frames.is_empty());
355        assert_eq!(end, ReceiveEnd::Again);
356
357        let mut multiple = VecDeque::from([
358            ReceiveStatus::Frame(1),
359            ReceiveStatus::Frame(2),
360            ReceiveStatus::Again,
361        ]);
362        let (frames, end) = drain_available(|| Ok::<_, ()>(multiple.pop_front().unwrap())).unwrap();
363        assert_eq!(frames, VecDeque::from([1, 2]));
364        assert_eq!(end, ReceiveEnd::Again);
365
366        let mut delayed = VecDeque::from([ReceiveStatus::Frame(3), ReceiveStatus::Eof]);
367        let (frames, end) = drain_available(|| Ok::<_, ()>(delayed.pop_front().unwrap())).unwrap();
368        assert_eq!(frames, VecDeque::from([3]));
369        assert_eq!(end, ReceiveEnd::Eof);
370    }
371
372    #[test]
373    fn passthrough_yields_packets_and_invalid_stream_is_an_error() {
374        initialize(log::Level::Error).unwrap();
375        let mut input = input_buffer_with_format(crate::tests::encoded_ivf(2), "ivf").unwrap();
376        assert!(matches!(
377            Decoder::new_with_video(input.as_mut(), 2, FrameProcess::Decode),
378            Err(FFmpegError::StreamNotFound(2))
379        ));
380
381        let mut input = input_buffer_with_format(crate::tests::encoded_ivf(2), "ivf").unwrap();
382        let packets = Decoder::new_with_video(input.as_mut(), 0, FrameProcess::Passthrough)
383            .unwrap()
384            .map(|packet| match packet.unwrap() {
385                Frame::Packet(packet) => packet.size(),
386                Frame::Frame(_) => panic!("unexpected decoded frame"),
387            })
388            .collect::<Vec<_>>();
389        assert!(!packets.is_empty());
390        assert!(packets.iter().all(|size| *size > 0));
391    }
392
393    #[test]
394    fn decoded_planes_expose_stride_and_truncated_input_is_an_error() {
395        initialize(log::Level::Error).unwrap();
396        let mut input = input_buffer_with_format(crate::tests::encoded_ivf(1), "ivf").unwrap();
397        let frame = Decoder::new_with_video(input.as_mut(), 0, FrameProcess::Decode)
398            .unwrap()
399            .next()
400            .unwrap()
401            .unwrap();
402        let Frame::Frame(StreamFrame::Video(frame)) = frame else {
403            panic!("unexpected frame type");
404        };
405        assert!(frame.stride(0) >= frame.width() as usize);
406        assert!(frame.data(0).len() >= frame.stride(0) * frame.height() as usize);
407
408        let mut bytes = crate::tests::encoded_ivf(4);
409        bytes.truncate(bytes.len() - 8);
410        let mut input = input_buffer_with_format(bytes, "ivf").unwrap();
411        let result = Decoder::new_with_video(input.as_mut(), 0, FrameProcess::Decode)
412            .unwrap()
413            .collect::<FFmpegResult<Vec<_>>>();
414        assert!(matches!(result, Err(FFmpegError::Decoder { .. })));
415    }
416}