diff --git a/src/Cli/Commands/CacheHandlerCommand.php b/src/Cli/Commands/CacheHandlerCommand.php index 1a732314..c1e51478 100644 --- a/src/Cli/Commands/CacheHandlerCommand.php +++ b/src/Cli/Commands/CacheHandlerCommand.php @@ -8,6 +8,7 @@ use XcVm\Core\Config\SettingsManager; use XcVm\Core\Config\SettingsRepository; use XcVm\Domain\Line\LineService; use XcVm\Domain\Server\ServerRepository; +use XcVm\Infrastructure\Signal\SignalQueue; /** * CacheHandlerCommand — cache handler command @@ -75,8 +76,7 @@ class CacheHandlerCommand implements CommandInterface { try { $rUpdatedLines = array(); - foreach (glob(SIGNALS_TMP_PATH . 'cache_*') as $rFileMD5) { - list($rKey, $rData) = json_decode(file_get_contents($rFileMD5), true); + foreach (SignalQueue::pending() as list($rFileMD5, $rKey, $rData)) { list($rHeader) = explode('/', $rKey); switch ($rHeader) { case 'restream_block_user': diff --git a/src/Core/Auth/BruteforceGuard.php b/src/Core/Auth/BruteforceGuard.php index f533f941..3ea56538 100644 --- a/src/Core/Auth/BruteforceGuard.php +++ b/src/Core/Auth/BruteforceGuard.php @@ -102,7 +102,7 @@ class BruteforceGuard { private static function blockIP(string $ip, string $reason, bool $useCachedMode = false): void { if ($useCachedMode && !empty($GLOBALS['rCached'])) { $signalKey = (stripos($reason, 'BRUTEFORCE') !== false ? 'bruteforce_attack' : 'flood_attack'); - \XcVm\Infrastructure\Redis\RedisManager::setSignal($signalKey . '/' . $ip, 1); + \XcVm\Infrastructure\Signal\SignalQueue::push($signalKey . '/' . $ip, 1); } else { $db = self::getDB(); if ($db) { diff --git a/src/Domain/User/UserRepository.php b/src/Domain/User/UserRepository.php index 4f20e18b..7b48bcb4 100644 --- a/src/Domain/User/UserRepository.php +++ b/src/Domain/User/UserRepository.php @@ -7,6 +7,7 @@ use XcVm\Core\GeoIP\GeoIPService; use XcVm\Core\Util\GeoIP; use XcVm\Domain\Bouquet\BouquetService; use XcVm\Domain\Security\BlocklistService; +use XcVm\Infrastructure\Signal\SignalQueue; /** * UserRepository — user repository @@ -252,19 +253,16 @@ class UserRepository { $db = self::db(); $rUserIDs = $rPermissions['direct_reports']; - if (!$rIncludeSelf) { - } else { + if ($rIncludeSelf) { $rUserIDs[] = $rUserInfo['id']; } $rReturn = array(); - if (0 >= count($rUserIDs)) { - } else { + if (0 < count($rUserIDs)) { $db->query('SELECT * FROM `users` WHERE `owner_id` IN (' . implode(',', array_map('intval', $rUserIDs)) . ') ORDER BY `username` ASC;'); - if (0 >= $db->num_rows()) { - } else { + if (0 < $db->num_rows()) { foreach ($db->get_rows() as $rRow) { $rReturn[intval($rRow['id'])] = $rRow; } @@ -322,8 +320,7 @@ class UserRepository { $db = self::db(); $db->query('SELECT * FROM `lines` WHERE `id` = ?;', $rID); - if ($db->num_rows() != 1) { - } else { + if ($db->num_rows() == 1) { return $db->get_row(); } return null; @@ -339,8 +336,7 @@ class UserRepository { $db = self::db(); $db->query('SELECT * FROM `users` WHERE `id` = ?;', $rID); - if ($db->num_rows() != 1) { - } else { + if ($db->num_rows() == 1) { return $db->get_row(); } return null; @@ -358,11 +354,9 @@ class UserRepository { $rReturn = array(); $db->query('SELECT * FROM `users` ORDER BY `username` ASC;'); - if (0 >= $db->num_rows()) { - } else { + if (0 < $db->num_rows()) { foreach ($db->get_rows() as $rRow) { - if (!(!$rOwner || $rRow['owner_id'] == $rOwner || $rRow['id'] == $rOwner && $rIncludeSelf)) { - } else { + if (!$rOwner || $rRow['owner_id'] == $rOwner || $rRow['id'] == $rOwner && $rIncludeSelf) { $rReturn[intval($rRow['id'])] = $rRow; } } @@ -409,7 +403,7 @@ class UserRepository { $rUserInfo['forced_country'] = GeoIPService::getIPInfo($rIP)['registered_country']['iso_code']; if ($rCached) { - file_put_contents(SIGNALS_TMP_PATH . 'cache_' . md5('forced_country/' . $rUserInfo['id']), json_encode(array('forced_country/' . $rUserInfo['id'], $rUserInfo['forced_country']))); + SignalQueue::push('forced_country/' . $rUserInfo['id'], $rUserInfo['forced_country']); } else { $db->query('UPDATE `lines` SET `forced_country` = ? WHERE `id` = ?', $rUserInfo['forced_country'], $rUserInfo['id']); } @@ -418,7 +412,6 @@ class UserRepository { $rUserInfo = self::decodeUserFields($rUserInfo); $rUserInfo['output_formats'] = self::resolveOutputFormats($db, $rCached, $rUserInfo['allowed_outputs']); - $rUserInfo['con_isp_name'] = null; $rUserInfo['isp_violate'] = 0; $rUserInfo['isp_is_server'] = 0; @@ -442,8 +435,7 @@ class UserRepository { if (self::ispChanged($rUserInfo['con_isp_name'], $rUserInfo['isp_violate'], $rUserInfo['isp_desc'])) { if ($rCached) { - $rSignalKey = 'isp/' . $rUserInfo['id']; - file_put_contents(SIGNALS_TMP_PATH . 'cache_' . md5($rSignalKey), json_encode(array($rSignalKey, json_encode(array($rUserInfo['con_isp_name'], $rUserInfo['isp_asn']))))); + SignalQueue::push('isp/' . $rUserInfo['id'], json_encode(array($rUserInfo['con_isp_name'], $rUserInfo['isp_asn']))); } else { $db->query('UPDATE `lines` SET `isp_desc` = ?, `as_number` = ? WHERE `id` = ?', $rUserInfo['con_isp_name'], $rUserInfo['isp_asn'], $rUserInfo['id']); } @@ -486,7 +478,7 @@ class UserRepository { if ($rSettings['county_override_1st'] == 1 && empty($rUserInfo['forced_country']) && !empty($rIP) && $rUserInfo['max_connections'] == 1) { $rUserInfo['forced_country'] = GeoIP::getCountry($rIP)['registered_country']['iso_code']; if ($rCached) { - \XcVm\Infrastructure\Redis\RedisManager::setSignal('forced_country/' . $rUserInfo['id'], $rUserInfo['forced_country']); + SignalQueue::push('forced_country/' . $rUserInfo['id'], $rUserInfo['forced_country']); } else { $db->query('UPDATE `lines` SET `forced_country` = ? WHERE `id` = ?', $rUserInfo['forced_country'], $rUserInfo['id']); } @@ -512,7 +504,7 @@ class UserRepository { } if (self::ispChanged($rUserInfo['con_isp_name'], $rUserInfo['isp_violate'], $rUserInfo['isp_desc'])) { if ($rCached) { - \XcVm\Infrastructure\Redis\RedisManager::setSignal('isp/' . $rUserInfo['id'], json_encode(array($rUserInfo['con_isp_name'], $rUserInfo['isp_asn']))); + SignalQueue::push('isp/' . $rUserInfo['id'], json_encode(array($rUserInfo['con_isp_name'], $rUserInfo['isp_asn']))); } else { $db->query('UPDATE `lines` SET `isp_desc` = ?, `as_number` = ? WHERE `id` = ?', $rUserInfo['con_isp_name'], $rUserInfo['isp_asn'], $rUserInfo['id']); } diff --git a/src/Infrastructure/Redis/RedisManager.php b/src/Infrastructure/Redis/RedisManager.php index 82087388..5dfc583f 100644 --- a/src/Infrastructure/Redis/RedisManager.php +++ b/src/Infrastructure/Redis/RedisManager.php @@ -97,16 +97,15 @@ class RedisManager { /** - * Write a signal to the filesystem cache. - * - * Stores a JSON-encoded [key, data] pair in SIGNALS_TMP_PATH. + * @deprecated Signals now live in {@see \XcVm\Infrastructure\Signal\SignalQueue}. + * Kept as a thin back-compat alias; call SignalQueue::push() directly. * * @param string $rKey Signal key. * @param mixed $rData Signal payload. * @return void */ public static function setSignal(string $rKey, $rData): void { - file_put_contents(SIGNALS_TMP_PATH . 'cache_' . md5($rKey), json_encode(array($rKey, $rData))); + \XcVm\Infrastructure\Signal\SignalQueue::push($rKey, $rData); } /** diff --git a/src/Infrastructure/Signal/SignalQueue.php b/src/Infrastructure/Signal/SignalQueue.php new file mode 100644 index 00000000..a20cd9cd --- /dev/null +++ b/src/Infrastructure/Signal/SignalQueue.php @@ -0,0 +1,67 @@ +` and holding `json_encode([key, data])`. The key's first + * '/'-separated segment selects the consumer action, e.g. `isp/` or + * `forced_country/`. + * + * @package XC_VM_Infrastructure_Signal + * @author Divarion_D + * @copyright 2025-2026 Vateron Media + * @link https://github.com/Vateron-Media/XC_VM + * @license AGPL-3.0 https://www.gnu.org/licenses/agpl-3.0.html + */ +final class SignalQueue { + + /** Filename prefix shared by every queued signal. */ + public const PREFIX = 'cache_'; + + /** + * Filesystem path of the signal file backing $rKey. + * + * @param string $rKey Signal key. + * @return string Absolute path under SIGNALS_TMP_PATH. + */ + public static function pathFor(string $rKey): string { + return SIGNALS_TMP_PATH . self::PREFIX . md5($rKey); + } + + /** + * Queue a signal for the background worker. Re-queuing the same key + * overwrites the pending record (idempotent per key). + * + * @param string $rKey Signal key; its first '/'-segment routes the action. + * @param mixed $rData JSON-encodable payload. + * @return void + */ + public static function push(string $rKey, $rData): void { + file_put_contents(self::pathFor($rKey), json_encode(array($rKey, $rData))); + } + + /** + * Every queued signal as a `[file, key, data]` tuple (filesystem order). + * Malformed records are skipped. The caller deletes each file once handled. + * + * @return array + */ + public static function pending(): array { + $rOut = array(); + foreach (glob(SIGNALS_TMP_PATH . self::PREFIX . '*') ?: array() as $rFile) { + $rDecoded = json_decode((string) @file_get_contents($rFile), true); + if (is_array($rDecoded) && array_key_exists(0, $rDecoded) && array_key_exists(1, $rDecoded)) { + $rOut[] = array($rFile, $rDecoded[0], $rDecoded[1]); + } + } + return $rOut; + } +} diff --git a/src/Public/stream/auth.php b/src/Public/stream/auth.php index 5c382ab0..f8ed8300 100644 --- a/src/Public/stream/auth.php +++ b/src/Public/stream/auth.php @@ -9,7 +9,7 @@ use XcVm\Core\Util\Encryption; use XcVm\Domain\Security\BlocklistService; use XcVm\Domain\Stream\ConnectionTracker; use XcVm\Domain\User\UserRepository; -use XcVm\Infrastructure\Redis\RedisManager; +use XcVm\Infrastructure\Signal\SignalQueue; use XcVm\Streaming\Balancer\ProxySelector; use XcVm\Streaming\Delivery\OffAirHandler; use XcVm\Streaming\Delivery\StreamRedirector; @@ -391,7 +391,7 @@ if ($rExtension) { if ($rRestreamDetect) { if ($rSettings['detect_restream_block_user']) { if ($rCached) { - RedisManager::setSignal('restream_block_user/' . $rUserInfo['id'] . '/' . $rStreamID . '/' . $rIP, 1); + SignalQueue::push('restream_block_user/' . $rUserInfo['id'] . '/' . $rStreamID . '/' . $rIP, 1); } else { $db->query('UPDATE `lines` SET `admin_enabled` = 0 WHERE `id` = ?;', $rUserInfo['id']); } @@ -414,7 +414,7 @@ if ($rExtension) { if (($rType == 'live' && $rSettings['show_expiring_video'] && !$rUserInfo['is_trial'] && !is_null($rUserInfo['exp_date']) && $rUserInfo['exp_date'] - 86400 * 7 <= time() && (86400 <= time() - $rUserInfo['last_expiration_video'] || !$rUserInfo['last_expiration_video']))) { if ($rCached) { - RedisManager::setSignal('expiring/' . $rUserInfo['id'], time()); + SignalQueue::push('expiring/' . $rUserInfo['id'], time()); } else { $db->query('UPDATE `lines` SET `last_expiration_video` = ? WHERE `id` = ?;', time(), $rUserInfo['id']); } diff --git a/tests/Unit/SignalQueueTest.php b/tests/Unit/SignalQueueTest.php new file mode 100644 index 00000000..680deff4 --- /dev/null +++ b/tests/Unit/SignalQueueTest.php @@ -0,0 +1,67 @@ + holding + * [key, data]) that the consumer's switch depends on. + */ +final class SignalQueueTest extends TestCase { + + public static function setUpBeforeClass(): void { + if (!defined('SIGNALS_TMP_PATH')) { + define('SIGNALS_TMP_PATH', sys_get_temp_dir() . '/xcvm-sig-test/'); + } + if (!is_dir(SIGNALS_TMP_PATH)) { + mkdir(SIGNALS_TMP_PATH, 0777, true); + } + } + + protected function setUp(): void { + foreach (glob(SIGNALS_TMP_PATH . 'cache_*') ?: array() as $rFile) { + unlink($rFile); + } + } + + public function testPushWritesCanonicalFileFormat(): void { + SignalQueue::push('isp/42', json_encode(array('Comcast', 7922))); + $rPath = SIGNALS_TMP_PATH . 'cache_' . md5('isp/42'); + $this->assertFileExists($rPath); + $this->assertSame(array('isp/42', json_encode(array('Comcast', 7922))), json_decode(file_get_contents($rPath), true)); + } + + public function testPathForMatchesPushLocation(): void { + $this->assertSame(SIGNALS_TMP_PATH . 'cache_' . md5('forced_country/5'), SignalQueue::pathFor('forced_country/5')); + } + + public function testPushIsIdempotentPerKey(): void { + SignalQueue::push('expiring/9', 100); + SignalQueue::push('expiring/9', 200); // same key overwrites the pending record + $rPending = SignalQueue::pending(); + $this->assertCount(1, $rPending); + $this->assertSame(200, $rPending[0][2]); + } + + public function testPendingReturnsFileKeyDataTuples(): void { + SignalQueue::push('forced_country/1', 'DE'); + SignalQueue::push('isp/2', json_encode(array('ISP', 1))); + $rByKey = array(); + foreach (SignalQueue::pending() as list($rFile, $rKey, $rData)) { + $this->assertFileExists($rFile); + $rByKey[$rKey] = $rData; + } + $this->assertSame('DE', $rByKey['forced_country/1']); + $this->assertSame(json_encode(array('ISP', 1)), $rByKey['isp/2']); + } + + public function testPendingSkipsMalformedRecords(): void { + file_put_contents(SIGNALS_TMP_PATH . 'cache_bad', 'not-json'); + SignalQueue::push('isp/3', 1); + $rPending = SignalQueue::pending(); + $this->assertCount(1, $rPending); + $this->assertSame('isp/3', $rPending[0][1]); + } +}