From 6a791713ea97649fbb23a09eeada45bac5adf742 Mon Sep 17 00:00:00 2001 From: Yuri Kuznetsov Date: Wed, 16 Jun 2021 12:43:26 +0300 Subject: [PATCH] job group processing --- .../Espo/Core/Job/AbstractGroupJob.php | 135 ++++++++++++++++++ application/Espo/Jobs/ProcessJobGroup.php | 34 +++++ .../Resources/metadata/app/scheduledJobs.json | 3 + .../metadata/entityDefs/ScheduledJob.json | 5 + tests/integration/Espo/Core/Job/JobTest.php | 4 +- 5 files changed, 178 insertions(+), 3 deletions(-) create mode 100644 application/Espo/Core/Job/AbstractGroupJob.php create mode 100644 application/Espo/Jobs/ProcessJobGroup.php diff --git a/application/Espo/Core/Job/AbstractGroupJob.php b/application/Espo/Core/Job/AbstractGroupJob.php new file mode 100644 index 0000000000..b42edb9122 --- /dev/null +++ b/application/Espo/Core/Job/AbstractGroupJob.php @@ -0,0 +1,135 @@ +jobManager = $jobManager; + $this->entityManager = $entityManager; + $this->config = $config; + } + + public function run(JobData $data): void + { + $limit = $this->config->get('jobGroupMaxPortion') ?? self::PORTION_NUMBER; + + $group = $data->get('group'); + + $this->jobManager->processGroup($group, $limit); + } + + public function prepare(ScheduledJobData $data, DateTimeImmutable $executeTime): void + { + $groupList = []; + + $query = $this->entityManager + ->getQueryBuilder() + ->select('group') + ->from(JobEntity::ENTITY_TYPE) + ->where([ + 'status' => JobStatus::PENDING, + 'queue' => null, + 'group!=' => null, + 'executeTime<=' => $executeTime->format(DateTime::SYSTEM_DATE_TIME_FORMAT), + ]) + ->groupBy('group') + ->build(); + + $sth = $this->entityManager->getQueryExecutor()->execute($query); + + while ($row = $sth->fetch()) { + $group = $row['group']; + + if ($group === null) { + continue; + } + + $groupList[] = $group; + } + + if (!count($groupList)) { + return; + } + + foreach ($groupList as $group) { + $existingJob = $this->entityManager + ->getRDBRepository(JobEntity::ENTITY_TYPE) + ->select('id') + ->where([ + 'scheduledJobId' => $data->getId(), + 'targetGroup' => $group, + 'status' => [ + JobStatus::RUNNING, + JobStatus::READY, + JobStatus::PENDING, + ], + ]) + ->findOne(); + + if ($existingJob) { + continue; + } + + $name = $data->getName() . ' :: ' . $group; + + $this->entityManager->createEntity(JobEntity::ENTITY_TYPE, [ + 'scheduledJobId' => $data->getId(), + 'executeTime' => $executeTime->format(DateTime::SYSTEM_DATE_TIME_FORMAT), + 'name' => $name, + 'data' => [ + 'group' => $group, + ], + 'targetGroup' => $group, + ]); + } + } +} diff --git a/application/Espo/Jobs/ProcessJobGroup.php b/application/Espo/Jobs/ProcessJobGroup.php new file mode 100644 index 0000000000..99519c3aed --- /dev/null +++ b/application/Espo/Jobs/ProcessJobGroup.php @@ -0,0 +1,34 @@ + 'group-1', ]); - // @todo Change to `$this->jobManager->process();`. - $this->jobManager->processGroup('group-0', 100); - $this->jobManager->processGroup('group-1', 100); + $this->jobManager->process(); $job1Reloaded = $this->entityManager->getEntity('Job', $job1->getId()); $job2Reloaded = $this->entityManager->getEntity('Job', $job2->getId());