diff --git a/vendor/magento/module-message-queue/Console/StartConsumerCommand.php b/vendor/magento/module-message-queue/Console/StartConsumerCommand.php index fc2207dcd7c..8ea6290a2a4 100644 --- a/vendor/magento/module-message-queue/Console/StartConsumerCommand.php +++ b/vendor/magento/module-message-queue/Console/StartConsumerCommand.php @@ -79,20 +79,15 @@ class StartConsumerCommand extends Command $singleThread = $input->getOption(self::OPTION_SINGLE_THREAD); - if ($singleThread && $this->lockManager->isLocked(md5($consumerName))) { //phpcs:ignore + if ($singleThread && !$this->lockManager->lock(md5($consumerName),0)) { //phpcs:ignore $output->writeln('Consumer with the same name is running'); return \Magento\Framework\Console\Cli::RETURN_FAILURE; } - if ($singleThread) { - $this->lockManager->lock(md5($consumerName)); //phpcs:ignore - } - $this->appState->setAreaCode($areaCode ?? 'global'); $consumer = $this->consumerFactory->get($consumerName, $batchSize); $consumer->process($numberOfMessages); - if ($singleThread) { $this->lockManager->unlock(md5($consumerName)); //phpcs:ignore } @@ -163,7 +158,7 @@ To specify the number of messages per batch for the batch consumer: To specify the preferred area: %command.full_name% someConsumer --area-code='adminhtml' - + To do not run multiple copies of one consumer simultaneously: %command.full_name% someConsumer --single-thread' diff --git a/vendor/magento/framework/Cache/LockGuardedCacheLoader.php b/vendor/magento/framework/Cache/LockGuardedCacheLoader.php index bca23e0dcf3..4937112a444 100644 --- a/vendor/magento/framework/Cache/LockGuardedCacheLoader.php +++ b/vendor/magento/framework/Cache/LockGuardedCacheLoader.php @@ -131,7 +131,7 @@ class LockGuardedCacheLoader return $dataCollector(); } - if ($this->locker->lock($lockName, $this->lockTimeout / 1000)) { + if ($this->locker->lock($lockName, 0)) { try { $data = $dataCollector(); $dataSaver($data); diff --git a/vendor/magento/framework/Lock/Backend/Cache.php b/vendor/magento/framework/Lock/Backend/Cache.php index 612d8541281..ae777a6701c 100644 --- a/vendor/magento/framework/Lock/Backend/Cache.php +++ b/vendor/magento/framework/Lock/Backend/Cache.php @@ -32,6 +32,27 @@ class Cache implements \Magento\Framework\Lock\LockManagerInterface private $lockSign; /** + * How many microseconds to wait before re-try to acquire a lock + * + * @var int + */ + private $sleepCycle = 100000; + + /** + * Lifetime of lock data in seconds. + * + * @var int + */ + private $defaultLifetime = 7200; + + /** + * Array for keeping all lock attempt to release them on destruct. + * + * @var string[] + */ + private $lockArrayState = []; + + /** * @param FrontendInterface $cache */ public function __construct(FrontendInterface $cache) @@ -49,18 +70,21 @@ class Cache implements \Magento\Framework\Lock\LockManagerInterface $this->lockSign = $this->generateLockSign(); } - $data = $this->cache->load($this->getIdentifier($name)); - - if (false !== $data) { - return false; + $skipDeadline = $timeout < 0; + $deadline = microtime(true) + $timeout; + while ($this->cache->load($this->getIdentifier($name))) { + if (!$skipDeadline && $deadline <= microtime(true)) { + return false; + } + usleep($this->sleepCycle); } - $timeout = $timeout <= 0 ? null : $timeout; - $this->cache->save($this->lockSign, $this->getIdentifier($name), [], $timeout); + $this->cache->save($this->lockSign, $this->getIdentifier($name), [], $this->defaultLifetime); $data = $this->cache->load($this->getIdentifier($name)); if ($data === $this->lockSign) { + $this->lockArrayState[$name] = 1; return true; } @@ -85,6 +109,7 @@ class Cache implements \Magento\Framework\Lock\LockManagerInterface $removeResult = false; if ($data === $this->lockSign) { $removeResult = (bool)$this->cache->remove($this->getIdentifier($name)); + unset($this->lockArrayState[$name]); } return $removeResult; @@ -131,4 +156,26 @@ class Cache implements \Magento\Framework\Lock\LockManagerInterface return $sign; } + + /** + * Destruct method should release all locks that left. + * + * @return void + */ + public function __destruct() + { + $this->releaseLocks(); + } + + /** + * Release all locks that were not removed with unlock method. + * + * @return void + */ + private function releaseLocks() + { + foreach ($this->lockArrayState as $name => $value) { + $this->unlock($name); + } + } }