diff --git a/src/Cli/Commands/MonitorCommand.php b/src/Cli/Commands/MonitorCommand.php index 2baa6b8e..4d512a8b 100644 --- a/src/Cli/Commands/MonitorCommand.php +++ b/src/Cli/Commands/MonitorCommand.php @@ -12,6 +12,7 @@ use XcVm\Domain\Server\ServerRepository; use XcVm\Domain\Stream\StreamProcess; use XcVm\Domain\Stream\StreamSorter; use XcVm\Streaming\Codec\FFprobeRunner; +use XcVm\Streaming\Fanout\FanoutClient; /** * `monitor [restart]` — the per-stream watchdog. @@ -81,6 +82,21 @@ class MonitorCommand implements CommandInterface { } $rStreamInfo = $db->get_row(); + + // Stand down if the fanout daemon is already supervising this stream's + // encoder. Two watchdogs on one ffmpeg fight: each reads the other's kill + // as a stream failure and restarts, so the channel flaps indefinitely. + // + // The daemon is asked rather than assumed, and an unreachable daemon + // answers false — so this watchdog keeps doing its job whenever the + // daemon is not actually holding the stream. That is the rollback path, + // and it is why the check is here rather than at the call sites: however + // a monitor gets started, it defers to whoever really owns the process. + if (FanoutClient::isSupervised($rStreamID)) { + echo "Stream is supervised by the fanout daemon; monitor standing down.\n"; + return 0; + } + $db->query('UPDATE `streams_servers` SET `monitor_pid` = ? WHERE `server_stream_id` = ?', getmypid(), $rStreamInfo['server_stream_id']); if (SettingsManager::get('enable_cache')) { diff --git a/src/Domain/Stream/StreamProcess.php b/src/Domain/Stream/StreamProcess.php index 8a98b5b7..5eed6e30 100644 --- a/src/Domain/Stream/StreamProcess.php +++ b/src/Domain/Stream/StreamProcess.php @@ -894,7 +894,11 @@ class StreamProcess { $rFFMPEG .= '{MAP} -individual_header_trailer 0 -f hls -hls_time ' . intval($rSegmentSettings['seg_time']) . ' -hls_list_size ' . intval($rStream['stream_info']['delay_minutes']) * 6 . ' -hls_delete_threshold 4 -start_number ' . $rSegmentStart . ' -hls_flags delete_segments+discont_start+omit_endlist -hls_segment_type mpegts -hls_segment_filename "' . DELAY_PATH . intval($rStreamID) . '_%d.ts" "' . DELAY_PATH . intval($rStreamID) . '_.m3u8" '; } - $rFFMPEG .= ' >/dev/null 2>>' . STREAMS_PATH . intval($rStreamID) . '.errors & echo $! > ' . STREAMS_PATH . intval($rStreamID) . '_.pid'; + // NB: the redirect-and-background tail is NOT appended here any more. + // buildLive returns the BARE command so it can be handed to the fanout + // daemon to supervise (it has to be the process's parent, and a + // backgrounded command leaves it nothing to supervise). The legacy path + // appends liveRedirectTail() itself, so its behaviour is unchanged. $ffprobeContainer = (isset($rFFProbeOutput['container']) && is_string($rFFProbeOutput['container'])) ? $rFFProbeOutput['container'] : ''; @@ -921,6 +925,37 @@ class StreamProcess { return $rFFMPEG; } + /** + * The shell tail that runs a live command the legacy way: stderr appended to + * the stream's .errors file, the process backgrounded, and its pid written to + * `/_.pid`. + * + * Split out of buildLive so the same command can either be run here or handed + * to the fanout daemon, which supervises the process itself and therefore + * must NOT have it backgrounded — it does the redirection and the pid file on + * its own. Keeping the two in one place is what stops them drifting. + * + * @param int $rStreamID Stream id. + * @return string Shell fragment to append to a bare live command. + */ + public static function liveRedirectTail($rStreamID): string { + return ' >/dev/null 2>>' . STREAMS_PATH . intval($rStreamID) . '.errors & echo $! > ' . STREAMS_PATH . intval($rStreamID) . '_.pid'; + } + + /** + * Whether this node hands live encoders to the fanout daemon to supervise + * instead of running them itself under the PHP watchdog. + * + * Off unless `daemon_supervise` is set, so a panel that has never heard of + * the setting keeps the legacy behaviour exactly — no migration needed, and + * clearing the setting is the rollback. + * + * @return bool + */ + public static function daemonSupervises(): bool { + return (bool) SettingsManager::get('daemon_supervise'); + } + public static function createChannelItem($rStreamID, $rSource) { global $rSettings, $rServers, $rFFMPEG_CPU, $rFFMPEG_GPU; $db = self::db(); @@ -1526,7 +1561,27 @@ class StreamProcess { 'ingestSock' => $rIngestSock, )); - shell_exec($rFFMPEG); + // Hand the encoder to the daemon when this node supervises there. + // The daemon starts it, watches it and restarts it, writing the + // same pid file the rest of PHP still reads — so everything below + // (and stopStream, and isStreamRunning) keeps working unchanged. + // A daemon that is unreachable or declines falls through to the + // legacy shell_exec, which is the rollback path. + $rHandedOver = false; + if (self::daemonSupervises() && !$rDelayActive) { + $rHandedOver = FanoutClient::supervise($rStreamID, $rFFMPEG, (string) $rRealSource, array( + 'stop_failures' => intval($rSettings['stop_failures']), + 'stream_fail_sleep' => intval($rSettings['stream_fail_sleep']), + 'on_demand' => (bool) $rStream['server_info']['on_demand'], + 'on_demand_failure_exit' => (bool) $rSettings['on_demand_failure_exit'], + )); + } + if (!$rHandedOver) { + $rFFMPEG .= self::liveRedirectTail($rStreamID); + shell_exec($rFFMPEG); + } + // Record what actually ran: the bare command when the daemon owns + // the process, the backgrounded one when this node does. file_put_contents(STREAMS_PATH . $rStreamID . '_.ffmpeg', $rFFMPEG); // Wait briefly for PID file to be written, with retry diff --git a/src/Streaming/Fanout/FanoutClient.php b/src/Streaming/Fanout/FanoutClient.php index 83c4678b..04bd20cd 100644 --- a/src/Streaming/Fanout/FanoutClient.php +++ b/src/Streaming/Fanout/FanoutClient.php @@ -457,4 +457,126 @@ class FanoutClient { return $rCode >= 200 && $rCode < 300; } + + /** + * Hand a stream's encoder to the daemon to run and supervise (ADR 0002 in + * the daemon repo, M1). The daemon starts the command, watches it, restarts + * it on death and applies the panel's own give-up rules; PHP keeps building + * the command and stays the system of record. + * + * $rCmd must be buildLive()'s BARE command — no `>/dev/null 2>>… & echo $!` + * tail. Backgrounding it would leave the daemon nothing to supervise, and + * redirection plus the pid file are the daemon's job once it owns the + * process. StreamProcess::liveRedirectTail() supplies that tail for the + * legacy path, which still runs whenever this returns false. + * + * @param int $rStreamID Stream id. + * @param string $rCmd Bare ffmpeg command line. + * @param string $rLabel Source label to echo back in logs / current_source. + * @param array $rPolicy Restart policy (stop_failures, stream_fail_sleep, + * on_demand, on_demand_failure_exit, start_timeout_sec). + * @return bool True when the daemon accepted the stream. + */ + public static function supervise(int $rStreamID, string $rCmd, string $rLabel, array $rPolicy = []): bool { + if ($rCmd === '') { + return false; + } + $rBody = json_encode([ + 'sources' => [['label' => $rLabel, 'cmd' => $rCmd]], + 'policy' => $rPolicy, + 'pid_path' => STREAMS_PATH . $rStreamID . '_.pid', + 'errors_path' => STREAMS_PATH . $rStreamID . '.errors', + 'log_path' => LOGS_TMP_PATH . 'stream_log.log', + 'server_id' => SERVER_ID, + ]); + if ($rBody === false) { + return false; + } + return self::monitorCall('PUT', '/monitor/' . $rStreamID, $rBody) !== null; + } + + /** + * Stop supervising a stream and kill its encoder. Idempotent. + * + * @param int $rStreamID Stream id. + * @return bool True when the daemon acknowledged. + */ + public static function releaseSupervision(int $rStreamID): bool { + return self::monitorCall('DELETE', '/monitor/' . $rStreamID, null) !== null; + } + + /** + * Current supervision state for a stream, or null when the daemon is + * unreachable or is not supervising it. + * + * @param int $rStreamID Stream id. + * @return array|null Keys: running, pid, source, restarts, failures, + * uptime_ms, gave_up, last_error. + */ + public static function monitorState(int $rStreamID): ?array { + $rBody = self::monitorCall('GET', '/monitor/' . $rStreamID, null); + if ($rBody === null) { + return null; + } + $rJson = json_decode($rBody, true); + return is_array($rJson) ? $rJson : null; + } + + /** + * Whether the daemon is supervising this stream's encoder. The PHP watchdog + * uses this to stand down: two supervisors restarting the same ffmpeg would + * fight, each reading the other's kill as a stream failure. + * + * Fails CLOSED — an unreachable daemon reports false, so the PHP watchdog + * keeps doing its job. That is the rollback path. + * + * @param int $rStreamID Stream id. + * @return bool + */ + public static function isSupervised(int $rStreamID): bool { + $rState = self::monitorState($rStreamID); + return is_array($rState) && !empty($rState['supervised']); + } + + /** + * Issue a control request against an arbitrary daemon path and return the + * response body, or null when the daemon is unreachable or answers non-2xx. + * + * Distinct from call() above, which is hard-wired to /streams/. Returning + * the body (rather than a bool) is what GET /monitor/ needs. + * + * @param string $rMethod HTTP method. + * @param string $rPath Absolute daemon path, e.g. "/monitor/5". + * @param string|null $rBody JSON request body, if any. + * @return string|null Response body on 2xx, else null. + */ + private static function monitorCall(string $rMethod, string $rPath, ?string $rBody): ?string { + if (!function_exists('curl_init') || !defined('FANOUT_CTL_SOCK') || !file_exists(FANOUT_CTL_SOCK)) { + return null; + } + + $rCurl = curl_init(); + curl_setopt_array($rCurl, [ + CURLOPT_UNIX_SOCKET_PATH => FANOUT_CTL_SOCK, + CURLOPT_URL => 'http://localhost' . $rPath, + CURLOPT_CUSTOMREQUEST => $rMethod, + CURLOPT_RETURNTRANSFER => true, + CURLOPT_CONNECTTIMEOUT => 2, + CURLOPT_TIMEOUT => 3, + ]); + + if ($rBody !== null) { + curl_setopt($rCurl, CURLOPT_POSTFIELDS, $rBody); + curl_setopt($rCurl, CURLOPT_HTTPHEADER, ['Content-Type: application/json']); + } + + $rResponse = curl_exec($rCurl); + $rCode = curl_getinfo($rCurl, CURLINFO_HTTP_CODE); + curl_close($rCurl); + + if ($rCode < 200 || $rCode >= 300) { + return null; + } + return is_string($rResponse) ? $rResponse : ''; + } }