diff --git a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java index 0069310b5d..7806215a51 100644 --- a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java +++ b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java @@ -98,7 +98,7 @@ public class JobManagerTest extends AbstractControllerTest { await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { Job job = findJobById(jobId); assertThat(job.getStatus()).isEqualTo(JobStatus.RUNNING); - assertThat(job.getResult().getSuccessfulCount()).isBetween(1, tasksCount - 1); + assertThat(job.getResult().getSuccessfulCount()).isBetween(0, tasksCount - 1); assertThat(job.getResult().getTotalCount()).isEqualTo(tasksCount); }); await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { diff --git a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java index 59093ad802..c1b8f911ff 100644 --- a/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java +++ b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest_EntityPartitioningStrategy.java @@ -27,17 +27,16 @@ import org.thingsboard.server.dao.service.DaoSqlTest; public class JobManagerTest_EntityPartitioningStrategy extends JobManagerTest { /* - * Some tests are overridden because they are based on - * tenant partitioning strategy (subsequent tasks processing within a tenant) - * */ + * Some tests are overridden because they are based on + * tenant partitioning strategy (subsequent tasks processing within a tenant) + * */ @Override - public void testCancelJob_simulateTaskProcessorRestart() throws Exception { + public void testCancelJob_simulateTaskProcessorRestart() { } @Override public void testSubmitJob_generalError() { - } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java index bd2c73d0d8..6f6d924deb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java @@ -55,11 +55,17 @@ public abstract class JobResult implements Serializable { public void processTaskResult(TaskResult taskResult) { if (taskResult.isSuccess()) { - successfulCount++; + if (totalCount == null || successfulCount < totalCount) { + successfulCount++; + } } else if (taskResult.isDiscarded()) { - discardedCount++; + if (totalCount == null || discardedCount < totalCount) { + discardedCount++; + } } else { - failedCount++; + if (totalCount == null || failedCount < totalCount) { + failedCount++; + } if (results.size() < 100) { // preserving only first 100 errors, not reprocessing if there are more failures results.add(taskResult); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java b/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java index 360aa0063b..40200d6fed 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java @@ -87,11 +87,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi } job.getResult().setCancellationTs(System.currentTimeMillis()); JobStatus prevStatus = job.getStatus(); - if (job.getStatus() == QUEUED) { - job.setStatus(CANCELLED); // setting cancelled status right away, because we don't expect stats for cancelled tasks - } else if (job.getStatus() == PENDING) { - job.setStatus(RUNNING); - } + job.setStatus(CANCELLED); saveJob(tenantId, job, true, prevStatus); } @@ -145,7 +141,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi } } - if (job.getStatus() == RUNNING) { + if (job.getStatus().isOneOf(RUNNING, CANCELLED)) { if (result.getTotalCount() != null && result.getCompletedCount() >= result.getTotalCount()) { if (result.getCancellationTs() > 0) { job.setStatus(CANCELLED); @@ -193,7 +189,7 @@ public class DefaultJobService extends AbstractEntityService implements JobServi private void checkWaitingJobs(TenantId tenantId, JobType jobType) { Job queuedJob = jobDao.findOldestByTenantIdAndTypeAndStatusForUpdate(tenantId, jobType, QUEUED); - if (queuedJob == null) { + if (queuedJob == null || jobDao.existsByTenantIdAndTypeAndStatusOneOf(tenantId, jobType, PENDING, RUNNING)) { return; } queuedJob.setStatus(PENDING);