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}