66#include <fcntl.h>
77#include <inttypes.h>
88#include <pthread.h>
9+ #include <signal.h>
910#include <stdatomic.h>
1011#include <stdio.h>
1112#include <string.h>
1213#include <sys/stat.h>
1314#include <sys/wait.h>
15+ #include <time.h>
1416#include <unistd.h>
1517
1618#include <libavformat/avformat.h>
2830 * PATH already covers both. */
2931#define FFMPEG_BINARY "ffmpeg"
3032
31- /* Playback handlers run concurrently in the HTTP worker pool. Keep only one
32- * recording transcode active across all recordings: even different clips can
33- * exhaust the memory shared with recording and go2rtc on a small NVR. */
33+ /* Keep only one recording transcode active, including synchronous callers:
34+ * different clips can exhaust memory shared with recording on a small NVR. */
3435static pthread_mutex_t s_transcode_mutex = PTHREAD_MUTEX_INITIALIZER ;
36+ static atomic_bool s_stopping = false;
37+
38+ /* This mutex only protects job bookkeeping, never an encode or a slot wait.
39+ * A separate thread keeps both FFmpeg and s_transcode_mutex waits out of
40+ * libuv's shared pool (which also serves APIs and reads video file chunks). */
41+ static pthread_mutex_t s_job_mutex = PTHREAD_MUTEX_INITIALIZER ;
42+ static struct {
43+ pthread_t thread ;
44+ bool started ;
45+ bool active ;
46+ char original_path [512 ];
47+ char cache_path [512 ];
48+ int result ;
49+ struct timespec finished ;
50+ } s_job ;
3551
3652static bool file_exists_nonempty (const char * path ) {
3753 struct stat st ;
@@ -91,11 +107,20 @@ static int run_and_wait(char *const argv[]) {
91107 }
92108
93109 int status = 0 ;
94- while (waitpid (pid , & status , 0 ) < 0 ) {
95- if (errno != EINTR ) {
110+ for (;;) {
111+ pid_t result = waitpid (pid , & status , WNOHANG );
112+ if (result == pid ) break ;
113+ if (result < 0 && errno != EINTR ) {
96114 log_error ("recording_transcode: waitpid failed: %s" , strerror (errno ));
97115 return -1 ;
98116 }
117+ if (atomic_load (& s_stopping )) {
118+ kill (pid , SIGKILL );
119+ while (waitpid (pid , & status , 0 ) < 0 && errno == EINTR ) {}
120+ return -1 ;
121+ }
122+ const struct timespec delay = { .tv_nsec = 100000000 };
123+ nanosleep (& delay , NULL );
99124 }
100125
101126 return (WIFEXITED (status ) && WEXITSTATUS (status ) == 0 ) ? 0 : -1 ;
@@ -105,6 +130,7 @@ static int run_and_wait(char *const argv[]) {
105130 * completed cache file. Recheck the cache because another request may have
106131 * produced it while this caller was waiting. */
107132static int transcode_cache_locked (const char * original_path , const char * cache_path ) {
133+ if (atomic_load (& s_stopping )) return -1 ;
108134 if (file_exists_nonempty (cache_path )) {
109135 return 0 ;
110136 }
@@ -193,7 +219,7 @@ static int transcode_cache_locked(const char *original_path, const char *cache_p
193219 unlink (tmp_path );
194220 }
195221 }
196- if (rc != 0 ) {
222+ if (rc != 0 && ! atomic_load ( & s_stopping ) ) {
197223 log_info ("recording_transcode: transcoding %s -> %s via software libx264" , original_path , cache_path );
198224 rc = run_and_wait (argv_software );
199225 }
@@ -234,3 +260,71 @@ int ensure_recording_transcode_cache(const char *original_path, const char *cach
234260 pthread_mutex_unlock (& s_transcode_mutex );
235261 return rc ;
236262}
263+
264+ static void * transcode_worker (void * unused ) {
265+ (void )unused ;
266+ int result = ensure_recording_transcode_cache (s_job .original_path , s_job .cache_path );
267+ pthread_mutex_lock (& s_job_mutex );
268+ s_job .result = result ;
269+ clock_gettime (CLOCK_MONOTONIC , & s_job .finished );
270+ s_job .active = false;
271+ pthread_mutex_unlock (& s_job_mutex );
272+ return NULL ;
273+ }
274+
275+ recording_transcode_status_t request_recording_transcode_cache (
276+ const char * original_path , const char * cache_path ) {
277+ if (!original_path || !original_path [0 ] || !cache_path || !cache_path [0 ] ||
278+ strlen (original_path ) >= sizeof (s_job .original_path ) ||
279+ strlen (cache_path ) >= sizeof (s_job .cache_path )) {
280+ return RECORDING_TRANSCODE_FAILED ;
281+ }
282+ if (file_exists_nonempty (cache_path )) return RECORDING_TRANSCODE_READY ;
283+
284+ pthread_mutex_lock (& s_job_mutex );
285+ recording_transcode_status_t status = RECORDING_TRANSCODE_PENDING ;
286+ if (atomic_load (& s_stopping )) {
287+ status = RECORDING_TRANSCODE_FAILED ;
288+ } else if (file_exists_nonempty (cache_path )) {
289+ status = RECORDING_TRANSCODE_READY ;
290+ } else if (!s_job .active ) {
291+ if (s_job .started ) {
292+ pthread_join (s_job .thread , NULL );
293+ s_job .started = false;
294+ }
295+ struct timespec now ;
296+ clock_gettime (CLOCK_MONOTONIC , & now );
297+ if (s_job .result != 0 &&
298+ strcmp (original_path , s_job .original_path ) == 0 &&
299+ strcmp (cache_path , s_job .cache_path ) == 0 &&
300+ now .tv_sec - s_job .finished .tv_sec < 30 ) {
301+ status = RECORDING_TRANSCODE_FAILED ;
302+ } else {
303+ strcpy (s_job .original_path , original_path );
304+ strcpy (s_job .cache_path , cache_path );
305+ s_job .active = true;
306+ int err = pthread_create (& s_job .thread , NULL , transcode_worker , NULL );
307+ if (err != 0 ) {
308+ log_error ("recording_transcode: failed to start worker: %s" , strerror (err ));
309+ s_job .active = false;
310+ s_job .result = -1 ;
311+ s_job .finished = now ;
312+ status = RECORDING_TRANSCODE_FAILED ;
313+ } else {
314+ s_job .started = true;
315+ }
316+ }
317+ }
318+ pthread_mutex_unlock (& s_job_mutex );
319+ return status ;
320+ }
321+
322+ void shutdown_recording_transcode (void ) {
323+ pthread_mutex_lock (& s_job_mutex );
324+ atomic_store (& s_stopping , true);
325+ bool join = s_job .started ;
326+ pthread_t thread = s_job .thread ;
327+ s_job .started = false;
328+ pthread_mutex_unlock (& s_job_mutex );
329+ if (join ) pthread_join (thread , NULL );
330+ }
0 commit comments