diff --git a/application/Espo/Core/Webhook/Queue.php b/application/Espo/Core/Webhook/Queue.php index 6c7689887e..32dcc47903 100644 --- a/application/Espo/Core/Webhook/Queue.php +++ b/application/Espo/Core/Webhook/Queue.php @@ -29,6 +29,7 @@ namespace Espo\Core\Webhook; +use Espo\Core\Field\DateTime; use Espo\Entities\User; use Espo\Entities\Webhook; use Espo\Entities\WebhookEventQueueItem; @@ -41,7 +42,6 @@ use Espo\ORM\EntityManager; use Espo\ORM\Query\Part\Condition as Cond; use Exception; -use DateTime; use stdClass; /** @@ -74,11 +74,10 @@ class Queue { $portionSize = $this->config->get('webhookQueueEventPortionSize', self::EVENT_PORTION_SIZE); + /** @var iterable $itemList */ $itemList = $this->entityManager ->getRDBRepository(WebhookEventQueueItem::ENTITY_TYPE) - ->where([ - 'isProcessed' => false, - ]) + ->where(['isProcessed' => false]) ->order('number') ->limit(0, $portionSize) ->find(); @@ -86,10 +85,7 @@ class Queue foreach ($itemList as $item) { $this->createQueueFromEvent($item); - $item->set([ - 'isProcessed' => true, - ]); - + $item->setIsProcessed(); $this->entityManager->saveEntity($item); } } @@ -186,8 +182,7 @@ class Queue $user = null; if ($webhook->getUserId()) { - /** @var ?User $user */ - $user = $this->entityManager->getEntityById(User::ENTITY_TYPE, $webhook->getUserId()); + $user = $this->entityManager->getRDBRepositoryByClass(User::class)->getById($webhook->getUserId()); if (!$user) { foreach ($itemList as $item) { @@ -352,11 +347,11 @@ class Queue protected function succeedQueueItem(WebhookQueueItem $item): void { - $item->set([ - 'attempts' => $item->getAttempts() + 1, - 'status' => WebhookQueueItem::STATUS_SUCCESS, - 'processedAt' => DateTimeUtil::getSystemNowString(), - ]); + + $item + ->setAttempts($item->getAttempts() + 1) + ->setStatus(WebhookQueueItem::STATUS_SUCCESS) + ->setProcessedAt(DateTime::createNow()); $this->entityManager->saveEntity($item); } @@ -372,18 +367,14 @@ class Queue $maxAttemptsNumber = 0; } - $dt = new DateTime(); - $dt->modify($period); + $processAt = DateTime::createNow()->modify($period); - $item->set([ - 'attempts' => $attempts, - 'processAt' => $dt->format(DateTimeUtil::SYSTEM_DATE_TIME_FORMAT), - ]); + $item->setAttempts($attempts); + $item->setProcessAt($processAt); if ($attempts >= $maxAttemptsNumber) { - $item->set('status', WebhookQueueItem::STATUS_FAILED); - /** @noinspection PhpRedundantOptionalArgumentInspection */ - $item->set('processAt', null); + $item->setStatus(WebhookQueueItem::STATUS_FAILED); + $item->setProcessAt(null); } $this->entityManager->saveEntity($item); diff --git a/application/Espo/Entities/WebhookEventQueueItem.php b/application/Espo/Entities/WebhookEventQueueItem.php index a042b76bfd..5683b0be61 100644 --- a/application/Espo/Entities/WebhookEventQueueItem.php +++ b/application/Espo/Entities/WebhookEventQueueItem.php @@ -29,7 +29,16 @@ namespace Espo\Entities; -class WebhookEventQueueItem extends \Espo\Core\ORM\Entity +use Espo\Core\ORM\Entity; + +class WebhookEventQueueItem extends Entity { public const ENTITY_TYPE = 'WebhookEventQueueItem'; + + public function setIsProcessed(bool $isProcessed = true): self + { + $this->set('isProcessed', $isProcessed); + + return $this; + } } diff --git a/application/Espo/Entities/WebhookQueueItem.php b/application/Espo/Entities/WebhookQueueItem.php index 4bb8502d8a..33a49e1cbb 100644 --- a/application/Espo/Entities/WebhookQueueItem.php +++ b/application/Espo/Entities/WebhookQueueItem.php @@ -29,6 +29,7 @@ namespace Espo\Entities; +use Espo\Core\Field\DateTime; use Espo\Core\ORM\Entity; class WebhookQueueItem extends Entity @@ -43,4 +44,46 @@ class WebhookQueueItem extends Entity { return $this->get('attempts') ?? 0; } + + public function setStatus(string $status): self + { + $this->set('status', $status); + + return $this; + } + + public function setAttempts(?int $attempts): self + { + $this->set('attempts', $attempts); + + return $this; + } + + public function setProcessAt(?DateTime $processAt): self + { + if (!$processAt) { + $this->set('processAt', $processAt); + + return $this; + } + + $this->set('processAt', $processAt->toString()); + + return $this; + } + + public function setProcessedAt(?DateTime $processedAt): self + { + if (!$processedAt) { + $this->set('processedAt', $processedAt); + + return $this; + } + + $this->set('processedAt', $processedAt->toString()); + + return $this; + } + + }