diff --git a/vendor/magento/module-sales/Model/EmailSenderHandler.php b/vendor/magento/module-sales/Model/EmailSenderHandler.php index 59c39a0936a32..3620565956d30 100644 --- a/vendor/magento/module-sales/Model/EmailSenderHandler.php +++ b/vendor/magento/module-sales/Model/EmailSenderHandler.php @@ -14,12 +14,15 @@ use Magento\Sales\Model\ResourceModel\Collection\AbstractCollection; use Magento\Sales\Model\ResourceModel\EntityAbstract; use Magento\Store\Model\StoreManagerInterface; +use Magento\Framework\Lock\LockManagerInterface; /** * Sales emails sending * * Performs handling of cron jobs related to sending emails to customers * after creation/modification of Order, Invoice, Shipment or Creditmemo. + * + * @SuppressWarnings(PHPMD.CouplingBetweenObjects) */ class EmailSenderHandler { @@ -28,6 +31,23 @@ class EmailSenderHandler */ public const XML_PATH_ASYNC_SENDING_ATTEMPTS = 'sales_email/general/async_sending_attempts'; + /** + * Configuration path for stale in-progress async email claim timeout (minutes). + */ + public const XML_PATH_STALE_CLAIM_MINUTES = 'sales_email/general/stale_claim_minutes'; + + private const LOCK_PREFIX = 'sales_async_email_'; + + /** + * Default minutes after which a stale in-progress claim may be reclaimed by another worker. + */ + private const DEFAULT_STALE_CLAIM_MINUTES = 10; + + /** + * email_sent value used while an entity is being processed by an async email worker. + */ + public const EMAIL_SENT_PROCESSING = 2; + /** * Email sender model. * @@ -73,6 +93,11 @@ class EmailSenderHandler */ private $configValueFactory; + /** + * @var LockManagerInterface + */ + private $lockManager; + /** * @var string */ @@ -87,6 +112,7 @@ class EmailSenderHandler * @param StoreManagerInterface|null $storeManager * @param ValueFactory|null $configValueFactory * @param string|null $modifyStartFromDate + * @param LockManagerInterface|null $lockManager */ public function __construct( Sender $emailSender, @@ -97,6 +123,7 @@ public function __construct( ?StoreManagerInterface $storeManager = null, ?ValueFactory $configValueFactory = null, ?string $modifyStartFromDate = null, + ?LockManagerInterface $lockManager = null, ) { $this->emailSender = $emailSender; $this->entityResource = $entityResource; @@ -109,6 +136,7 @@ public function __construct( ->get(StoreManagerInterface::class); $this->configValueFactory = $configValueFactory ?: ObjectManager::getInstance()->get(ValueFactory::class); + $this->lockManager = $lockManager ?: ObjectManager::getInstance()->get(LockManagerInterface::class); $this->modifyStartFromDate = $modifyStartFromDate ?: $this->modifyStartFromDate; } @@ -121,13 +149,7 @@ public function sendEmails() { if ($this->globalConfig->getValue('sales_email/general/async_sending')) { $this->entityCollection->addFieldToFilter('send_email', ['eq' => 1]); - $this->entityCollection->addFieldToFilter( - 'email_sent', - [ - ['null' => true], - ['lteq' => -1] - ] - ); + $this->addPendingEmailSentFilter($this->entityCollection); $this->filterCollectionByStartFromDate($this->entityCollection); $this->entityCollection->setPageSize( $this->globalConfig->getValue('sales_email/general/sending_limit') @@ -137,6 +159,7 @@ public function sendEmails() $stores = $this->getStores(clone $this->entityCollection); $maxSendAttempts = $this->globalConfig->getValue(self::XML_PATH_ASYNC_SENDING_ATTEMPTS); + $staleClaimMinutes = $this->getStaleClaimMinutes(); /** @var \Magento\Store\Model\Store $store */ foreach ($stores as $store) { @@ -149,19 +172,31 @@ public function sendEmails() /** @var \Magento\Sales\Model\AbstractModel $item */ foreach ($entityCollection->getItems() as $item) { - $sendAttempts = $item->getEmailSent() ?? -$maxSendAttempts; - $isEmailSent = $this->emailSender->send($item, true); - - if ($isEmailSent) { - $sendAttempts = 1; - } else { - $sendAttempts++; + if (!$this->tryClaimForAsyncEmailSend((int)$item->getId(), $staleClaimMinutes)) { + continue; + } + $this->entityResource->load($item, $item->getId()); + $sendAttempts = $this->resolveSendAttempts($item->getEmailSent(), $maxSendAttempts); + $lockName = self::LOCK_PREFIX . $this->entityResource->getMainTable() . '_' . $item->getId(); + if (!$this->lockManager->lock($lockName, 0)) { + continue; } + try { + $isEmailSent = $this->emailSender->send($item, true); - $this->entityResource->saveAttribute( - $item->setEmailSent($sendAttempts), - 'email_sent' - ); + if ($isEmailSent) { + $sendAttempts = 1; + } else { + $sendAttempts++; + } + + $this->entityResource->saveAttribute( + $item->setEmailSent($sendAttempts), + 'email_sent' + ); + } finally { + $this->lockManager->unlock($lockName); + } } } } @@ -190,6 +225,90 @@ private function getStores( return $stores; } + /** + * Restrict collection to entities that still need async email processing. + * + * @param AbstractCollection $collection + * @return void + */ + private function addPendingEmailSentFilter(AbstractCollection $collection): void + { + $staleClaimMinutes = $this->getStaleClaimMinutes(); + $staleThreshold = date('Y-m-d H:i:s', strtotime(sprintf('-%d minutes', $staleClaimMinutes))); + $collection->getSelect()->where( + '(main_table.email_sent IS NULL OR main_table.email_sent = 0 OR main_table.email_sent <= ? ' + . 'OR (main_table.email_sent = ? AND main_table.updated_at < ?))', + -1, + self::EMAIL_SENT_PROCESSING, + $staleThreshold + ); + } + + /** + * Resolve the retry counter for an entity pending async email delivery. + * + * @param mixed $emailSent + * @param int $maxSendAttempts + * @return int + */ + private function resolveSendAttempts(mixed $emailSent, int $maxSendAttempts): int + { + if ($emailSent === null || (int)$emailSent === self::EMAIL_SENT_PROCESSING) { + return -$maxSendAttempts; + } + + return (int)$emailSent; + } + + /** + * Atomically claim an entity for async email sending. + * + * Sets email_sent to the in-progress status only when the row is still pending or has a stale claim. + * + * @param int $entityId + * @param int $staleClaimMinutes + * @return bool + */ + private function tryClaimForAsyncEmailSend(int $entityId, int $staleClaimMinutes): bool + { + $connection = $this->entityResource->getConnection(); + $mainTable = $this->entityResource->getMainTable(); + $staleThreshold = date('Y-m-d H:i:s', strtotime(sprintf('-%d minutes', $staleClaimMinutes))); + $pendingCondition = implode( + ' OR ', + [ + 'email_sent IS NULL', + 'email_sent = 0', + $connection->quoteInto('email_sent <= ?', -1), + '(' . $connection->quoteInto('email_sent = ?', self::EMAIL_SENT_PROCESSING) + . ' AND ' . $connection->quoteInto('updated_at < ?', $staleThreshold) . ')', + ] + ); + $where = $connection->quoteInto($this->entityResource->getIdFieldName() . ' = ?', $entityId) + . ' AND send_email = 1 AND (' . $pendingCondition . ')'; + $data = ['email_sent' => self::EMAIL_SENT_PROCESSING]; + if ($connection->tableColumnExists($mainTable, 'updated_at')) { + $data['updated_at'] = gmdate('Y-m-d H:i:s'); + } + return $connection->update($mainTable, $data, $where) === 1; + } + + /** + * Get configured stale claim timeout in minutes. + * + * @return int + */ + private function getStaleClaimMinutes(): int + { + $staleClaimMinutes = $this->globalConfig->getValue(self::XML_PATH_STALE_CLAIM_MINUTES); + + if ($staleClaimMinutes === null || $staleClaimMinutes === '') { + return self::DEFAULT_STALE_CLAIM_MINUTES; + } + + return max(1, (int)$staleClaimMinutes); + } + /** * Filter collection by start from date * diff --git a/vendor/magento/module-sales/etc/config.xml b/vendor/magento/module-sales/etc/config.xml index aa74c06b1ce96..6541ee94e2dad 100644 --- a/vendor/magento/module-sales/etc/config.xml +++ b/vendor/magento/module-sales/etc/config.xml @@ -38,6 +38,7 @@ 0 3 50 + 10 1