diff --git a/src/Cli/Commands/DelayCommand.php b/src/Cli/Commands/DelayCommand.php index 3b04927b..88b4a1e9 100644 --- a/src/Cli/Commands/DelayCommand.php +++ b/src/Cli/Commands/DelayCommand.php @@ -6,6 +6,7 @@ use XcVm\Cli\CommandInterface; use XcVm\Core\Config\SettingsManager; use XcVm\Core\Process\ProcessManager; use XcVm\Domain\Stream\StreamProcess; +use XcVm\Streaming\Fanout\IngestFeeder; /** * DelayCommand — delay command @@ -77,8 +78,22 @@ class DelayCommand implements CommandInterface { if (file_exists($rPlaylistOld)) { $rOldSegments = $this->getSegments($rPlaylistOld, -1); } + // Clients are served only by the xc_fanout daemon (ADR 0003, Phase E), and a + // delayed stream's encoder output is the undelayed one — so nothing fed the + // daemon the delayed stream and a delayed channel could not be watched at + // all. The segments this worker publishes are now pushed into the daemon's + // ingest as they go out, paced over their duration so TS viewers get a + // steady stream; the daemon re-segments them for HLS. + $rFeeder = IngestFeeder::forStream($rStreamID, (bool) SettingsManager::get('encrypt_hls'), static function (string $rLine) use ($rStreamID) { + @file_put_contents(STREAMS_PATH . $rStreamID . '.errors', '[Delay] ' . $rLine . "\n", FILE_APPEND | LOCK_EX); + }); + $rFeeder->connect(); + $rFedSegment = null; + $rFeedQueue = array(); + $rFeedCurrent = null; + $rPrevMD5 = null; - $rMD5 = md5(file_get_contents($rPlaylistDelay)); + $rMD5 = md5((string) @file_get_contents($rPlaylistDelay)); while (ProcessManager::isStreamRunning($rPID, $rStreamID) && file_exists($rPlaylistDelay)) { if ($rMD5 != $rPrevMD5) { if (file_exists(STREAMS_PATH . $rStreamID . '_.dur')) { @@ -103,18 +118,110 @@ class DelayCommand implements CommandInterface { $rData .= '#EXTINF:' . $rSegment['seconds'] . ',' . "\n" . $rSegment['file'] . "\n"; } file_put_contents($rPlaylist, $rData, LOCK_EX); + $this->queueForDaemon($rM3U8['segments'], $rFedSegment, $rFeedQueue); $rMD5 = $rPrevMD5; $this->deleteSegments($rStreamID, $rSequence - 2); $this->cleanUpSegments($rStreamID, $rDelayDuration); } } - usleep(1000); - $rPrevMD5 = md5(file_get_contents($rPlaylistDelay)); + $this->pumpDaemon($rFeeder, $rFeedQueue, $rFeedCurrent); + // 50 ms: fine enough to pace the daemon feed and to publish a new + // delayed segment promptly (this used to spin every 1 ms, hashing the + // playlist a thousand times a second). + usleep(50000); + $rPrevMD5 = md5((string) @file_get_contents($rPlaylistDelay)); } + $rFeeder->close(); // the daemon keeps the stream; viewers wait for the restart's feed return 0; } + /** Segments sent at once when the worker starts, so viewers get data immediately. */ + private const FEED_SEED_SEGMENTS = 2; + + /** Queued segments past which the feed stops pacing and catches up. */ + private const FEED_BACKLOG_SEGMENTS = 3; + + /** + * Queue the segments just published that the daemon has not been fed yet (all + * but the newest FEED_SEED_SEGMENTS are skipped on the worker's first pass). + * + * @param array $rSegments Published segments, oldest first ({seconds, file}). + * @param int|null $rFedSegment Highest segment number queued so far. + * @param array $rQueue Pending {data, dur, burst} entries. + * @return void + */ + private function queueForDaemon(array $rSegments, ?int &$rFedSegment, array &$rQueue): void { + $rNew = array(); + foreach ($rSegments as $rSegment) { + if (preg_match('/_(\d+)\.ts$/', (string) ($rSegment['file'] ?? ''), $rMatch)) { + $rNumber = intval($rMatch[1]); + if ($rFedSegment === null || $rNumber > $rFedSegment) { + $rNew[$rNumber] = $rSegment; + } + } + } + if (count($rNew) === 0) { + return; + } + ksort($rNew); + $rSeed = ($rFedSegment === null); + if ($rSeed) { + $rNew = array_slice($rNew, -self::FEED_SEED_SEGMENTS, null, true); + } + foreach ($rNew as $rNumber => $rSegment) { + $rData = @file_get_contents(STREAMS_PATH . $rSegment['file']); + if (!is_string($rData) || $rData === '') { + $rData = @file_get_contents(DELAY_PATH . $rSegment['file']); + } + if (is_string($rData) && strlen($rData) >= 188) { + $rData = substr($rData, 0, strlen($rData) - strlen($rData) % 188); // whole packets + $rQueue[] = array('data' => $rData, 'dur' => max(0.5, floatval($rSegment['seconds'])), 'burst' => $rSeed); + } + $rFedSegment = $rNumber; + } + } + + /** + * Feed the daemon: the current segment is released in whole packets spread + * over 90% of its duration (so the feed never falls behind the playlist); the + * start-up seed, or a backlog of queued segments, is sent at once. + * + * @param IngestFeeder $rFeeder The daemon feed. + * @param array $rQueue Pending {data, dur, burst} entries. + * @param array|null $rCurrent The segment being paced ({data, sent, start, dur, burst}). + * @return void + */ + private function pumpDaemon(IngestFeeder $rFeeder, array &$rQueue, ?array &$rCurrent): void { + $rNow = microtime(true); + if ($rCurrent === null && count($rQueue) > 0) { + $rItem = array_shift($rQueue); + $rCurrent = array('data' => $rItem['data'], 'sent' => 0, 'start' => $rNow, 'dur' => $rItem['dur'], 'burst' => $rItem['burst']); + } + if ($rCurrent === null) { + $rFeeder->flush(); + return; + } + + $rLength = strlen($rCurrent['data']); + if ($rCurrent['burst'] || count($rQueue) >= self::FEED_BACKLOG_SEGMENTS) { + $rTarget = $rLength; + } else { + $rTarget = (int) min($rLength, ceil($rLength * ($rNow - $rCurrent['start']) / ($rCurrent['dur'] * 0.9))); + $rTarget -= $rTarget % 188; + } + $rChunk = $rTarget - $rCurrent['sent']; + if ($rChunk > 0) { + $rFeeder->write(substr($rCurrent['data'], $rCurrent['sent'], $rChunk)); + $rCurrent['sent'] += $rChunk; + } else { + $rFeeder->flush(); + } + if ($rCurrent['sent'] >= $rLength) { + $rCurrent = null; + } + } + private function cleanUpSegments($rStreamID, $rDelayDuration): void { shell_exec('find ' . DELAY_PATH . intval($rStreamID) . '_*' . ' -type f -cmin +' . $rDelayDuration . ' -delete'); } diff --git a/src/Cli/Commands/LlodCommand.php b/src/Cli/Commands/LlodCommand.php index 326c33f4..89b5e840 100644 --- a/src/Cli/Commands/LlodCommand.php +++ b/src/Cli/Commands/LlodCommand.php @@ -3,7 +3,10 @@ namespace XcVm\Cli\Commands; use XcVm\Cli\CommandInterface; +use XcVm\Core\Process\ProcessManager; +use XcVm\Core\Util\StreamUtils; use XcVm\Streaming\Fanout\FanoutClient; +use XcVm\Streaming\Fanout\IngestFeeder; /** * LlodCommand — llod command @@ -112,12 +115,12 @@ class LlodCommand implements CommandInterface { echo "Starting LLOD processing...\n\n"; - $this->startLlod($rStreamID, $rStreamSources, $rStreamArguments, $rRequestPrebuffer, $rSegListSize, $rSegDeleteThreshold, $rSegTime, $rSegmentStatus, $rFP, $rSegmentFile); + $this->startLlod($rStreamID, $rStreamSources, $rStreamArguments, $rRequestPrebuffer, $rSegListSize, $rSegDeleteThreshold, $rSegTime, !empty($rSettings['encrypt_hls']), $rSegmentStatus, $rFP, $rSegmentFile); return 0; } - private function startLlod($rStreamID, $rStreamSources, $rStreamArguments, $rRequestPrebuffer, $rSegListSize, $rSegDeleteThreshold, $rSegTime, &$rSegmentStatus, &$rFP, &$rSegmentFile): void { + private function startLlod($rStreamID, $rStreamSources, $rStreamArguments, $rRequestPrebuffer, $rSegListSize, $rSegDeleteThreshold, $rSegTime, $rEncryptHLS, &$rSegmentStatus, &$rFP, &$rSegmentFile): void { // Keyframe-aligned segmentation. The previous version cut segments on a // fixed wall-clock timer (SEGMENT_DURATION) at an arbitrary TS-packet // boundary, so a segment frequently started mid-GOP (no IDR keyframe) or @@ -139,20 +142,8 @@ class LlodCommand implements CommandInterface { } } - $ua = $rStreamArguments['user_agent']['value'] ?? 'Mozilla/5.0'; - - $context = stream_context_create([ - 'http' => [ - 'timeout' => TIMEOUT, - 'user_agent' => $ua, - ], - 'ssl' => [ - 'verify_peer' => false, - 'verify_peer_name' => false, - ] - ]); - - $rFP = $this->getActiveStream($rStreamID, $rStreamSources, $context); + $rSniffed = ''; + $rFP = $this->getActiveStream($rStreamID, $rStreamSources, $rStreamArguments, $rRequestPrebuffer, $rSniffed); if (!$rFP) { echo "No active stream\n"; return; @@ -160,21 +151,19 @@ class LlodCommand implements CommandInterface { stream_set_blocking($rFP, true); - // LLOD v3 daemon feed (ADR 0003). LLOD reads MPEG-TS itself (no ffmpeg), - // so mirror the raw bytes into the xc_fanout daemon's push-fed ingest - // socket — the same stream then fans out via /live/ and in-RAM /hls, - // and live.php's isStreamFed() routes viewers to the daemon. The write is - // non-blocking + best-effort so a daemon stall never slows LLOD's own - // segmenting (its primary job); the daemon resyncs on PAT/PMT after any - // dropped bytes. Null / failed connect ⇒ legacy-only, no behaviour change. - $rDaemonSock = FanoutClient::registerIngest($rStreamID); - $rDaemonConn = null; - if ($rDaemonSock !== null) { - $rDaemonConn = @stream_socket_client('unix://' . $rDaemonSock, $rDErrno, $rDErrstr, 2); - if ($rDaemonConn) { - stream_set_blocking($rDaemonConn, false); - } - } + // LLOD v3 daemon feed (ADR 0003). LLOD reads MPEG-TS itself (no ffmpeg) and + // pushes it into the xc_fanout daemon's ingest socket: the daemon is what + // serves this channel's viewers (/live/ and the in-RAM /hls) — there is + // no other client path since Phase E, so the feed is the delivery, and the + // on-disk segments below only serve timeshift/thumbnails/the monitor. + // IngestFeeder keeps short writes from tearing packets, re-registers and + // redials after a daemon restart, and carries the HLS key when encrypted + // HLS is on (the playlist declares it; without the key the daemon served + // plain segments no player could decrypt). + $rFeeder = IngestFeeder::forStream($rStreamID, !empty($rEncryptHLS), function (string $rLine) use ($rStreamID) { + $this->writeError($rStreamID, '[LLOD] ' . $rLine); + }); + $rFeeder->connect(); shell_exec('rm -f ' . STREAMS_PATH . escapeshellarg($rStreamID) . '_*.ts'); @@ -191,12 +180,13 @@ class LlodCommand implements CommandInterface { $lastData = time(); $firstDataAt = microtime(true); - $buffer = ''; + $buffer = $rSniffed; // bytes read while identifying the source as MPEG-TS while (!feof($rFP)) { $data = fread($rFP, BUFFER_SIZE); if ($data === '' || $data === false) { + $rFeeder->flush(); // drain a backlog / reconnect while the source is quiet if (time() - $lastData > TIMEOUT) { $this->writeError($rStreamID, '[LLOD] stream timeout'); break; @@ -208,17 +198,9 @@ class LlodCommand implements CommandInterface { $lastData = time(); $buffer .= $data; - // Mirror to the daemon (best-effort, non-blocking). On a write error - // (daemon gone) stop mirroring; live.php then falls back to legacy. - if ($rDaemonConn) { - if (@fwrite($rDaemonConn, $data) === false) { - @fclose($rDaemonConn); - $rDaemonConn = null; - } - } - $len = strlen($buffer); $off = 0; + $rFeed = ''; // whole packets for the daemon // Process only whole 188-byte TS packets; keep any remainder buffered. while ($len - $off >= PACKET_SIZE) { @@ -239,6 +221,7 @@ class LlodCommand implements CommandInterface { $pkt = substr($buffer, $off, PACKET_SIZE); $off += PACKET_SIZE; + $rFeed .= $pkt; // every aligned packet goes to the daemon, from the first $hdr = $this->parseTsHeader($pkt); $pid = $hdr['pid']; @@ -304,6 +287,8 @@ class LlodCommand implements CommandInterface { fwrite($rSegmentFile, $pkt); } + $rFeeder->write($rFeed); + // Retain the partial trailing packet for the next read. $buffer = ($off >= $len) ? '' : substr($buffer, $off); } @@ -314,9 +299,7 @@ class LlodCommand implements CommandInterface { if (is_resource($rFP)) { fclose($rFP); } - if (is_resource($rDaemonConn)) { - fclose($rDaemonConn); - } + $rFeeder->close(); // Drop the daemon ingest on a clean exit; stopStream() is the backstop // when LLOD is killed mid-loop. FanoutClient::unregister($rStreamID); @@ -454,13 +437,93 @@ class LlodCommand implements CommandInterface { return $pcrPid !== 0x1FFF ? $pcrPid : null; } - private function getActiveStream($rStreamID, $rURLs, $rContext) { + /** + * The HTTP stream context for one source, from the stream's fetch arguments — + * the same user agent, extra headers, cookie and proxy the ffmpeg path sends + * (it used to send only a user agent, and no default one when the argument + * was left empty), plus the prebuffer request an XC_VM source understands + * when request_prebuffer is on. + * + * @param string $rURL Source URL. + * @param array $rStreamArguments argument_key => {value, argument_default_value}. + * @param mixed $rRequestPrebuffer The request_prebuffer setting. + * @return resource + */ + private function sourceContext($rURL, $rStreamArguments, $rRequestPrebuffer) { + $rArg = static function (string $rKey) use ($rStreamArguments): string { + $rValue = $rStreamArguments[$rKey]['value'] ?? ''; + if ($rValue === '' || $rValue === null) { + $rValue = $rStreamArguments[$rKey]['argument_default_value'] ?? ''; + } + return trim((string) $rValue); + }; + + $rHeaders = array(); + foreach (preg_split('/\r\n|\r|\n/', $rArg('headers')) as $rLine) { + if (trim($rLine) !== '' && strpos($rLine, ':') !== false) { + $rHeaders[] = trim($rLine); + } + } + if ($rArg('cookie') !== '') { + $rHeaders[] = 'Cookie: ' . $rArg('cookie'); + } + if (!empty($rRequestPrebuffer) && StreamUtils::detectXC_VM($rURL)) { + $rHeaders[] = 'X-XC_VM-Prebuffer: 1'; + } + + $rHTTP = array( + 'timeout' => TIMEOUT, + 'user_agent' => ($rArg('user_agent') !== '' ? $rArg('user_agent') : 'Mozilla/5.0'), + ); + if (count($rHeaders) > 0) { + $rHTTP['header'] = implode("\r\n", $rHeaders); + } + $rProxy = $rArg('proxy'); + if ($rProxy !== '') { + $rHTTP['proxy'] = (strpos($rProxy, '://') === false ? 'tcp://' : '') . $rProxy; + $rHTTP['request_fulluri'] = true; + } + + return stream_context_create(array( + 'http' => $rHTTP, + 'ssl' => array('verify_peer' => false, 'verify_peer_name' => false), + )); + } + + /** + * Read up to $rLength bytes from a just-opened source to recognise MPEG-TS by + * its content: two sync bytes one packet apart. + * + * @param resource $rFP Source stream. + * @param string $rSniffed Receives the bytes read (they are part of the stream). + * @return bool + */ + private function looksLikeMpegTs($rFP, string &$rSniffed): bool { + $rSniffed = ''; + $rDeadline = time() + TIMEOUT; + while (strlen($rSniffed) < PACKET_SIZE * 4 && !feof($rFP) && time() < $rDeadline) { + $rChunk = fread($rFP, PACKET_SIZE * 4 - strlen($rSniffed)); + if ($rChunk === false || $rChunk === '') { + usleep(10000); + continue; + } + $rSniffed .= $rChunk; + } + for ($i = 0; $i + PACKET_SIZE < strlen($rSniffed) && $i < PACKET_SIZE; $i++) { + if ($rSniffed[$i] === "\x47" && $rSniffed[$i + PACKET_SIZE] === "\x47") { + return true; + } + } + return false; + } + + private function getActiveStream($rStreamID, $rURLs, $rStreamArguments, $rRequestPrebuffer, string &$rSniffed) { echo "Trying to get active stream from " . count($rURLs) . " URL(s)\n"; foreach ($rURLs as $index => $rURL) { echo "\nAttempting source " . ($index + 1) . "/" . count($rURLs) . ": $rURL\n"; - $rFP = @fopen($rURL, 'rb', false, $rContext); + $rFP = @fopen($rURL, 'rb', false, $this->sourceContext($rURL, $rStreamArguments, $rRequestPrebuffer)); if ($rFP) { echo "Connection successful\n"; @@ -490,7 +553,13 @@ class LlodCommand implements CommandInterface { echo " $key: $value\n"; } - $rContentType = $rHeaders['Content-Type'] ?? ''; + // Header names are case-insensitive (HTTP/2-style lower-case is common). + $rContentType = ''; + foreach ($rHeaders as $rKey => $rValue) { + if (is_string($rKey) && strcasecmp($rKey, 'Content-Type') === 0) { + $rContentType = $rValue; + } + } echo "Content-Type: $rContentType\n"; if (stripos($rContentType, 'video/mp2t') !== false) { @@ -499,7 +568,17 @@ class LlodCommand implements CommandInterface { return $rFP; } - $contentTypeInfo = $rHeaders['Content-Type'] ?? 'unknown'; + // Many TS sources label it application/octet-stream, or nothing at + // all: recognise it by content rather than refusing a playable + // source. An HLS playlist or an HTML error page fails the check. + if ($this->looksLikeMpegTs($rFP, $rSniffed)) { + echo "Content is MPEG-TS (labelled '$rContentType')\n"; + echo "=== getActiveStream() successful ===\n\n"; + return $rFP; + } + $rSniffed = ''; + + $contentTypeInfo = ($rContentType !== '' ? $rContentType : 'unknown'); $this->writeError($rStreamID, "[LLOD] Source isn't MPEG-TS: " . $rURL . ' - ' . $contentTypeInfo); fclose($rFP); } else { @@ -596,7 +675,11 @@ class LlodCommand implements CommandInterface { $m3u8 .= "{$rStreamID}_{$seg}.ts\n"; } - if (@file_put_contents(STREAMS_PATH . $rStreamID . '_.m3u8', $m3u8, LOCK_EX) === false) { + // Write-then-rename: readers (the monitor's staleness check, thumbnails, + // timeshift) never see a half-written playlist. + $rTmp = STREAMS_PATH . $rStreamID . '_.m3u8.tmp'; + if (@file_put_contents($rTmp, $m3u8) === false || !@rename($rTmp, STREAMS_PATH . $rStreamID . '_.m3u8')) { + @unlink($rTmp); $this->writeError($rStreamID, '[LLOD] Failed to write playlist file'); return; } @@ -611,41 +694,18 @@ class LlodCommand implements CommandInterface { @file_put_contents(STREAMS_PATH . $rStreamID . '.errors', $logMessage, FILE_APPEND | LOCK_EX); } + /** + * One segmenter per stream: end any other LLOD process for this stream (one a + * killed monitor left behind). This used to read `_.monitor`, which holds the + * MONITOR's pid — never an LLOD one — so an old segmenter was never found and + * two could write the same segment files and feed the daemon at once. + */ private function checkRunning($rStreamID): void { echo "Checking for existing process for stream $rStreamID\n"; - clearstatcache(true); - $monitorFile = STREAMS_PATH . $rStreamID . '_.monitor'; - $rPID = null; - if (file_exists($monitorFile)) { - $rPID = intval(file_get_contents($monitorFile)); - echo "Monitor file found, PID: $rPID\n"; - } else { - echo "No monitor file found\n"; - } - if (empty($rPID)) { - $killCmd = "kill -9 `ps -ef | grep 'LLOD\\[" . intval($rStreamID) . "\\]' | grep -v grep | awk '{print \$2}'`"; - echo "No PID from monitor, executing kill command: $killCmd\n"; - shell_exec($killCmd); - } else { - if (file_exists('/proc/' . $rPID)) { - echo "Process directory exists: /proc/$rPID\n"; - $cmdlineFile = '/proc/' . $rPID . '/cmdline'; - if (file_exists($cmdlineFile)) { - $rCommand = trim(file_get_contents($cmdlineFile)); - echo "Process command line: $rCommand\n"; - $expectedCommand = 'LLOD[' . $rStreamID . ']'; - if ($rCommand === $expectedCommand && 0 < $rPID) { - echo "Killing existing process PID: $rPID\n"; - posix_kill($rPID, 9); - } else { - echo "Process command doesn't match expected: '$rCommand' != '$expectedCommand'\n"; - } - } else { - echo "Command line file not found\n"; - } - } else { - echo "Process directory doesn't exist, process not running\n"; - } + $rTerms = array('LLOD[' . intval($rStreamID) . ']', 'console.php llod ' . intval($rStreamID) . ' '); + foreach (ProcessManager::findProcessPIDs($rTerms) as $rPID) { + echo "Killing existing LLOD process PID: $rPID\n"; + @posix_kill($rPID, 9); } } } diff --git a/src/Cli/Commands/LoopbackCommand.php b/src/Cli/Commands/LoopbackCommand.php index e12f06c3..f115d317 100644 --- a/src/Cli/Commands/LoopbackCommand.php +++ b/src/Cli/Commands/LoopbackCommand.php @@ -5,6 +5,7 @@ namespace XcVm\Cli\Commands; use XcVm\Cli\CommandInterface; use XcVm\Core\Config\ConfigReader; use XcVm\Streaming\Fanout\FanoutClient; +use XcVm\Streaming\Fanout\IngestFeeder; /** * LoopbackCommand — loopback command @@ -127,24 +128,18 @@ class LoopbackCommand implements CommandInterface { stream_set_blocking($rFP, true); // Loopback daemon feed (ADR 0003, Phase G restream-from-origin). Loopback - // reads MPEG-TS from the parent server itself (no ffmpeg), so — like LLOD v3 - // — mirror the bytes into the xc_fanout daemon's push-fed ingest socket: the - // same stream then fans out via /live/ and in-RAM /hls on this LB, and - // live.php's isStreamFed() routes viewers to the daemon instead of the PHP - // byte path (the whole point on an LB too). The write is non-blocking + - // best-effort so a daemon stall never slows loopback's own segmenting; the - // daemon resyncs on PAT/PMT after any dropped bytes. Null/failed connect ⇒ - // legacy-only, no behaviour change. We mirror the SANITIZED buffer (below), - // not the raw read, because admin/live interleaves 0xFF padding that would - // otherwise break the daemon's packet parsing. - $rDaemonSock = FanoutClient::registerIngest($rStreamID); - $rDaemonConn = null; - if ($rDaemonSock !== null) { - $rDaemonConn = @stream_socket_client('unix://' . $rDaemonSock, $rDErrno, $rDErrstr, 2); - if ($rDaemonConn) { - stream_set_blocking($rDaemonConn, false); - } - } + // reads MPEG-TS from the parent server itself (no ffmpeg) and pushes it into + // the xc_fanout daemon's ingest socket; the daemon is what serves this LB's + // viewers (/live/ and the in-RAM /hls) — there is no other client path + // since Phase E, so this feed is the channel's delivery. IngestFeeder keeps + // short writes from tearing packets, re-registers and redials after a daemon + // restart, and carries the HLS key when encrypted HLS is on. We feed the + // SANITIZED buffer (below), not the raw read, because admin/live interleaves + // 0xFF padding that would otherwise break the daemon's packet parsing. + $rFeeder = IngestFeeder::forStream($rStreamID, !empty($rSettings['encrypt_hls']), function (string $rLine) use ($rStreamID) { + $this->writeError($rStreamID, '[Loopback] ' . $rLine); + }); + $rFeeder->connect(); $rExcessBuffer = $rPrebuffer = $rBuffer = $rPacket = ''; $rPATHeaders = array(); @@ -207,16 +202,12 @@ class LoopbackCommand implements CommandInterface { } } $rPacketNum = floor(strlen($rBuffer) / PACKET_SIZE); + if (0 == $rPacketNum) { + $rFeeder->flush(); // drain a backlog / reconnect even when nothing new arrived + } if (0 < $rPacketNum) { - // Mirror the sanitized whole-packet buffer to the daemon (best-effort, - // non-blocking). On a write error (daemon gone) stop mirroring; viewers - // then fall back to the legacy on-disk HLS this loop still writes. - if ($rDaemonConn) { - if (@fwrite($rDaemonConn, $rBuffer) === false) { - @fclose($rDaemonConn); - $rDaemonConn = null; - } - } + // Feed the sanitized whole-packet buffer to the daemon (see above). + $rFeeder->write($rBuffer); foreach (str_split($rBuffer, PACKET_SIZE) as $rPacket) { list(, $rHeader) = unpack('N', substr($rPacket, 0, 4)); $rSync = $rHeader >> 24 & 255; @@ -318,9 +309,7 @@ class LoopbackCommand implements CommandInterface { if (time() - $rLastPacket < TIMEOUT) { $this->writeError($rStreamID, '[Loopback] Connection to source closed unexpectedly.'); } - if ($rDaemonConn) { - @fclose($rDaemonConn); - } + $rFeeder->close(); FanoutClient::unregister($rStreamID); fclose($rSegmentFile); fclose($rFP); @@ -382,7 +371,11 @@ class LoopbackCommand implements CommandInterface { $rHLS .= '#EXTINF:' . round((isset($rSegmentDuration[$rSegment]) ? $rSegmentDuration[$rSegment] : 10), 0) . '.000000,' . "\n" . $rStreamID . '_' . $rSegment . '.ts' . "\n"; } } - file_put_contents(STREAMS_PATH . $rStreamID . '_.m3u8', $rHLS); + // Write-then-rename: readers never see a half-written playlist. + $rTmp = STREAMS_PATH . $rStreamID . '_.m3u8.tmp'; + if (@file_put_contents($rTmp, $rHLS) === false || !@rename($rTmp, STREAMS_PATH . $rStreamID . '_.m3u8')) { + @unlink($rTmp); + } } private function writeError($rStreamID, $rError): void { diff --git a/src/Cli/CronJobs/StreamsCronJob.php b/src/Cli/CronJobs/StreamsCronJob.php index 88a1b4d5..af9928b1 100644 --- a/src/Cli/CronJobs/StreamsCronJob.php +++ b/src/Cli/CronJobs/StreamsCronJob.php @@ -135,9 +135,9 @@ class StreamsCronJob implements CommandInterface { $rMigrate = $rStates !== null && !empty($rStates['accepting']) && StreamProcess::supervisionEnabled(); 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); + $db->query('SELECT t2.stream_display_name, t2.delay_minutes, 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 GROUP BY `stream_id`) 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 { - $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, clients.online_clients, clients_hls.online_clients_hls, 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 online_clients FROM `lines_live` WHERE `server_id` = ? AND `hls_end` = 0 GROUP BY stream_id) AS clients ON clients.stream_id = t1.stream_id 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` LEFT JOIN (SELECT stream_id, COUNT(*) as online_clients_hls FROM `lines_live` WHERE `server_id` = ? AND `container` = 'hls' AND `hls_end` = 0 GROUP BY stream_id) AS clients_hls ON clients_hls.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, SERVER_ID, SERVER_ID); + $db->query("SELECT t2.stream_display_name, t2.delay_minutes, 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, clients.online_clients, clients_hls.online_clients_hls, 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 online_clients FROM `lines_live` WHERE `server_id` = ? AND `hls_end` = 0 GROUP BY stream_id) AS clients ON clients.stream_id = t1.stream_id 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 GROUP BY `stream_id`) AS `servers_attached` ON `servers_attached`.`stream_id` = t1.`stream_id` LEFT JOIN (SELECT stream_id, COUNT(*) as online_clients_hls FROM `lines_live` WHERE `server_id` = ? AND `container` = 'hls' AND `hls_end` = 0 GROUP BY stream_id) AS clients_hls ON clients_hls.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, SERVER_ID, SERVER_ID); } if ($db->num_rows() > 0) { @@ -192,6 +192,9 @@ class StreamsCronJob implements CommandInterface { if ($rQueue == 0 && $rAdminQueue == 0 && $rStream['online_clients'] == 0 && (file_exists(STREAMS_PATH . $rStream['stream_id'] . '_.m3u8') || SettingsManager::getInt('on_demand_wait_time') < time() - intval($rStream['stream_started']) || $rStream['stream_status'] == 1)) { echo 'Stop on-demand stream...' . "\n\n"; StreamProcess::stopStream($rStream['stream_id'], true); + // Stopped: nothing below applies (it would start a thumbnail + // and a TV archive worker for the stream just stopped). + continue; } } @@ -232,8 +235,17 @@ class StreamsCronJob implements CommandInterface { // stamp so a stream whose ingest keeps failing is not // restart-looped every cron tick. Proxy streams have no // local ffmpeg, so they never reach this running branch. - if (FanoutClient::daemonStreamMissing($rStream['stream_id'])) { - $rRefeedStamp = STREAMS_PATH . $rStream['stream_id'] . '_.refeed'; + // Only an ffmpeg producer needs the restart: a broken tee slave + // stays broken. The PHP relays (LLOD, loopback) and a delayed + // stream's DelayCommand re-register and redial the daemon by + // themselves (IngestFeeder), and restarting a delayed stream + // would throw its buffer away. + $rSelfFeeding = intval($rStream['delay_minutes'] ?? 0) > 0 || ProcessManager::producerKind($rPID) === 'php'; + if (!$rSelfFeeding && FanoutClient::daemonStreamMissing($rStream['stream_id'])) { + // The stamp lives outside STREAMS_PATH: the restart below + // runs `rm -f _*` there, which deleted a `_.refeed` + // stamp — and with it the 120 s throttle it was meant to be. + $rRefeedStamp = SIGNALS_TMP_PATH . 'refeed_' . intval($rStream['stream_id']); if (!file_exists($rRefeedStamp) || time() - filemtime($rRefeedStamp) > 120) { echo 'Daemon lost stream ' . $rStream['stream_id'] . ' (restarted) — re-feeding...' . "\n\n"; touch($rRefeedStamp); diff --git a/src/Domain/Stream/StreamProcess.php b/src/Domain/Stream/StreamProcess.php index d7a3483c..a594deba 100644 --- a/src/Domain/Stream/StreamProcess.php +++ b/src/Domain/Stream/StreamProcess.php @@ -8,6 +8,7 @@ use XcVm\Core\Http\CurlClient; use XcVm\Core\Process\ProcessManager; use XcVm\Core\Util\StreamUtils; use XcVm\Streaming\Fanout\FanoutClient; +use XcVm\Streaming\Fanout\IngestFeeder; /** * StreamProcess — stream process @@ -289,12 +290,14 @@ class StreamProcess { if (!empty($rSubtitles) && !empty($rSubtitles['files']) && is_array($rSubtitles['files'])) { $rCount = count($rSubtitles['files']); for ($i = 0; $i < $rCount; $i++) { - $rSubtitleFile = escapeshellarg($rSubtitles['files'][$i]); $rInputCharset = escapeshellarg($rSubtitles['charset'][$i]); if ($rSubtitles['location'] == SERVER_ID) { - $rSubtitlesImport .= '-sub_charenc ' . $rInputCharset . ' -i ' . $rSubtitleFile . ' '; + $rSubtitlesImport .= '-sub_charenc ' . $rInputCharset . ' -i ' . escapeshellarg($rSubtitles['files'][$i]) . ' '; } else { - $rSubtitlesImport .= '-sub_charenc ' . $rInputCharset . ' -i "' . $rServers[$rSubtitles['location']]['api_url'] . '&action=getFile&filename=' . urlencode($rSubtitleFile) . '" '; + // URL-encode the raw path, then quote the whole URL for the shell. + // (Encoding the already shell-quoted path sent the quotes along, + // so the remote server looked up a filename that does not exist.) + $rSubtitlesImport .= '-sub_charenc ' . $rInputCharset . ' -i ' . escapeshellarg($rServers[$rSubtitles['location']]['api_url'] . '&action=getFile&filename=' . urlencode($rSubtitles['files'][$i])) . ' '; } } for ($i = 0; $i < $rCount; $i++) { @@ -467,6 +470,31 @@ class StreamProcess { return $rOptions . ' -f tee "' . $rHls . '|' . $rDaemon . '"'; } + /** + * The output options a live-on-demand (LLOD) start adds for low latency. + * + * The encoder tune is chosen per encoder: `-tune zerolatency` is an x264/x265 + * option, which NVENC rejects as an unknown tune value (failing the start) and + * a stream copy ignores; NVENC's equivalent is `-zerolatency 1`. + * + * @param array $rTranscodeAttributes Resolved transcode attributes. + * @return string Options for the {LLOD} placeholder. + */ + private static function llodOutputOptions(array $rTranscodeAttributes): string { + $rCodec = $rTranscodeAttributes['-vcodec'] ?? 'copy'; + if (is_array($rCodec)) { + $rCodec = $rCodec['cmd'] ?? ($rCodec['val'] ?? ''); + } + $rCodec = strtolower(trim((string) $rCodec)); + $rTune = ''; + if (in_array($rCodec, array('libx264', 'libx265'), true)) { + $rTune = '-tune zerolatency '; + } elseif (substr($rCodec, -6) === '_nvenc') { + $rTune = '-zerolatency 1 '; + } + return $rTune . '-strict experimental'; + } + /** * Wrap an FLV output target (local RTMP relay or external push URL) with the * shared `-f flv -flvflags no_duration_filesize` options. Extracted from the @@ -783,9 +811,12 @@ class StreamProcess { $rReconnect = (!$rLoopback && is_string($rSource) && preg_match('#^https?://#i', $rSource)) ? '-reconnect 1 -reconnect_streamed 1 -reconnect_delay_max 5 ' : ''; - // +discardcorrupt tolerates corrupt packets from the source; keep it scoped - // to the on-demand (LLOD) path where it originally shipped. - $rLLODInputFlags = $rReconnect . (($rLLOD && !$rLoopback) ? '-fflags +discardcorrupt ' : ''); + // LLOD input flags: +discardcorrupt tolerates corrupt packets from the + // source, +nobuffer stops the demuxer holding back what it read during + // stream analysis. Both are demuxer (input) flags — +nobuffer used to sit + // among the output options, where it does nothing. Scoped to the on-demand + // (LLOD) path where they shipped. + $rLLODInputFlags = $rReconnect . (($rLLOD && !$rLoopback) ? '-fflags +discardcorrupt+nobuffer ' : ''); // Command-template defaults: only the non-custom_ffmpeg branch below // assigns these, yet the {MAP}/{GEN_PTS}/{READ_NATIVE} substitution and @@ -856,7 +887,7 @@ class StreamProcess { $rFFMPEG = ((stripos($rStream['stream_info']['custom_ffmpeg'], 'nvenc') !== false ? $rFFMPEG_GPU : $rFFMPEG_CPU)) . ' -y -nostdin -hide_banner -loglevel ' . (($rSettings['ffmpeg_warnings'] ? 'warning' : 'error')) . ' -progress "' . $rProgressFile . '" ' . $rStream['stream_info']['custom_ffmpeg']; } - $rLLODOptions = ($rLLOD && !$rLoopback ? '-tune zerolatency -fflags nobuffer -flags low_delay -strict experimental -threads 0' : ''); + $rLLODOptions = ($rLLOD && !$rLoopback ? self::llodOutputOptions($rStream['stream_info']['transcode_attributes']) : ''); $rOutputs = array(); if ($rLoopback) { @@ -874,11 +905,11 @@ class StreamProcess { $rInitTime = min(2, intval($rSegmentSettings['seg_time'])); // When the xc_fanout daemon accepted an ingest registration (reachable), // tee the HLS output to it too (ADR 0003, A2). Never for delay, whose HLS - // goes to its own directory. startStream() registers no ingest for a - // loopback stream (the PHP relay feeds the daemon for those); a supervised - // one does, because the supervisor confirms and judges a stream by the - // bytes the daemon receives. If the daemon was unreachable - // ($data['ingestSock'] is null) the original on-disk-only HLS runs. + // goes to its own directory and whose DelayCommand feeds the daemon the + // delayed segments. A loopback stream tees too: clients are served only by + // the daemon, so an ffmpeg loopback that did not feed it could not be + // watched. If the daemon was unreachable ($data['ingestSock'] is null) the + // on-disk-only HLS runs. if (!$rDelayActive && !empty($data['ingestSock'])) { // The tee muxer needs an EXPLICIT -map — plain single outputs use // ffmpeg's automatic stream selection, but tee does not ("Output file @@ -1227,11 +1258,10 @@ class StreamProcess { list($rProbesize, $rAnalyseDuration, $rTimeout) = self::resolveProbeSettings($rStream['server_info']['on_demand'], $rInfo['probesize_ondemand'], $rLLOD, $rSettings); self::writeStreamKeyIv($rStreamID); - $rEncKey = $rEncIV = null; - if (!empty($rSettings['encrypt_hls']) && !$rLoopback) { - $rEncKey = @bin2hex((string) @file_get_contents(STREAMS_PATH . $rStreamID . '_.key')); - $rEncIV = @bin2hex((string) @file_get_contents(STREAMS_PATH . $rStreamID . '_.iv')); - } + // Loopback included: the playlist declares AES-128 whenever encrypt_hls is + // on (HLSGenerator::tokenizeDaemonPlaylist), so a daemon fed without the + // key served plain segments no player could decrypt. + [$rEncKey, $rEncIV] = !empty($rSettings['encrypt_hls']) ? IngestFeeder::streamKey($rStreamID) : array(null, null); $rIngestSock = FanoutClient::registerIngest($rStreamID, $rEncKey, $rEncIV); if ($rIngestSock === null) { return null; // no daemon to feed: the stream runs the legacy way @@ -1968,11 +1998,13 @@ class StreamProcess { if ($db->num_rows() > 0) { $rStream['server_info'] = $db->get_row(); if ($rStream['server_info']['parent_id'] != 0) { + // The key first: the relay hands it to the daemon when it + // registers its ingest, moments after it starts. + self::writeStreamKeyIv($rStreamID); shell_exec(PHP_BIN . ' ' . MAIN_HOME . 'console.php loopback ' . intval($rStreamID) . ' ' . intval($rStream['server_info']['parent_id']) . ' >/dev/null 2>/dev/null & echo $! > ' . STREAMS_PATH . intval($rStreamID) . '_.pid'); $rPID = intval(file_get_contents(STREAMS_PATH . $rStreamID . '_.pid')); $rLoopURL = (!is_null($rServers[SERVER_ID]['private_url_ip']) && !is_null($rServers[$rStream['server_info']['parent_id']]['private_url_ip']) ? $rServers[$rStream['server_info']['parent_id']]['private_url_ip'] : $rServers[$rStream['server_info']['parent_id']]['public_url_ip']); $rCurrentSource = $rLoopURL . 'admin/live?stream=' . intval($rStreamID) . '&password=' . urlencode($rSettings['live_streaming_pass']) . '&extension=ts'; - self::writeStreamKeyIv($rStreamID); $db->query('UPDATE `streams_servers` SET `delay_available_at` = ?,`to_analyze` = 0,`stream_started` = ?,`stream_info` = ?,`stream_status` = 2,`pid` = ?,`progress_info` = ?,`current_source` = ? WHERE `stream_id` = ? AND `server_id` = ?', null, time(), null, $rPID, json_encode(array()), $rCurrentSource, $rStreamID, SERVER_ID); self::updateStream($rStreamID); return array('main_pid' => $rPID, 'stream_source' => $rLoopURL . 'admin/live?stream=' . intval($rStreamID) . '&password=' . urlencode($rSettings['live_streaming_pass']) . '&extension=ts', 'delay_enabled' => false, 'parent_id' => 0, 'delay_start_at' => null, 'playlist' => STREAMS_PATH . $rStreamID . '_.m3u8', 'transcode' => false, 'offset' => 0); @@ -2001,9 +2033,11 @@ class StreamProcess { foreach ($rStreamArguments as $rStreamArgument) { $rArgumentMap[$rStreamArgument['argument_key']] = array('value' => $rStreamArgument['value'], 'argument_default_value' => $rStreamArgument['argument_default_value']); } + // The key first: the segmenter hands it to the daemon when it registers + // its ingest, moments after it starts. + self::writeStreamKeyIv($rStreamID); shell_exec(PHP_BIN . ' ' . MAIN_HOME . 'console.php llod ' . intval($rStreamID) . ' "' . base64_encode(json_encode($rSources)) . '" "' . base64_encode(json_encode($rArgumentMap)) . '" >/dev/null 2>/dev/null & echo $! > ' . STREAMS_PATH . intval($rStreamID) . '_.pid'); $rPID = intval(file_get_contents(STREAMS_PATH . $rStreamID . '_.pid')); - self::writeStreamKeyIv($rStreamID); $db->query('UPDATE `streams_servers` SET `delay_available_at` = ?,`to_analyze` = 0,`stream_started` = ?,`stream_info` = ?,`stream_status` = 2,`pid` = ?,`progress_info` = ?,`current_source` = ? WHERE `stream_id` = ? AND `server_id` = ?', null, time(), null, $rPID, json_encode(array()), $rSources[0], $rStreamID, SERVER_ID); self::updateStream($rStreamID); return array('main_pid' => $rPID, 'stream_source' => $rSources[0], 'delay_enabled' => false, 'parent_id' => 0, 'delay_start_at' => null, 'playlist' => STREAMS_PATH . $rStreamID . '_.m3u8', 'transcode' => false, 'offset' => 0); @@ -2208,6 +2242,10 @@ class StreamProcess { break; } + } else { + // LLOD skips the probe, so nothing above can pick a source: + // start on the first one rather than falling through to the last. + break; } } if (!($rStream['server_info']['on_demand'] && $rLLOD)) { @@ -2257,21 +2295,20 @@ class StreamProcess { // Assemble the live ffmpeg command (pure). buildLive() does all the // transcode-attribute resolution and {TEMPLATE} substitution internally, so // it is fed the raw stream row. - // Register a daemon ingest for standard live streams (ADR 0003, - // A2). The daemon starts listening on the returned socket, then - // buildLive tees the HLS output into it. Null when the daemon is - // unreachable → buildLive emits the on-disk-only HLS (rollback). - // Generate the stream's HLS key/iv up-front so, when encrypt_hls - // is on, we can hand them to the daemon at ingest registration - // and it encrypts the HLS segments it serves (ADR 0003, Phase B - // encrypted) — matching the panel's #EXT-X-KEY. + // Register a daemon ingest (ADR 0003, A2). The daemon starts + // listening on the returned socket, then buildLive tees the HLS + // output into it. Null when the daemon is unreachable → + // buildLive emits the on-disk-only HLS. Loopback streams tee too + // (clients are served only by the daemon); a delayed stream does + // not — its encoder output is the undelayed one, and DelayCommand + // feeds the daemon the delayed segments instead. + // The stream's HLS key/iv are generated up-front so, when + // encrypt_hls is on, the daemon gets them at registration and + // encrypts the HLS segments it serves (ADR 0003, Phase B) — + // matching the panel's #EXT-X-KEY. self::writeStreamKeyIv($rStreamID); - $rEncKey = $rEncIV = null; - if (!empty($rSettings['encrypt_hls']) && !$rLoopback && !$rDelayActive) { - $rEncKey = @bin2hex((string) @file_get_contents(STREAMS_PATH . $rStreamID . '_.key')); - $rEncIV = @bin2hex((string) @file_get_contents(STREAMS_PATH . $rStreamID . '_.iv')); - } - $rIngestSock = (!$rLoopback && !$rDelayActive) ? FanoutClient::registerIngest($rStreamID, $rEncKey, $rEncIV) : null; + [$rEncKey, $rEncIV] = (!empty($rSettings['encrypt_hls']) && !$rDelayActive) ? IngestFeeder::streamKey(intval($rStreamID)) : array(null, null); + $rIngestSock = !$rDelayActive ? FanoutClient::registerIngest(intval($rStreamID), $rEncKey, $rEncIV) : null; $rFFMPEG = self::buildLive(array( 'stream' => $rStream, 'settings' => $rSettings, 'servers' => $rServers, diff --git a/tests/Unit/StreamProcessBuildLiveTest.php b/tests/Unit/StreamProcessBuildLiveTest.php index c7e54cc4..2fbf5be7 100644 --- a/tests/Unit/StreamProcessBuildLiveTest.php +++ b/tests/Unit/StreamProcessBuildLiveTest.php @@ -139,11 +139,40 @@ final class StreamProcessBuildLiveTest extends TestCase { $this->assertStringContainsString('-map 0 -copy_unknown', $out); } - /** The legacy path registers no ingest for loopback and must stay on-disk only. */ - public function testLegacyLoopbackStaysOnDiskOnly(): void { + /** With no daemon ingest (daemon unreachable) a loopback stays on-disk only. */ + public function testLoopbackWithoutIngestStaysOnDiskOnly(): void { $this->assertStringNotContainsString('-f tee', $this->build(['loopback' => true])); } + /** A self-launched (unsupervised) loopback feeds the daemon too — clients are daemon-only. */ + public function testUnsupervisedLoopbackTeesIntoTheDaemon(): void { + $out = $this->build(['loopback' => true, 'ingestSock' => '/run/ingest/42.sock']); + $this->assertStringContainsString('-f tee', $out); + $this->assertStringContainsString('unix:/run/ingest/42.sock', $out); + } + + // ── LLOD ─────────────────────────────────────────────────── + + /** +nobuffer is a demuxer flag: it belongs before -i, not among the output options. */ + public function testLlodLatencyFlagsSitOnTheInput(): void { + $out = $this->build(['llod' => true]); + $input = substr($out, 0, (int) strpos($out, ' -i ')); + $this->assertStringContainsString('-fflags +discardcorrupt+nobuffer', $input); + $this->assertStringNotContainsString('-fflags nobuffer', $out); + } + + /** A stream copy gets no encoder tune; x264 gets zerolatency; NVENC its own option. */ + public function testLlodTuneMatchesTheEncoder(): void { + $this->assertStringNotContainsString('-tune', $this->build(['llod' => true])); + + $x264 = $this->build(['llod' => true, 'stream_info' => ['enable_transcode' => 1, 'transcode_profile_id' => 1, 'profile_options' => json_encode(['-vcodec' => 'libx264'])]]); + $this->assertStringContainsString('-tune zerolatency', $x264); + + $nvenc = $this->build(['llod' => true, 'stream_info' => ['enable_transcode' => 1, 'transcode_profile_id' => 1, 'profile_options' => json_encode(['-vcodec' => 'h264_nvenc'])]]); + $this->assertStringNotContainsString('-tune zerolatency', $nvenc, 'NVENC rejects the x264 tune value'); + $this->assertStringContainsString('-zerolatency 1', $nvenc); + } + // ── custom_ffmpeg branch ─────────────────────────────────── public function testCustomFfmpegBypassesTemplate(): void { diff --git a/tests/Unit/StreamProcessSubtitleTest.php b/tests/Unit/StreamProcessSubtitleTest.php index 1fb87a4e..36829c5b 100644 --- a/tests/Unit/StreamProcessSubtitleTest.php +++ b/tests/Unit/StreamProcessSubtitleTest.php @@ -96,4 +96,26 @@ final class StreamProcessSubtitleTest extends TestCase { $this->assertStringContainsString('action=getFile', $import); $this->assertStringContainsString('filename=', $import); } + + /** + * The remote filename is URL-encoded from the raw path — not from the shell- + * quoted one, which put encoded quotes (%27) into the name the remote server + * looked up. + */ + public function testRemoteSubtitleFilenameCarriesNoShellQuotes(): void { + if (PHP_OS_FAMILY === 'Windows') { + $this->markTestSkipped('escapeshellarg() rewrites % on Windows; the panel runs on Linux'); + } + $servers = [2 => ['api_url' => 'http://node2/api?key=abc']]; + $json = json_encode([ + 'location' => 2, + 'files' => ['/subs/My Movie.srt'], + 'charset' => ['UTF-8'], + 'names' => ['Remote'], + ]); + [$import] = $this->build($json, $servers); + + $this->assertStringContainsString('filename=' . urlencode('/subs/My Movie.srt'), $import); + $this->assertStringNotContainsString('%27', $import); + } }