mirror of
https://github.com/Vateron-Media/XC_VM.git
synced 2026-10-04 04:02:30 +02:00
feat(fanout): reconcile supervised streams back into the panel
Closes the two gaps left by the supervision work: the panel could not see what the daemon knew, and it kept trying to start a watchdog for streams that already had one. The daemon has no database access -- XC_VM's credentials are not available to it -- so `cron:streams` is the only path by which its view reaches the panel. It now asks once per pass which streams the daemon supervises, then for each reads GET /monitor/<id> and writes back `stream_status`, `pid`, `current_source`, and the codecs, resolution and measured bitrate the daemon derived from the bytes it is fanning out. Only fields the daemon could actually determine are written: it reports an unknown as unknown, and overwriting a correct value with a blank would be worse than leaving it alone. `monitor_pid` names a PHP process, and isMonitorAlive checks the pid's /proc/<pid>/exe is the PHP binary -- so for a supervised stream it always answers false. Two callers drew the wrong conclusion from that: * cron:streams started a PHP monitor on every single pass, which then immediately stood down again. * admin/live.php, on the on-demand connect path, started one and then waited its full three seconds for a `_.monitor` file that would never appear. That is latency paid on the viewer's connect, for nothing. Both now ask who is watching before concluding nobody is. An unreachable daemon falls back to exactly the old behaviour in every case. That distinction is load-bearing: supervisedIDs() returns null rather than an empty list precisely so a blinking socket is not read as "supervising nothing", which would start a PHP monitor for every stream on the node at once. The supervised set is fetched once per cron pass rather than per stream, since the loop runs over every live stream on the node and each call is a socket round-trip. Verified: php -l clean on every changed file. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
8e24a2427a
commit
cea545fed0
@@ -36,6 +36,71 @@ class StreamsCronJob implements CommandInterface {
|
||||
return 'Cron: check live streams, monitors, on-demand, rogue PIDs';
|
||||
}
|
||||
|
||||
/**
|
||||
* Copy what the fanout daemon knows about a supervised stream back into
|
||||
* `streams_servers`, so the panel's own view stays true while the daemon owns
|
||||
* the process.
|
||||
*
|
||||
* The daemon cannot write to the database — XC_VM's credentials are not
|
||||
* available to it — so this is the only path by which its view reaches the
|
||||
* panel. It reports the codecs, picture size and MEASURED bitrate derived
|
||||
* from the bytes it is fanning out, which is both cheaper and more accurate
|
||||
* than the per-segment ffprobe this replaces.
|
||||
*
|
||||
* Only fields the daemon could actually determine are written. It reports an
|
||||
* unknown as unknown rather than guessing, and overwriting a correct value
|
||||
* with a blank would be worse than leaving it alone.
|
||||
*
|
||||
* @param object $db Database handle.
|
||||
* @param array $rStream The `streams_servers` row being reconciled.
|
||||
* @return void
|
||||
*/
|
||||
private function reconcileSupervisedStream($db, array $rStream): void {
|
||||
$rState = FanoutClient::monitorState((int) $rStream['stream_id']);
|
||||
if (!is_array($rState)) {
|
||||
return; // daemon unreachable, or it no longer holds this stream
|
||||
}
|
||||
|
||||
// stream_status: 0 = running, 1 = failed. The daemon knows which, and it
|
||||
// has given up entirely when `gave_up` is set.
|
||||
$rStatus = (!empty($rState['running']) && empty($rState['gave_up'])) ? 0 : 1;
|
||||
$rPID = isset($rState['pid']) ? (int) $rState['pid'] : 0;
|
||||
|
||||
$rSets = array('`stream_status` = ?');
|
||||
$rArgs = array($rStatus);
|
||||
if ($rPID > 0) {
|
||||
$rSets[] = '`pid` = ?';
|
||||
$rArgs[] = $rPID;
|
||||
}
|
||||
if (!empty($rState['source'])) {
|
||||
$rSets[] = '`current_source` = ?';
|
||||
$rArgs[] = (string) $rState['source'];
|
||||
}
|
||||
|
||||
$rMeta = isset($rState['meta']) && is_array($rState['meta']) ? $rState['meta'] : array();
|
||||
if (!empty($rMeta['video_codec'])) {
|
||||
$rSets[] = '`video_codec` = ?';
|
||||
$rArgs[] = (string) $rMeta['video_codec'];
|
||||
}
|
||||
if (!empty($rMeta['audio_codec'])) {
|
||||
$rSets[] = '`audio_codec` = ?';
|
||||
$rArgs[] = (string) $rMeta['audio_codec'];
|
||||
}
|
||||
if (!empty($rMeta['height'])) {
|
||||
$rSets[] = '`resolution` = ?';
|
||||
$rArgs[] = StreamSorter::getNearest(array(240, 360, 480, 576, 720, 1080, 1440, 2160), (int) $rMeta['height']);
|
||||
}
|
||||
if (!empty($rMeta['bitrate_kbps'])) {
|
||||
$rSets[] = '`bitrate` = ?';
|
||||
$rArgs[] = (int) $rMeta['bitrate_kbps'];
|
||||
}
|
||||
|
||||
$rArgs[] = $rStream['server_stream_id'];
|
||||
$db->query('UPDATE `streams_servers` SET ' . implode(', ', $rSets) . ' WHERE `server_stream_id` = ?', ...$rArgs);
|
||||
|
||||
echo 'Supervised by daemon (pid ' . $rPID . ', ' . ($rStatus === 0 ? 'running' : 'failed') . ")\n";
|
||||
}
|
||||
|
||||
public function execute(array $rArgs): int {
|
||||
if (!$this->assertRunAsXcVm()) {
|
||||
return 1;
|
||||
@@ -62,6 +127,22 @@ class StreamsCronJob implements CommandInterface {
|
||||
$rActivePIDs = array();
|
||||
$rStreamIDs = array();
|
||||
|
||||
// Which streams the fanout daemon is supervising, asked once per pass
|
||||
// rather than once per stream: it is a socket round-trip, and this loop
|
||||
// runs over every live stream on the node.
|
||||
//
|
||||
// null means the daemon could not be asked. That is deliberately NOT the
|
||||
// same as "it supervises nothing": treating an unreachable daemon as
|
||||
// supervising nothing would make this cron start a PHP monitor for every
|
||||
// stream on the node the moment the socket blinked.
|
||||
$rSupervised = FanoutClient::supervisedIDs();
|
||||
$rSupervisedSet = array();
|
||||
if (is_array($rSupervised)) {
|
||||
foreach ($rSupervised as $rID) {
|
||||
$rSupervisedSet[(string) $rID] = true;
|
||||
}
|
||||
}
|
||||
|
||||
if ($rRedis) {
|
||||
$db->query('SELECT t2.stream_display_name, t1.stream_started, t1.stream_info, t2.fps_restart, t1.stream_status, t1.progress_info, t1.stream_id, t1.monitor_pid, t1.on_demand, t1.server_stream_id, t1.pid, servers_attached.attached, t2.vframes_server_id, t2.vframes_pid, t2.tv_archive_server_id, t2.tv_archive_pid FROM `streams_servers` t1 INNER JOIN `streams` t2 ON t2.id = t1.stream_id AND t2.direct_source = 0 INNER JOIN `streams_types` t3 ON t3.type_id = t2.type LEFT JOIN (SELECT `stream_id`, COUNT(*) AS `attached` FROM `streams_servers` WHERE `parent_id` = ? AND `pid` IS NOT NULL AND `pid` > 0 AND `monitor_pid` IS NOT NULL AND `monitor_pid` > 0) AS `servers_attached` ON `servers_attached`.`stream_id` = t1.`stream_id` WHERE (t1.pid IS NOT NULL OR t1.stream_status <> 0 OR t1.to_analyze = 1) AND t1.server_id = ? AND t3.live = 1', SERVER_ID, SERVER_ID);
|
||||
} else {
|
||||
@@ -73,7 +154,18 @@ class StreamsCronJob implements CommandInterface {
|
||||
echo 'Stream ID: ' . $rStream['stream_id'] . "\n";
|
||||
$rStreamIDs[] = $rStream['stream_id'];
|
||||
|
||||
if (ProcessManager::isMonitorAlive($rStream['monitor_pid'], $rStream['stream_id']) || $rStream['on_demand']) {
|
||||
// `monitor_pid` names a PHP process, and isMonitorAlive checks the
|
||||
// pid's /proc/<pid>/exe is the PHP binary. For a stream the daemon
|
||||
// supervises there IS no PHP monitor, so that check answers false
|
||||
// and the branch below would spawn one on every single pass — which
|
||||
// would then immediately stand down again. Something IS watching the
|
||||
// stream; ask who before concluding nobody is.
|
||||
$rDaemonMonitored = isset($rSupervisedSet[(string) $rStream['stream_id']]);
|
||||
if ($rDaemonMonitored) {
|
||||
$this->reconcileSupervisedStream($db, $rStream);
|
||||
}
|
||||
|
||||
if ($rDaemonMonitored || ProcessManager::isMonitorAlive($rStream['monitor_pid'], $rStream['stream_id']) || $rStream['on_demand']) {
|
||||
if ($rStream['on_demand'] == 1 && $rStream['attached'] == 0) {
|
||||
if ($rRedis) {
|
||||
$rCount = 0;
|
||||
|
||||
@@ -11,6 +11,7 @@ use XcVm\Domain\Stream\AdminStreamToken;
|
||||
use XcVm\Domain\Stream\ConnectionTracker;
|
||||
use XcVm\Domain\Stream\StreamProcess;
|
||||
use XcVm\Infrastructure\Database\DatabaseFactory;
|
||||
use XcVm\Streaming\Fanout\FanoutClient;
|
||||
|
||||
/**
|
||||
* Admin live stream handler
|
||||
@@ -86,7 +87,14 @@ if (0 < $db->num_rows()) {
|
||||
$rChannelInfo['pid'] = null;
|
||||
|
||||
if ($rChannelInfo['on_demand'] == 1) {
|
||||
if (!ProcessManager::isMonitorAlive($rChannelInfo['monitor_pid'], $rStreamID)) {
|
||||
// A stream the fanout daemon supervises already has a watchdog; it will
|
||||
// restart the encoder itself. Starting a PHP monitor here would stand
|
||||
// down immediately and never write `_.monitor`, so the wait below would
|
||||
// burn its full three seconds on every on-demand request for that
|
||||
// stream — latency paid on the viewer's connect, for nothing.
|
||||
if (FanoutClient::isSupervised($rStreamID)) {
|
||||
$rChannelInfo['monitor_pid'] = null;
|
||||
} elseif (!ProcessManager::isMonitorAlive($rChannelInfo['monitor_pid'], $rStreamID)) {
|
||||
StreamProcess::startMonitor($rStreamID);
|
||||
|
||||
for ($rRetries = 0; !file_exists(STREAMS_PATH . intval($rStreamID) . '_.monitor') && $rRetries < 300; $rRetries++) {
|
||||
|
||||
Reference in New Issue
Block a user