From 6c03a00f573906a37082cb91ed6f8d698e2b2cd7 Mon Sep 17 00:00:00 2001 From: Yuri Kuznetsov Date: Sun, 18 Apr 2021 19:43:54 +0300 Subject: [PATCH] cs fix --- application/Espo/Core/WebSocket/Pusher.php | 10 +- application/Espo/Core/WebSocket/Sender.php | 2 +- .../Espo/Core/WebSocket/ServerStarter.php | 4 +- .../Espo/Core/WebSocket/Submission.php | 2 +- .../Espo/Core/WebSocket/Subscriber.php | 2 +- .../Espo/Core/WebSocket/ZeroMQSender.php | 2 +- .../Espo/Core/WebSocket/ZeroMQSubscriber.php | 2 +- application/Espo/Core/Webhook/Manager.php | 35 ++-- application/Espo/Core/Webhook/Queue.php | 153 ++++++++++++------ application/Espo/Core/Webhook/Sender.php | 18 ++- 10 files changed, 155 insertions(+), 75 deletions(-) diff --git a/application/Espo/Core/WebSocket/Pusher.php b/application/Espo/Core/WebSocket/Pusher.php index 067545c9be..2bacae9549 100644 --- a/application/Espo/Core/WebSocket/Pusher.php +++ b/application/Espo/Core/WebSocket/Pusher.php @@ -155,7 +155,7 @@ class Pusher implements WampServerInterface } } - protected function getCategoryData(string $topicId) : array + protected function getCategoryData(string $topicId): array { $arr = explode('.', $topicId); @@ -174,7 +174,7 @@ class Pusher implements WampServerInterface return $data; } - protected function getParamsFromTopicId(string $topicId) : array + protected function getParamsFromTopicId(string $topicId): array { $arr = explode('.', $topicId); @@ -196,7 +196,7 @@ class Pusher implements WampServerInterface return $params; } - protected function getAccessCheckCommandForTopic(ConnectionInterface $connection, $topic) : ?string + protected function getAccessCheckCommandForTopic(ConnectionInterface $connection, $topic): ?string { $topicId = $topic->getId(); @@ -381,7 +381,7 @@ class Pusher implements WampServerInterface { } - public function onMessageReceive(string $message) : void + public function onMessageReceive(string $message): void { $data = json_decode($message); @@ -443,7 +443,7 @@ class Pusher implements WampServerInterface } } - protected function log(string $msg) : void + protected function log(string $msg): void { echo "[" . date('Y-m-d H:i:s') . "] " . $msg . "\n"; } diff --git a/application/Espo/Core/WebSocket/Sender.php b/application/Espo/Core/WebSocket/Sender.php index 56efebec7a..c7851d00a2 100644 --- a/application/Espo/Core/WebSocket/Sender.php +++ b/application/Espo/Core/WebSocket/Sender.php @@ -37,5 +37,5 @@ interface Sender /** * Send a message to a web-socket server. */ - public function send(string $message) : void; + public function send(string $message): void; } diff --git a/application/Espo/Core/WebSocket/ServerStarter.php b/application/Espo/Core/WebSocket/ServerStarter.php index 6c40c055e8..7c6ffff517 100644 --- a/application/Espo/Core/WebSocket/ServerStarter.php +++ b/application/Espo/Core/WebSocket/ServerStarter.php @@ -90,7 +90,7 @@ class ServerStarter /** * Start a web-socket server. */ - public function start() : void + public function start(): void { $loop = EventLoopFactory::create(); @@ -118,7 +118,7 @@ class ServerStarter $loop->run(); } - protected function getSslParams() : array + protected function getSslParams(): array { $sslParams = [ 'local_cert' => $this->config->get('webSocketSslCertificateFile'), diff --git a/application/Espo/Core/WebSocket/Submission.php b/application/Espo/Core/WebSocket/Submission.php index 0894d36c2e..027b3e6a72 100644 --- a/application/Espo/Core/WebSocket/Submission.php +++ b/application/Espo/Core/WebSocket/Submission.php @@ -48,7 +48,7 @@ class Submission /** * Submit to a web-socket server. */ - public function submit(string $topic, ?string $userId = null, ?StdClass $data = null) : void + public function submit(string $topic, ?string $userId = null, ?StdClass $data = null): void { if (!$data) { $data = (object) []; diff --git a/application/Espo/Core/WebSocket/Subscriber.php b/application/Espo/Core/WebSocket/Subscriber.php index 044b347c1e..c748a5dc2f 100644 --- a/application/Espo/Core/WebSocket/Subscriber.php +++ b/application/Espo/Core/WebSocket/Subscriber.php @@ -39,5 +39,5 @@ interface Subscriber /** * Subscribe messages to Pusher::onMessageReceive method. */ - public function subscribe(Pusher $pusher, LoopInterface $loop) : void; + public function subscribe(Pusher $pusher, LoopInterface $loop): void; } diff --git a/application/Espo/Core/WebSocket/ZeroMQSender.php b/application/Espo/Core/WebSocket/ZeroMQSender.php index 1d04fc937c..67a87eb84d 100644 --- a/application/Espo/Core/WebSocket/ZeroMQSender.php +++ b/application/Espo/Core/WebSocket/ZeroMQSender.php @@ -43,7 +43,7 @@ class ZeroMQSender implements Sender $this->config = $config; } - public function send(string $message) : void + public function send(string $message): void { $dsn = $this->config->get('webSocketSubmissionDsn', 'tcp://localhost:5555'); diff --git a/application/Espo/Core/WebSocket/ZeroMQSubscriber.php b/application/Espo/Core/WebSocket/ZeroMQSubscriber.php index 3c5ea0551a..6e93ec75dc 100644 --- a/application/Espo/Core/WebSocket/ZeroMQSubscriber.php +++ b/application/Espo/Core/WebSocket/ZeroMQSubscriber.php @@ -38,7 +38,7 @@ use ZMQ; class ZeroMQSubscriber implements Subscriber { - public function subscribe(Pusher $pusher, LoopInterface $loop) : void + public function subscribe(Pusher $pusher, LoopInterface $loop): void { $context = new ZMQContext($loop); diff --git a/application/Espo/Core/Webhook/Manager.php b/application/Espo/Core/Webhook/Manager.php index 97f62c3dd5..dd3379ebea 100644 --- a/application/Espo/Core/Webhook/Manager.php +++ b/application/Espo/Core/Webhook/Manager.php @@ -48,10 +48,13 @@ class Manager private $data = null; - protected $config; - protected $dataCache; - protected $entityManager; - protected $fieldUtil; + private $config; + + private $dataCache; + + private $entityManager; + + private $fieldUtil; public function __construct( Config $config, @@ -67,7 +70,7 @@ class Manager $this->loadData(); } - private function loadData() + private function loadData(): void { if ($this->config->get('useCache')) { if ($this->dataCache->has($this->cacheKey)) { @@ -84,12 +87,12 @@ class Manager } } - private function storeDataToCache() + private function storeDataToCache(): void { $this->dataCache->store($this->cacheKey, $this->data); } - private function buildData() + private function buildData(): void { $data = []; @@ -112,7 +115,7 @@ class Manager /** * Add an event. To cache the information that at least one webhook for this event exists. */ - public function addEvent(string $event) + public function addEvent(string $event): void { $this->data[$event] = true; @@ -124,7 +127,7 @@ class Manager /** * Remove an event. If no webhooks with this event left, then it will be removed from the cache. */ - public function removeEvent(string $event) + public function removeEvent(string $event): void { $notExists = !$this->entityManager ->getRepository('Webhook') @@ -144,12 +147,12 @@ class Manager } } - protected function eventExists(string $event) : bool + protected function eventExists(string $event): bool { return isset($this->data[$event]); } - protected function logDebugEvent(string $event, Entity $entity) + protected function logDebugEvent(string $event, Entity $entity): void { $GLOBALS['log']->debug("Webhook: {$event} on record {$entity->id}."); } @@ -157,7 +160,7 @@ class Manager /** * Process 'create' event. */ - public function processCreate(Entity $entity) + public function processCreate(Entity $entity): void { $event = $entity->getEntityType() . '.create'; @@ -178,7 +181,7 @@ class Manager /** * Process 'delete' event. */ - public function processDelete(Entity $entity) + public function processDelete(Entity $entity): void { $event = $entity->getEntityType() . '.delete'; @@ -201,7 +204,7 @@ class Manager /** * Process 'update' event. */ - public function processUpdate(Entity $entity) + public function processUpdate(Entity $entity): void { $event = $entity->getEntityType() . '.update'; @@ -230,6 +233,7 @@ class Manager 'targetId' => $entity->id, 'data' => $data, ]); + $this->logDebugEvent($event, $entity); } @@ -241,6 +245,7 @@ class Manager } $attributeList = $this->fieldUtil->getActualAttributeList($entity->getEntityType(), $field); + $isChanged = false; foreach ($attributeList as $attribute) { @@ -257,7 +262,9 @@ class Manager if ($isChanged) { $itemData = (object) []; + $itemData->id = $entity->id; + $attributeList = $this->fieldUtil->getAttributeList($entity->getEntityType(), $field); foreach ($attributeList as $attribute) { diff --git a/application/Espo/Core/Webhook/Queue.php b/application/Espo/Core/Webhook/Queue.php index cc4a877ec8..ed875ecc21 100644 --- a/application/Espo/Core/Webhook/Queue.php +++ b/application/Espo/Core/Webhook/Queue.php @@ -38,12 +38,16 @@ use Espo\Entities\{ use Espo\Core\{ AclManager, Utils\Config, - Utils\DateTime, + Utils\DateTime as DateTimeUtil, ORM\EntityManager, }; +use Exception; +use DateTime; + /** - * Groups ocurred events into portions and sends them. A portion contains multiple events of the same webhook. + * Groups ocurred events into portions and sends them. A portion contains + * multiple events of the same webhook. */ class Queue { @@ -57,20 +61,27 @@ class Queue const FAIL_ATTEMPT_PERIOD = '10 minutes'; - protected $sender; - protected $config; - protected $entityManager; - protected $aclManager; + private $sender; - public function __construct(Sender $sender, Config $config, EntityManager $entityManager, AclManager $aclManager) - { + private $config; + + private $entityManager; + + private $aclManager; + + public function __construct( + Sender $sender, + Config $config, + EntityManager $entityManager, + AclManager $aclManager + ) { $this->sender = $sender; $this->config = $config; $this->entityManager = $entityManager; $this->aclManager = $aclManager; } - public function process() + public function process(): void { $this->processEvents(); $this->processSending(); @@ -80,25 +91,36 @@ class Queue { $portionSize = $this->config->get('webhookQueueEventPortionSize', self::EVENT_PORTION_SIZE); - $itemList = $this->entityManager->getRepository('WebhookEventQueueItem')->where([ - 'isProcessed' => false, - ])->order('number')->limit(0, $portionSize)->find(); + $itemList = $this->entityManager + ->getRepository('WebhookEventQueueItem') + ->where([ + 'isProcessed' => false, + ]) + ->order('number') + ->limit(0, $portionSize) + ->find(); foreach ($itemList as $item) { $this->createQueueFromEvent($item); + $item->set([ 'isProcessed' => true, ]); + $this->entityManager->saveEntity($item); } } protected function createQueueFromEvent(WebhookEventQueueItem $item) { - $webhookList = $this->entityManager->getRepository('Webhook')->where([ - 'event' => $item->get('event'), - 'isActive' => true, - ])->order('createdAt')->find(); + $webhookList = $this->entityManager + ->getRepository('Webhook') + ->where([ + 'event' => $item->get('event'), + 'isActive' => true, + ]) + ->order('createdAt') + ->find(); foreach ($webhookList as $webhook) { $this->entityManager->createEntity('WebhookQueueItem', [ @@ -129,7 +151,7 @@ class Queue 'status' => 'Pending', 'OR' => [ ['processAt' => null], - ['processAt<=' => DateTime::getSystemNowString()], + ['processAt<=' => DateTimeUtil::getSystemNowString()], ], ], 'groupBy' => ['webhookId'], @@ -143,16 +165,22 @@ class Queue foreach ($groupedItemList as $group) { $webhookId = $group->get('webhookId'); - $itemList = $this->entityManager->getRepository('WebhookQueueItem')->where([ - 'webhookId' => $webhookId, - 'status' => 'Pending', - 'OR' => [ - ['processAt' => null], - ['processAt<=' => DateTime::getSystemNowString()], - ], - ])->order('number')->limit(0, $batchSize)->find(); + $itemList = $this->entityManager + ->getRepository('WebhookQueueItem') + ->where([ + 'webhookId' => $webhookId, + 'status' => 'Pending', + 'OR' => [ + ['processAt' => null], + ['processAt<=' => DateTimeUtil::getSystemNowString()], + ], + ]) + ->order('number') + ->limit(0, $batchSize) + ->find(); $webhook = $this->entityManager->getEntity('Webhook', $webhookId); + if (!$webhook || !$webhook->get('isActive')) { foreach ($itemList as $item) { $this->deleteQueueItem($item); @@ -160,43 +188,58 @@ class Queue } $forbiddenAttributeList = []; + $user = null; + if ($webhook->get('userId')) { $user = $this->entityManager->getEntity('User', $webhook->get('userId')); + if (!$user) { foreach ($itemList as $item) { $this->deleteQueueItem($item); } + continue; - } else { - $forbiddenAttributeList = $this->aclManager->getScopeForbiddenAttributeList($user, $webhook->get('entityType')); + } + else { + $forbiddenAttributeList = $this->aclManager + ->getScopeForbiddenAttributeList($user, $webhook->get('entityType')); } } $actualItemList = []; $dataList = []; + foreach ($itemList as $item) { $targetType = $item->get('targetType'); $target = null; + if ($this->entityManager->hasRepository($targetType)) { - $target = $this->entityManager->getRepository($targetType)->where([ - 'id' => $item->get('targetId') - ])->findOne(['withDeleted' => true]); + $target = $this->entityManager + ->getRepository($targetType) + ->where([ + 'id' => $item->get('targetId') + ]) + ->findOne(['withDeleted' => true]); } + if (!$target) { $this->deleteQueueItem($item); + continue; } if ($user) { if (!$this->aclManager->check($user, $target)) { $this->deleteQueueItem($item); + continue; } } $data = $item->get('data') ?? (object) []; + $data = clone $data; foreach ($forbiddenAttributeList as $attribute) { @@ -204,9 +247,13 @@ class Queue } $actualItemList[] = $item; + $dataList[] = $data; } - if (empty($dataList)) continue; + + if (empty($dataList)) { + continue; + } $this->send($webhook, $dataList, $actualItemList); } @@ -216,21 +263,29 @@ class Queue { try { $code = $this->sender->send($webhook, $dataList); - } catch (\Exception $e) { + } + catch (Exception $e) { $this->failQueueItemList($itemList, true); - $GLOBALS['log']->error("Webhook Queue: Webhook {$webhook->id} sending failed. Error: " . $e->getMessage()); + + $GLOBALS['log']->error( + "Webhook Queue: Webhook {$webhook->id} sending failed. Error: " . $e->getMessage() + ); return; } if ($code >= 200 && $code < 400) { $this->succeedQueueItemList($itemList); - } else if ($code === 410) { + } + else if ($code === 410) { $this->dropWebhook($webhook); - } else if (in_array($code, [0, 401, 403, 404, 405, 408, 500, 503])) { + } + else if (in_array($code, [0, 401, 403, 404, 405, 408, 500, 503])) { $this->failQueueItemList($itemList); - } else if ($code >= 400 && $code < 500) { + } + else if ($code >= 400 && $code < 500) { $this->failQueueItemList($itemList, true); - } else { + } + else { $this->failQueueItemList($itemList, true); } @@ -263,10 +318,14 @@ class Queue protected function dropWebhook(Webhook $webhook) { - $itemList = $this->entityManager->getRepository('WebhookQueueItem')->where([ - 'status' => 'Pending', - 'webhookId' => $webhook->id, - ])->order('number')->find(); + $itemList = $this->entityManager + ->getRepository('WebhookQueueItem') + ->where([ + 'status' => 'Pending', + 'webhookId' => $webhook->id, + ]) + ->order('number') + ->find(); foreach ($itemList as $item) { $this->deleteQueueItem($item); @@ -280,7 +339,7 @@ class Queue $item->set([ 'attempts' => $item->get('attempts') + 1, 'status' => 'Success', - 'processedAt' => DateTime::getSystemNowString(), + 'processedAt' => DateTimeUtil::getSystemNowString(), ]); $this->entityManager->saveEntity($item); @@ -289,17 +348,21 @@ class Queue protected function failQueueItem(WebhookQueueItem $item, bool $force = false) { $attempts = $item->get('attempts') + 1; + $maxAttemptsNumber = $this->config->get('webhookMaxAttemptNumber', self::MAX_ATTEMPT_NUMBER); $period = $this->config->get('webhookFailAttemptPeriod', self::FAIL_ATTEMPT_PERIOD); - if ($force) $maxAttemptsNumber = 0; + if ($force) { + $maxAttemptsNumber = 0; + } + + $dt = new DateTime(); - $dt = new \DateTime(); $dt->modify($period); $item->set([ 'attempts' => $attempts, - 'processAt' => $dt->format(DateTime::$systemDateTimeFormat), + 'processAt' => $dt->format(DateTimeUtil::$systemDateTimeFormat), ]); if ($attempts >= $maxAttemptsNumber) { diff --git a/application/Espo/Core/Webhook/Sender.php b/application/Espo/Core/Webhook/Sender.php index 5f77858572..098ae7d25b 100644 --- a/application/Espo/Core/Webhook/Sender.php +++ b/application/Espo/Core/Webhook/Sender.php @@ -44,19 +44,21 @@ class Sender const TIMEOUT = 10; - protected $config; + private $config; public function __construct(Config $config) { $this->config = $config; } - public function send(Webhook $webhook, array $dataList) : int + public function send(Webhook $webhook, array $dataList): int { $payload = json_encode($dataList); $signature = null; + $secretKey = $webhook->get('secretKey'); + if ($secretKey) { $signature = $this->buildSignature($webhook, $payload, $secretKey); } @@ -65,13 +67,16 @@ class Sender $timeout = $this->config->get('webhookTimeout', self::TIMEOUT); $headerList = []; + $headerList[] = 'Content-Type: application/json'; $headerList[] = 'Content-Length: ' . strlen($payload); + if ($signature) { $headerList[] = 'X-Signature: ' . $signature; } $handler = curl_init($webhook->get('url')); + curl_setopt($handler, \CURLOPT_RETURNTRANSFER, true); curl_setopt($handler, \CURLOPT_FOLLOWLOCATION, true); curl_setopt($handler, \CURLOPT_SSL_VERIFYPEER, true); @@ -86,8 +91,13 @@ class Sender $code = curl_getinfo($handler, \CURLINFO_HTTP_CODE); - if (!is_numeric($code)) $code = 0; - if (!is_int($code)) $code = intval($code); + if (!is_numeric($code)) { + $code = 0; + } + + if (!is_int($code)) { + $code = intval($code); + } if ($errorNumber = curl_errno($handler)) { if (in_array($errorNumber, [\CURLE_OPERATION_TIMEDOUT, \CURLE_OPERATION_TIMEOUTED])) {