mirror of
https://github.com/Vateron-Media/XC_VM.git
synced 2026-10-04 04:02:30 +02:00
feat(fanout): health policy, source forcing and supervision reconcile
Completes the panel side of moving the per-stream watchdog into the daemon (XC_VM_Fanout M5). * daemonHealthPolicy() translates this stream's watchdog settings -- seg_time, audio_restart_loss, fps_restart/fps_threshold/fps_delay and auto_restart -- into the daemon's health policy. These are the same conditions MonitorCommand.php checked in its inner loop; the daemon now judges them against the bytes it is already fanning out instead of by hashing a playlist, ffprobing a segment and reading a progress file. Each is omitted when the panel has it switched off. * supervise() now sends adopt_match: the stream's own HLS path, which appears in every live command buildLive produces and in no other process. It is how the daemon recognises an encoder of ours that outlived it, so a daemon restart resumes watching the running ffmpeg rather than starting a second one beside it. Sending a token the panel knows, rather than letting the daemon guess, is what makes that safe against pid reuse. * FanoutClient::supervisedIDs() reports what the daemon currently holds, so the panel can notice a restarted daemon and hand its streams back. Null (not an empty array) when unreachable, so a dead socket is not read as every stream needing re-registration. * FanoutClient::forceSource() replaces writing a <id>.force signal file. Still gated behind the `daemon_supervise` setting, absent on every existing install. Verified: php -l clean on the changed files. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
56c0619e1e
commit
8244ece100
@@ -956,6 +956,53 @@ class StreamProcess {
|
||||
return (bool) SettingsManager::get('daemon_supervise');
|
||||
}
|
||||
|
||||
/**
|
||||
* Translate this stream's watchdog settings into the daemon's health policy.
|
||||
*
|
||||
* These are the same conditions MonitorCommand.php checked in its inner loop;
|
||||
* the daemon judges them against the bytes it is already fanning out instead
|
||||
* of by hashing a playlist, ffprobing a segment and reading a progress file.
|
||||
* Each is omitted when the panel has it switched off, and an omitted check is
|
||||
* simply not made.
|
||||
*
|
||||
* @param array $rStream Stream row (stream_info + server_info).
|
||||
* @param array $rSettings Resolved panel settings.
|
||||
* @return array Health policy for FanoutClient::supervise().
|
||||
*/
|
||||
private static function daemonHealthPolicy($rStream, $rSettings): array {
|
||||
$rHealth = array();
|
||||
|
||||
// Stall: the panel restarted when the playlist stopped changing for
|
||||
// seg_time * 6. Same window, measured from the last byte published.
|
||||
$rSegTime = max(1, intval($rSettings['seg_time']));
|
||||
$rHealth['stall_sec'] = $rSegTime * 6;
|
||||
|
||||
if (!empty($rSettings['audio_restart_loss'])) {
|
||||
// The panel could only notice on its 300s ffprobe cycle; the daemon
|
||||
// sees the audio PID go quiet, so the window can be a real one.
|
||||
$rHealth['audio_loss_sec'] = 30;
|
||||
}
|
||||
|
||||
if (!empty($rStream['stream_info']['fps_restart'])) {
|
||||
// fps_threshold is a percentage in the panel and a fraction here.
|
||||
$rThreshold = floatval($rStream['stream_info']['fps_threshold'] ?: 100) / 100.0;
|
||||
if ($rThreshold > 0 && $rThreshold < 1) {
|
||||
$rHealth['fps_threshold'] = $rThreshold;
|
||||
$rHealth['fps_grace_sec'] = intval($rSettings['fps_delay']);
|
||||
}
|
||||
}
|
||||
|
||||
$rAutoRestart = json_decode((string) $rStream['stream_info']['auto_restart'], true);
|
||||
if (!empty($rAutoRestart['days']) && !empty($rAutoRestart['at'])) {
|
||||
$rHealth['auto_restart'] = array(
|
||||
'days' => array_values((array) $rAutoRestart['days']),
|
||||
'at' => (string) $rAutoRestart['at'],
|
||||
);
|
||||
}
|
||||
|
||||
return $rHealth;
|
||||
}
|
||||
|
||||
public static function createChannelItem($rStreamID, $rSource) {
|
||||
global $rSettings, $rServers, $rFFMPEG_CPU, $rFFMPEG_GPU;
|
||||
$db = self::db();
|
||||
@@ -1569,12 +1616,18 @@ class StreamProcess {
|
||||
// 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'],
|
||||
));
|
||||
$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'],
|
||||
),
|
||||
self::daemonHealthPolicy($rStream, $rSettings)
|
||||
);
|
||||
}
|
||||
if (!$rHandedOver) {
|
||||
$rFFMPEG .= self::liveRedirectTail($rStreamID);
|
||||
|
||||
@@ -477,17 +477,24 @@ class FanoutClient {
|
||||
* 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 {
|
||||
public static function supervise(int $rStreamID, string $rCmd, string $rLabel, array $rPolicy = [], array $rHealth = []): bool {
|
||||
if ($rCmd === '') {
|
||||
return false;
|
||||
}
|
||||
$rBody = json_encode([
|
||||
'sources' => [['label' => $rLabel, 'cmd' => $rCmd]],
|
||||
'policy' => $rPolicy,
|
||||
'health' => (object) $rHealth,
|
||||
'pid_path' => STREAMS_PATH . $rStreamID . '_.pid',
|
||||
'errors_path' => STREAMS_PATH . $rStreamID . '.errors',
|
||||
'log_path' => LOGS_TMP_PATH . 'stream_log.log',
|
||||
'server_id' => SERVER_ID,
|
||||
// How the daemon recognises an encoder of OURS that outlived it, so a
|
||||
// daemon restart resumes watching the running ffmpeg instead of
|
||||
// starting a second one beside it. The stream's own HLS path appears
|
||||
// in every live command buildLive produces and in no other process,
|
||||
// which is what makes it safe against pid reuse.
|
||||
'adopt_match' => STREAMS_PATH . $rStreamID . '_.m3u8',
|
||||
]);
|
||||
if ($rBody === false) {
|
||||
return false;
|
||||
@@ -538,6 +545,45 @@ class FanoutClient {
|
||||
return is_array($rState) && !empty($rState['supervised']);
|
||||
}
|
||||
|
||||
/**
|
||||
* Every stream this node's daemon is currently supervising.
|
||||
*
|
||||
* A daemon that restarted comes back supervising nothing, so the panel has to
|
||||
* hand its streams over again — at which point the daemon adopts the ffmpeg
|
||||
* processes that outlived it rather than starting duplicates. This is how the
|
||||
* panel notices there is anything to hand over.
|
||||
*
|
||||
* Returns null (not an empty array) when the daemon is unreachable, so a
|
||||
* caller can tell "supervising nothing" apart from "cannot be asked" and does
|
||||
* not treat a dead socket as every stream needing re-registration.
|
||||
*
|
||||
* @return array|null Stream ids as strings, or null if unreachable.
|
||||
*/
|
||||
public static function supervisedIDs(): ?array {
|
||||
$rBody = self::monitorCall('GET', '/monitors', null);
|
||||
if ($rBody === null) {
|
||||
return null;
|
||||
}
|
||||
$rJson = json_decode($rBody, true);
|
||||
return is_array($rJson) ? $rJson : null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Force a supervised stream onto a specific source (its index in the
|
||||
* stream_source list). Replaces writing a `<id>.force` signal file.
|
||||
*
|
||||
* @param int $rStreamID Stream id.
|
||||
* @param int $rIndex Source index.
|
||||
* @return bool True when the daemon accepted the switch.
|
||||
*/
|
||||
public static function forceSource(int $rStreamID, int $rIndex): bool {
|
||||
$rBody = json_encode(['index' => $rIndex]);
|
||||
if ($rBody === false) {
|
||||
return false;
|
||||
}
|
||||
return self::monitorCall('POST', '/monitor/' . $rStreamID . '/source', $rBody) !== null;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
|
||||
Reference in New Issue
Block a user