feat(fanout): hand live encoders to the daemon to supervise

Enables the daemon-side supervisor (XC_VM_Fanout M1, its ADR 0002): the
daemon starts, watches and restarts a stream's ffmpeg instead of a
resident PHP watchdog per channel doing it.

The daemon cannot read this panel's configuration -- database credentials
live inside the XC_VM C extension and are never handed out, not even to
PHP (Core/Config/ConfigReader.php) -- so the split is process ownership,
not config ownership. buildLive keeps composing the command here from the
database; the daemon only ever runs the command it is handed.

* buildLive no longer appends its own redirect-and-background tail. A
  backgrounded command leaves the daemon nothing to supervise: it has to
  be the process's parent to reap it and to signal it. The tail moves to
  liveRedirectTail(), which the legacy path appends itself, so that path
  is byte-for-byte what it was.

* startStream hands the bare command to the daemon when this node
  supervises there, passing the panel's own restart policy
  (stop_failures, stream_fail_sleep, on_demand, on_demand_failure_exit).
  The daemon writes the same <id>_.pid file, so stopStream,
  isStreamRunning and everything else downstream keep working untouched.

* FanoutClient gains supervise / releaseSupervision / monitorState /
  isSupervised over the existing control socket.

* MonitorCommand stands down when the daemon holds the stream. Two
  watchdogs on one ffmpeg fight -- each reads the other's kill as a
  failure and restarts -- so the channel would flap forever. The check
  asks the daemon rather than assuming, and an unreachable daemon answers
  false, so this watchdog keeps working whenever the daemon is not
  actually holding the process.

Off by default: daemonSupervises() reads the `daemon_supervise` setting,
absent on every existing install, so behaviour is unchanged until an
operator sets it. No migration, and clearing it is the rollback. Delay
streams are excluded for now (they write their own playlist).

Verified: php -l clean on the changed files and across all of src/.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
obscuremind
2026-09-10 19:15:17 +01:00
co-authored by Claude Opus 5
parent 34fb54973b
commit 56c0619e1e
3 changed files with 195 additions and 2 deletions
+16
View File
@@ -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 <stream_id> [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')) {
+57 -2
View File
@@ -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
* `<streams>/<id>_.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
+122
View File
@@ -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/<id>. Returning
* the body (rather than a bool) is what GET /monitor/<id> 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 : '';
}
}