Skip to main content

bg_movie_writer/
lib.rs

1// Copyright (C) The Strand-Braid Authors
2// SPDX-License-Identifier: MIT OR Apache-2.0
3
4use std::{
5    path::PathBuf,
6    sync::{Arc, Mutex},
7};
8
9use strand_dynamic_frame::DynamicFrameOwned;
10
11mod movie_writer_thread;
12
13/// Possible errors
14#[derive(Debug, thiserror::Error)]
15pub enum Error {
16    #[error("io error: {source}")]
17    IoError {
18        #[from]
19        source: std::io::Error,
20    },
21    #[error("mp4 writer error: {0}")]
22    Mp4WriterError(#[from] mp4_writer::Error),
23    #[error("WorkerDisconnected")]
24    WorkerDisconnected,
25    #[error("AlreadyClosed")]
26    AlreadyClosed,
27    #[error(transparent)]
28    RecvError(#[from] std::sync::mpsc::RecvError),
29    #[error("already done")]
30    AlreadyDone,
31    #[error("disconnected")]
32    Disconnected,
33    #[error("filename does not end with '.mp4'")]
34    FilenameDoesNotEndWithMp4,
35    #[error("ffmpeg rewriter error {0}")]
36    FfmpegReWriterError(#[from] ffmpeg_rewriter::Error),
37    #[error("error loading CUDA or nvidia-encode: {0}")]
38    NvEncLoad(String),
39    #[error("error starting nvidia-encode: {0}")]
40    NvEncStart(String),
41}
42
43type Result<T> = std::result::Result<T, Error>;
44
45/// From outside the worker thread, check if we received an error from the
46/// thread.
47macro_rules! poll_err {
48    ($err_rx: expr_2021) => {{
49        // Recover from a poisoned lock rather than panicking: a panic here
50        // would propagate out of `write`/`finish` and take down the process.
51        if let Some(e) = $err_rx.lock().unwrap_or_else(|e| e.into_inner()).take() {
52            return Err(e);
53        }
54    }};
55}
56
57/// A writer which will save a movie in a background thread.
58///
59/// [Self::new] will spawn the thread and the methods [Self::write] and
60/// [Self::finish] return immediately, even though their work is not done.
61pub struct BgMovieWriter {
62    tx: std::sync::mpsc::SyncSender<Msg>,
63    is_done: bool,
64    err_from_worker: Arc<Mutex<Option<Error>>>,
65}
66
67impl BgMovieWriter {
68    /// This spawns the writer thread.
69    ///
70    /// - `format_str_mp4` determines the filename used after formatting with
71    ///   [chrono::DateTime::format].
72    /// - `recording_config` specifies the recording method and configuration
73    /// - `queue_size` is the number of frames that can be buffered before
74    ///   frames will be dropped.
75    /// - `data_dir`, if specified, will be the directory location of the saved
76    ///   file.
77    pub fn new(
78        recording_config: strand_cam_remote_control::RecordingConfig,
79        queue_size: usize,
80        mp4_path: PathBuf,
81    ) -> Self {
82        // Create an Arc<Mutex<Option<Error>>> to hold a potential error from
83        // the to-be-spawned writer thread.
84        let err_to_launcher = Arc::new(Mutex::new(None));
85        let err_from_worker = err_to_launcher.clone();
86        // Create a channel to send data into the writer thread.
87        let (tx, rx) = std::sync::mpsc::sync_channel::<Msg>(queue_size);
88        // Spawn the writer thread
89        std::thread::spawn(move || {
90            // Runs until the movie is done.
91            movie_writer_thread::writer_thread_loop(recording_config, err_to_launcher, rx, mp4_path)
92        });
93        Self {
94            tx,
95            is_done: false,
96            err_from_worker,
97        }
98    }
99
100    /// Enqueue the frame and timestamp for writing to the background thread.
101    ///
102    /// If the background writer thread has previously encountered an error,
103    /// this will return that previously-encountered error.
104    pub fn write<TS>(&mut self, frame: Arc<DynamicFrameOwned>, timestamp: TS) -> Result<()>
105    where
106        TS: Into<chrono::DateTime<chrono::Local>>,
107    {
108        let timestamp = timestamp.into();
109        poll_err!(self.err_from_worker);
110        if self.is_done {
111            return Err(Error::AlreadyDone);
112        }
113        let msg = Msg::Write((frame, timestamp));
114        // This will only succeed if the channel is not full. It will not block.
115        match self.tx.try_send(msg) {
116            Ok(()) => {}
117            Err(std::sync::mpsc::TrySendError::Full(_msg)) => {
118                tracing::warn!("Dropping frame to save: channel full");
119            }
120            Err(std::sync::mpsc::TrySendError::Disconnected(_msg)) => {
121                return Err(Error::WorkerDisconnected);
122            }
123        }
124        Ok(())
125    }
126
127    /// Enqueue a message telling the background thread to finish writing.
128    ///
129    /// If the background writer thread has previously encountered an error,
130    /// this will return that previously-encountered error.
131    pub fn finish(&mut self) -> Result<()> {
132        poll_err!(self.err_from_worker);
133        self.is_done = true;
134        let tx = self.tx.clone();
135        // We want to send the finish message without fail, so spawn a new
136        // thread which blocks until the message can be sent. If we don't spawn
137        // a new thread, the writer thread could be busy and block. If we don't
138        // block on sending, a full channel could cause the finish message to be
139        // dropped.
140        std::thread::spawn(move || {
141            // If the receiver has disconnected (e.g. the writer thread already
142            // exited after an error), there is nothing to finish. Do not panic.
143            if tx.send(Msg::Finish).is_err() {
144                tracing::debug!("writer thread already gone; nothing to finish");
145            }
146        });
147        Ok(())
148    }
149}
150
151pub(crate) enum Msg {
152    Write((Arc<DynamicFrameOwned>, chrono::DateTime<chrono::Local>)),
153    Finish,
154}
155
156#[cfg(test)]
157mod tests {
158    use super::*;
159    use machine_vision_formats::PixFmt;
160    use strand_dynamic_frame::DynamicFrameOwned;
161
162    /// Regression test: a writer error in the background thread must be
163    /// reported as an `Err` from the launcher-side methods, never as a panic.
164    ///
165    /// Previously the error-reporting path panicked on the first error (and
166    /// poisoned the shared error mutex), which propagated out and took down the
167    /// whole process. Here we force `create_writer` to fail inside the worker
168    /// thread by giving the ffmpeg writer a filename that does not end in
169    /// `.mp4` (this fails before any external `ffmpeg` process is spawned).
170    #[test]
171    fn writer_error_is_reported_without_panic() {
172        let cfg = strand_cam_remote_control::RecordingConfig::default();
173        let bad_path = std::env::temp_dir().join("bg_movie_writer_test.not_mp4");
174        let mut wtr = BgMovieWriter::new(cfg, 10, bad_path);
175
176        let frame =
177            Arc::new(DynamicFrameOwned::from_buf(4, 4, 4, vec![0u8; 16], PixFmt::Mono8).unwrap());
178        let ts = chrono::Local::now();
179
180        // Enqueue frames until the worker's error surfaces on the launcher
181        // side. This must arrive as an `Err`, never as a panic.
182        let mut got_err = false;
183        for _ in 0..200 {
184            if wtr.write(frame.clone(), ts).is_err() {
185                got_err = true;
186                break;
187            }
188            std::thread::sleep(std::time::Duration::from_millis(5));
189        }
190        assert!(got_err, "expected the writer error to be reported");
191
192        // Further calls must still behave gracefully (the mutex must not be
193        // poisoned) and must not panic.
194        let _ = wtr.write(frame, ts);
195        wtr.finish().unwrap();
196    }
197}