ref
This commit is contained in:
@@ -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<WebhookEventQueueItem> $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);
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user