diff --git a/application/src/main/java/org/thingsboard/server/controller/JobController.java b/application/src/main/java/org/thingsboard/server/controller/JobController.java new file mode 100644 index 0000000000..f7bc5255d3 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/controller/JobController.java @@ -0,0 +1,73 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.controller; + +import io.swagger.v3.oas.annotations.Parameter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.security.access.prepost.PreAuthorize; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestParam; +import org.springframework.web.bind.annotation.RestController; +import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.dao.task.JobService; +import org.thingsboard.server.queue.util.TbCoreComponent; + +import java.util.UUID; + +import static org.thingsboard.server.controller.ControllerConstants.PAGE_NUMBER_DESCRIPTION; +import static org.thingsboard.server.controller.ControllerConstants.PAGE_SIZE_DESCRIPTION; +import static org.thingsboard.server.controller.ControllerConstants.SORT_ORDER_DESCRIPTION; +import static org.thingsboard.server.controller.ControllerConstants.SORT_PROPERTY_DESCRIPTION; + +@RestController +@TbCoreComponent +@RequestMapping("/api") +@RequiredArgsConstructor +@Slf4j +public class JobController extends BaseController { + + private final JobService jobService; + + @GetMapping("/job/{id}") + @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") + public Job getJobById(@PathVariable UUID id) throws ThingsboardException { + return jobService.findJobById(getTenantId(), new JobId(id)); + } + + @GetMapping("/jobs") + @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") + public PageData getJobs(@Parameter(description = PAGE_SIZE_DESCRIPTION, required = true) + @RequestParam int pageSize, + @Parameter(description = PAGE_NUMBER_DESCRIPTION, required = true) + @RequestParam int page, + @Parameter(description = "Case-insensitive 'substring' filter based on job's description") + @RequestParam(required = false) String textSearch, + @Parameter(description = SORT_PROPERTY_DESCRIPTION) + @RequestParam(required = false) String sortProperty, + @Parameter(description = SORT_ORDER_DESCRIPTION) + @RequestParam(required = false) String sortOrder) throws ThingsboardException { + PageLink pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); + return jobService.findJobsByTenantId(getTenantId(), pageLink); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java b/application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java index 3b63d3736f..b5f6c4665f 100644 --- a/application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java @@ -17,9 +17,11 @@ package org.thingsboard.server.service.job; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.ProfileEntityIdInfo; import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.job.CfReprocessingJobConfiguration; import org.thingsboard.server.common.data.job.CfReprocessingTask; @@ -39,28 +41,34 @@ public class CfReprocessingJobProcessor extends JobProcessor { private final DeviceService deviceService; private final AssetService assetService; + // fixme: multiple jobs with single type + @Transactional @Override - public void process(Job job, Consumer taskConsumer) { + public int process(Job job, Consumer taskConsumer) { CfReprocessingJobConfiguration configuration = job.getConfiguration(); CalculatedField calculatedField = configuration.getCalculatedField(); - EntityId entityId = calculatedField.getEntityId(); + EntityId cfEntityId = calculatedField.getEntityId(); - if (entityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) { - taskConsumer.accept(createTask(job, configuration, entityId)); + int tasksCount = 0; + if (cfEntityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) { + taskConsumer.accept(createTask(job, configuration, cfEntityId)); + tasksCount++; } else { - PageDataIterable entities; - if (entityId.getEntityType() == EntityType.DEVICE_PROFILE) { - entities = new PageDataIterable<>(pageLink -> deviceService.findProfileEntityIdInfosByTenantId(job.getTenantId(), pageLink), 512); - } else if (entityId.getEntityType() == EntityType.ASSET_PROFILE) { - entities = new PageDataIterable<>(pageLink -> assetService.findProfileEntityIdInfosByTenantId(job.getTenantId(), pageLink), 512); + PageDataIterable entities; + if (cfEntityId.getEntityType() == EntityType.DEVICE_PROFILE) { + entities = new PageDataIterable<>(pageLink -> deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(job.getTenantId(), (DeviceProfileId) cfEntityId, pageLink), 512); + } else if (cfEntityId.getEntityType() == EntityType.ASSET_PROFILE) { + entities = new PageDataIterable<>(pageLink -> assetService.findAssetIdsByTenantIdAndAssetProfileId(job.getTenantId(), (AssetProfileId) cfEntityId, pageLink), 512); } else { - throw new IllegalArgumentException("Unsupported CF entity type " + entityId.getEntityType()); + throw new IllegalArgumentException("Unsupported CF entity type " + cfEntityId.getEntityType()); + } + for (EntityId entityId : entities) { + taskConsumer.accept(createTask(job, configuration, entityId)); + tasksCount++; } - entities.forEach(device -> { - taskConsumer.accept(createTask(job, configuration, device.getEntityId())); - }); } + return tasksCount; } private Task createTask(Job job, CfReprocessingJobConfiguration configuration, EntityId entityId) { @@ -68,6 +76,7 @@ public class CfReprocessingJobProcessor extends JobProcessor { .tenantId(job.getTenantId()) .jobId(job.getId()) .key(entityId.getEntityType().getNormalName() + " " + entityId.getId()) + .retries(2) // 3 attempts in total .calculatedField(configuration.getCalculatedField()) .entityId(entityId) .startTs(configuration.getStartTs()) diff --git a/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java b/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java index 73367e2af5..64c1615310 100644 --- a/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java +++ b/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java @@ -18,18 +18,20 @@ package org.thingsboard.server.service.job; import jakarta.annotation.PreDestroy; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobStats; import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.job.Task; import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.dao.task.JobService; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; -import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsgMetadata; @@ -37,12 +39,15 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; +import org.thingsboard.server.queue.task.JobStatsService; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbCoreComponent; import java.util.Arrays; +import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.UUID; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.function.Function; @@ -54,23 +59,26 @@ import java.util.stream.Collectors; public class DefaultJobManager implements JobManager { private final JobService jobService; - private final TbCoreQueueFactory queueFactory; + private final JobStatsService jobStatsService; private final Map jobProcessors; private final Map>> taskProducers; - private final QueueConsumerManager> taskResultConsumer; + private final QueueConsumerManager> taskResultConsumer; private final ExecutorService consumerExecutor; - public DefaultJobManager(JobService jobService, TbCoreQueueFactory queueFactory, List jobProcessors) { + @Value("${queue.tasks.stats.processing_interval_ms:5000}") + private int statsProcessingInterval; + + public DefaultJobManager(JobService jobService, JobStatsService jobStatsService, TbCoreQueueFactory queueFactory, List jobProcessors) { this.jobService = jobService; - this.queueFactory = queueFactory; + this.jobStatsService = jobStatsService; this.jobProcessors = jobProcessors.stream().collect(Collectors.toMap(JobProcessor::getType, Function.identity())); this.taskProducers = Arrays.stream(JobType.values()).collect(Collectors.toMap(Function.identity(), queueFactory::createTaskProducer)); this.consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("task-result-consumer")); - this.taskResultConsumer = QueueConsumerManager.>builder() // fixme: should be consumer per partition - .name("tasks-results") - .msgPackProcessor(this::processResults) + this.taskResultConsumer = QueueConsumerManager.>builder() + .name("job-stats") + .msgPackProcessor(this::processStats) .pollInterval(125) - .consumerCreator(queueFactory::createTaskResultConsumer) + .consumerCreator(queueFactory::createJobStatsConsumer) .consumerExecutor(consumerExecutor) .build(); } @@ -82,10 +90,13 @@ public class DefaultJobManager implements JobManager { } @Override - public void submitJob(Job job) { + public Job submitJob(Job job) { job = jobService.createJob(job.getTenantId(), job); log.info("Submitting job: {}", job); - jobProcessors.get(job.getType()).process(job, this::submitTask); + + int tasksCount = jobProcessors.get(job.getType()).process(job, this::submitTask); + jobStatsService.reportAllTasksSubmitted(job.getId(), tasksCount); + return job; } private void submitTask(Task task) { @@ -110,21 +121,34 @@ public class DefaultJobManager implements JobManager { } @SneakyThrows - private void processResults(List> msgs, TbQueueConsumer> consumer) { - Map> results = msgs.stream() - .map(msg -> JacksonUtil.fromString(msg.getValue().getValue(), TaskResult.class)) - .collect(Collectors.groupingBy(TaskResult::getJobId)); - results.forEach((jobId, taskResults) -> { + private void processStats(List> msgs, TbQueueConsumer> consumer) { + Map stats = new HashMap<>(); + + for (TbProtoQueueMsg msg : msgs) { + JobStatsMsg statsMsg = msg.getValue(); + JobId jobId = new JobId(new UUID(statsMsg.getJobIdMSB(), statsMsg.getJobIdLSB())); + JobStats jobStats = stats.computeIfAbsent(jobId, JobStats::new); + + if (statsMsg.hasTaskResult()) { + TaskResult taskResult = JacksonUtil.fromString(statsMsg.getTaskResult().getValue(), TaskResult.class); + jobStats.getTaskResults().add(taskResult); + } + if (statsMsg.hasTotalTasksCount()) { + jobStats.setTotalTasksCount(statsMsg.getTotalTasksCount()); + } + } + + stats.forEach((jobId, jobStats) -> { try { - log.info("[{}] Processing task results: {}", jobId, taskResults); - jobService.reportTaskResults(jobId, taskResults); + log.info("[{}] Processing job stats: {}", jobId, stats); + jobService.processStats(jobId, jobStats); } catch (Exception e) { - log.warn("Failed to report task results for job {}: {}", jobId, taskResults, e); + log.warn("Failed to process job stats for {}: {}", jobId, jobStats, e); } }); consumer.commit(); - Thread.sleep(5000); + Thread.sleep(statsProcessingInterval); } @PreDestroy diff --git a/application/src/main/java/org/thingsboard/server/service/job/DummyJobProcessor.java b/application/src/main/java/org/thingsboard/server/service/job/DummyJobProcessor.java new file mode 100644 index 0000000000..bed8f3f25e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/DummyJobProcessor.java @@ -0,0 +1,64 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.job; + +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.job.DummyJobConfiguration; +import org.thingsboard.server.common.data.job.DummyTask; +import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobType; +import org.thingsboard.server.common.data.job.Task; + +import java.util.List; +import java.util.function.Consumer; + +@Component +@RequiredArgsConstructor +public class DummyJobProcessor extends JobProcessor { + + @Override + public int process(Job job, Consumer taskConsumer) { + DummyJobConfiguration configuration = job.getConfiguration(); + for (int number = 1; number <= configuration.getSuccessfulTasksCount(); number++) { + taskConsumer.accept(createTask(job, configuration, number, null)); + } + if (configuration.getErrors() != null) { + for (int number = 1; number <= configuration.getFailedTasksCount(); number++) { + taskConsumer.accept(createTask(job, configuration, number, configuration.getErrors())); + } + } + return configuration.getSuccessfulTasksCount() + configuration.getFailedTasksCount(); + } + + private Task createTask(Job job, DummyJobConfiguration configuration, int number, List errors) { + return DummyTask.builder() + .tenantId(job.getTenantId()) + .jobId(job.getId()) + .key("Task " + number) + .retries(configuration.getRetries()) + .number(number) + .processingTimeMs(configuration.getTaskProcessingTimeMs()) + .errors(errors) + .build(); + } + + @Override + public JobType getType() { + return JobType.DUMMY; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/job/JobManager.java b/application/src/main/java/org/thingsboard/server/service/job/JobManager.java index 5db52ef448..c78de2a5d9 100644 --- a/application/src/main/java/org/thingsboard/server/service/job/JobManager.java +++ b/application/src/main/java/org/thingsboard/server/service/job/JobManager.java @@ -19,6 +19,6 @@ import org.thingsboard.server.common.data.job.Job; public interface JobManager { - void submitJob(Job job); + Job submitJob(Job job); } diff --git a/application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java b/application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java index cf15209616..01f7291dd3 100644 --- a/application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java @@ -23,7 +23,7 @@ import java.util.function.Consumer; public abstract class JobProcessor { - public abstract void process(Job job, Consumer taskConsumer); + public abstract int process(Job job, Consumer taskConsumer); public abstract JobType getType(); diff --git a/application/src/main/java/org/thingsboard/server/service/job/task/DummyTaskProcessor.java b/application/src/main/java/org/thingsboard/server/service/job/task/DummyTaskProcessor.java new file mode 100644 index 0000000000..10aafc197d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/task/DummyTaskProcessor.java @@ -0,0 +1,44 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.job.task; + +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.job.DummyTask; +import org.thingsboard.server.common.data.job.JobType; +import org.thingsboard.server.queue.task.TaskProcessor; + +@Component +@RequiredArgsConstructor +public class DummyTaskProcessor extends TaskProcessor { + + @Override + protected void process(DummyTask task) throws Exception { + if (task.getProcessingTimeMs() > 0) { + Thread.sleep(task.getProcessingTimeMs()); + } + if (task.getErrors() != null && task.getAttempt() <= task.getErrors().size()) { + String error = task.getErrors().get(task.getAttempt() - 1); + throw new RuntimeException(error); + } + } + + @Override + public JobType getJobType() { + return JobType.DUMMY; + } + +} diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 6670140b30..322b037e2d 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1883,6 +1883,10 @@ queue: enabled: "${TB_QUEUE_EDGE_STATS_ENABLED:true}" # Statistics printing interval for Edge services print-interval-ms: "${TB_QUEUE_EDGE_STATS_PRINT_INTERVAL_MS:60000}" + tasks: + stats: + # Interval in milliseconds to process job stats + processing_interval_ms: "${TB_QUEUE_TASKS_STATS_PROCESSING_INTERVAL_MS:5000}" # Event configuration parameters event: 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 new file mode 100644 index 0000000000..f9ee2e2b13 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/job/JobManagerTest.java @@ -0,0 +1,125 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.job; + +import com.fasterxml.jackson.core.type.TypeReference; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.test.context.TestPropertySource; +import org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.job.DummyJobConfiguration; +import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobStatus; +import org.thingsboard.server.common.data.job.JobType; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.controller.AbstractControllerTest; +import org.thingsboard.server.dao.service.DaoSqlTest; + +import java.util.List; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; + +@DaoSqlTest +@TestPropertySource(properties = { + "queue.tasks.stats.processing_interval_ms=0" +}) +public class JobManagerTest extends AbstractControllerTest { + + @Autowired + private JobManager jobManager; + + @Before + public void setUp() throws Exception { + loginTenantAdmin(); + } + + @After + public void tearDown() throws Exception { + } + + @Test + public void testSubmitJob_allTasksSuccessful() { + int tasksCount = 5; + JobId jobId = jobManager.submitJob(Job.builder() + .tenantId(tenantId) + .type(JobType.DUMMY) + .key("test-job") + .description("test job") + .configuration(DummyJobConfiguration.builder() + .successfulTasksCount(tasksCount) + .taskProcessingTimeMs(1000) + .build()) + .build()).getId(); + + 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().getTotalCount()).isEqualTo(tasksCount); + }); + await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { + Job job = findJobById(jobId); + assertThat(job.getStatus()).isEqualTo(JobStatus.COMPLETED); + assertThat(job.getResult().getSuccessfulCount()).isEqualTo(tasksCount); + assertThat(job.getResult().getFailures()).isEmpty(); + }); + } + + @Test + public void testSubmitJob_someTasksPermanentlyFailed() { + int successfulTasks = 3; + int failedTasks = 2; + JobId jobId = jobManager.submitJob(Job.builder() + .tenantId(tenantId) + .type(JobType.DUMMY) + .key("test-job") + .description("test job") + .configuration(DummyJobConfiguration.builder() + .successfulTasksCount(successfulTasks) + .failedTasksCount(failedTasks) + .errors(List.of("error1", "error2", "error3")) + .retries(2) + .taskProcessingTimeMs(100) + .build()) + .build()).getId(); + + await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { + Job job = findJobById(jobId); + assertThat(job.getStatus()).isEqualTo(JobStatus.FAILED); + assertThat(job.getResult().getSuccessfulCount()).isEqualTo(successfulTasks); + assertThat(job.getResult().getFailedCount()).isEqualTo(failedTasks); + assertThat(job.getResult().getTotalCount()).isEqualTo(successfulTasks + failedTasks); + assertThat(job.getResult().getFailures().get("Task 1")).isEqualTo("error3"); // last error + assertThat(job.getResult().getFailures().get("Task 2")).isEqualTo("error3"); // last error + }); + } + + + + private Job findJobById(JobId jobId) throws Exception { + return doGet("/api/job/" + jobId, Job.class); + } + + private List findJobs() throws Exception { + return doGetTypedWithPageLink("/api/jobs?", new TypeReference>() {}, new PageLink(100, 0)).getData(); + } + +} \ No newline at end of file diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java index 17e7645605..5581f1eac7 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java @@ -17,18 +17,18 @@ package org.thingsboard.server.dao.task; import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.job.TaskResult; +import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobStats; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.job.Job; - -import java.util.List; public interface JobService { Job createJob(TenantId tenantId, Job job); - void reportTaskResults(JobId jobId, List results); + Job findJobById(TenantId tenantId, JobId jobId); + + void processStats(JobId jobId, JobStats jobStats); PageData findJobsByTenantId(TenantId tenantId, PageLink pageLink); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java index 4f5a33dacd..797dd0b639 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.data.job; +import jakarta.validation.constraints.NotNull; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; @@ -27,6 +28,7 @@ import org.thingsboard.server.common.data.cf.CalculatedField; @Builder public class CfReprocessingJobConfiguration implements JobConfiguration { + @NotNull private CalculatedField calculatedField; private long startTs; private long endTs; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java index 0e15b9473f..5c380c4dfa 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java @@ -15,20 +15,17 @@ */ package org.thingsboard.server.common.data.job; -import lombok.AllArgsConstructor; -import lombok.Builder; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import lombok.experimental.SuperBuilder; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.JobId; -import org.thingsboard.server.common.data.id.TenantId; @Data -@AllArgsConstructor @NoArgsConstructor @EqualsAndHashCode(callSuper = true) +@SuperBuilder public class CfReprocessingTask extends Task { private CalculatedField calculatedField; @@ -36,15 +33,6 @@ public class CfReprocessingTask extends Task { private long startTs; private long endTs; - @Builder - public CfReprocessingTask(TenantId tenantId, JobId jobId, String key, CalculatedField calculatedField, EntityId entityId, long startTs, long endTs) { - super(tenantId, jobId, key); - this.calculatedField = calculatedField; - this.entityId = entityId; - this.startTs = startTs; - this.endTs = endTs; - } - @Override public JobType getJobType() { return JobType.CF_REPROCESSING; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobConfiguration.java new file mode 100644 index 0000000000..7fe621c058 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobConfiguration.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.job; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.util.List; + +@Data +@AllArgsConstructor +@NoArgsConstructor +@Builder +public class DummyJobConfiguration implements JobConfiguration { + + private long taskProcessingTimeMs; + private int successfulTasksCount; + private int failedTasksCount; + private List errors; + private int retries; + + @Override + public JobType getType() { + return JobType.DUMMY; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobResult.java new file mode 100644 index 0000000000..031a733d51 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyJobResult.java @@ -0,0 +1,25 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.job; + +public class DummyJobResult extends JobResult { + + @Override + public JobType getJobType() { + return JobType.DUMMY; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyTask.java new file mode 100644 index 0000000000..dee97fc3c9 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/DummyTask.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.job; + +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +import lombok.experimental.SuperBuilder; + +import java.util.List; + +@Data +@NoArgsConstructor +@EqualsAndHashCode(callSuper = true) +@SuperBuilder +public class DummyTask extends Task { + + private int number; + private long processingTimeMs; + private List errors; // errors for each attempt + + @Override + public JobType getJobType() { + return JobType.DUMMY; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java index e0161cf98e..5a51a49239 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.common.data.job; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.NotNull; import lombok.Builder; import lombok.Data; import lombok.EqualsAndHashCode; @@ -29,22 +31,30 @@ import org.thingsboard.server.common.data.id.TenantId; @EqualsAndHashCode(callSuper = true) public class Job extends BaseData implements HasTenantId { + @NotNull private TenantId tenantId; + @NotNull private JobType type; + @NotBlank private String key; + @NotBlank + private String description; private JobStatus status; + @NotNull private JobConfiguration configuration; private JobResult result; @Builder - public Job(TenantId tenantId, JobType type, String key, JobConfiguration configuration) { + public Job(TenantId tenantId, JobType type, String key, String description, JobConfiguration configuration) { this.tenantId = tenantId; this.type = type; this.key = key; + this.description = description; this.configuration = configuration; this.status = JobStatus.PENDING; this.result = switch (type) { case CF_REPROCESSING -> new CfReprocessingJobResult(); + case DUMMY -> new DummyJobResult(); }; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java index 8899206c49..eccdfae6bb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java @@ -24,6 +24,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type") @JsonSubTypes({ @Type(name = "CF_REPROCESSING", value = CfReprocessingJobConfiguration.class), + @Type(name = "DUMMY", value = DummyJobConfiguration.class), }) public interface JobConfiguration { 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 406eb0f65e..07b6c4eadd 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 @@ -17,8 +17,10 @@ package org.thingsboard.server.common.data.job; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonSubTypes.Type; import com.fasterxml.jackson.annotation.JsonTypeInfo; import lombok.Data; +import lombok.NoArgsConstructor; import java.util.HashMap; import java.util.Map; @@ -26,14 +28,16 @@ import java.util.Map; @JsonIgnoreProperties(ignoreUnknown = true) @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType") @JsonSubTypes({ - @JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingJobResult.class), + @Type(name = "CF_REPROCESSING", value = CfReprocessingJobResult.class), + @Type(name = "DUMMY", value = DummyJobResult.class) }) @Data +@NoArgsConstructor public abstract class JobResult { private int successfulCount; private int failedCount; - private int totalCount; + private Integer totalCount = null; // set when all tasks are submitted private Map failures = new HashMap<>(); public abstract JobType getJobType(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobStats.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobStats.java new file mode 100644 index 0000000000..6491d0998a --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobStats.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.job; + +import lombok.Data; +import org.thingsboard.server.common.data.id.JobId; + +import java.util.ArrayList; +import java.util.List; + +@Data +public class JobStats { + private final JobId jobId; + private final List taskResults = new ArrayList<>(); + private Integer totalTasksCount; +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java index 60cac8173f..7c0d9972e1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java @@ -17,7 +17,8 @@ package org.thingsboard.server.common.data.job; public enum JobType { - CF_REPROCESSING; + CF_REPROCESSING, + DUMMY; public String getTasksTopic() { return "tasks." + name().toLowerCase(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java index afeaeba393..5ef735f5d3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java @@ -17,8 +17,11 @@ package org.thingsboard.server.common.data.job; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonSubTypes.Type; import com.fasterxml.jackson.annotation.JsonTypeInfo; +import lombok.AllArgsConstructor; import lombok.Data; +import lombok.experimental.SuperBuilder; import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.TenantId; @@ -26,19 +29,17 @@ import org.thingsboard.server.common.data.id.TenantId; @JsonIgnoreProperties(ignoreUnknown = true) @JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType") @JsonSubTypes({ - @JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingTask.class), + @Type(name = "CF_REPROCESSING", value = CfReprocessingTask.class), + @Type(name = "DUMMY", value = DummyTask.class) }) +@SuperBuilder +@AllArgsConstructor public abstract class Task { private TenantId tenantId; private JobId jobId; private String key; - - public Task(TenantId tenantId, JobId jobId, String key) { - this.tenantId = tenantId; - this.jobId = jobId; - this.key = key; - } + private int retries; public Task() { } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java index bfbef46180..ae5e6bbba9 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java @@ -19,8 +19,6 @@ import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; -import org.thingsboard.server.common.data.id.JobId; -import org.thingsboard.server.common.data.id.TenantId; @Data @AllArgsConstructor @@ -28,8 +26,6 @@ import org.thingsboard.server.common.data.id.TenantId; @Builder public class TaskResult { - private TenantId tenantId; - private JobId jobId; private boolean success; private TaskFailure failure; diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index de03b11b6a..032301206f 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -1851,6 +1851,13 @@ message TaskProto { string value = 1; // fixme: TMP, make more efficient } +message JobStatsMsg { + int64 jobIdMSB = 1; + int64 jobIdLSB = 2; + optional TaskResultProto taskResult = 3; + optional int32 totalTasksCount = 4; +} + message TaskResultProto { string value = 1; // fixme: TMP, make more efficient } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java index c8866e52b4..9160818278 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java @@ -28,6 +28,7 @@ import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; @@ -270,13 +271,13 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueProducer> createTaskResultProducer() { - return new InMemoryTbQueueProducer<>(storage, "tasks.results"); + public TbQueueProducer> createJobStatsProducer() { + return new InMemoryTbQueueProducer<>(storage, "jobs.stats"); } @Override - public TbQueueConsumer> createTaskResultConsumer() { - return new InMemoryTbQueueConsumer<>(storage, "tasks.results"); + public TbQueueConsumer> createJobStatsConsumer() { + return new InMemoryTbQueueConsumer<>(storage, "jobs.stats"); } @Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}") diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java index 809099aa8b..e4b51eaa1f 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java @@ -30,8 +30,8 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; -import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -670,23 +670,23 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueProducer> createTaskResultProducer() { - return TbKafkaProducerTemplate.>builder() - .clientId("task-result-producer-" + serviceInfoProvider.getServiceId()) - .defaultTopic(topicService.buildTopicName("tasks.results")) + public TbQueueProducer> createJobStatsProducer() { + return TbKafkaProducerTemplate.>builder() + .clientId("job-stats-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName("jobs.stats")) .settings(kafkaSettings) .admin(tasksAdmin) .build(); } @Override - public TbQueueConsumer> createTaskResultConsumer() { - return TbKafkaConsumerTemplate.>builder() + public TbQueueConsumer> createJobStatsConsumer() { + return TbKafkaConsumerTemplate.>builder() .settings(kafkaSettings) - .topic(topicService.buildTopicName("tasks.results")) - .clientId("task-result-consumer-" + serviceInfoProvider.getServiceId()) - .groupId(topicService.buildTopicName("task-result-consumer-group")) - .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskResultProto.parseFrom(msg.getData()), msg.getHeaders())) + .topic(topicService.buildTopicName("jobs.stats")) + .clientId("job-stats-consumer-" + serviceInfoProvider.getServiceId()) + .groupId(topicService.buildTopicName("job-stats-consumer-group")) + .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), JobStatsMsg.parseFrom(msg.getData()), msg.getHeaders())) .admin(tasksAdmin) .statsService(consumerStatsService) .build(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java index 71a9669ba4..85d1e8be14 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java @@ -25,8 +25,8 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.gen.js.JsInvokeProtos; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; @@ -536,13 +536,36 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { } @Override - public TbQueueConsumer> createTaskResultConsumer() { - return TbKafkaConsumerTemplate.>builder() + public TbQueueConsumer> createTaskConsumer(JobType jobType) { + return TbKafkaConsumerTemplate.>builder() .settings(kafkaSettings) - .topic(topicService.buildTopicName("tasks.results")) - .clientId("task-result-consumer-" + serviceInfoProvider.getServiceId()) - .groupId(topicService.buildTopicName("task-result-consumer-group")) - .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.TaskResultProto.parseFrom(msg.getData()), msg.getHeaders())) + .topic(topicService.buildTopicName(jobType.getTasksTopic())) + .clientId(jobType.name().toLowerCase() + "-task-consumer-" + serviceInfoProvider.getServiceId()) + .groupId(topicService.buildTopicName(jobType.name().toLowerCase() + "-task-consumer-group")) + .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TaskProto.parseFrom(msg.getData()), msg.getHeaders())) + .admin(tasksAdmin) + .statsService(consumerStatsService) + .build(); + } + + @Override + public TbQueueProducer> createJobStatsProducer() { + return TbKafkaProducerTemplate.>builder() + .clientId("job-stats-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName("jobs.stats")) + .settings(kafkaSettings) + .admin(tasksAdmin) + .build(); + } + + @Override + public TbQueueConsumer> createJobStatsConsumer() { + return TbKafkaConsumerTemplate.>builder() + .settings(kafkaSettings) + .topic(topicService.buildTopicName("jobs.stats")) + .clientId("job-stats-consumer-" + serviceInfoProvider.getServiceId()) + .groupId(topicService.buildTopicName("job-stats-consumer-group")) + .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), JobStatsMsg.parseFrom(msg.getData()), msg.getHeaders())) .admin(tasksAdmin) .statsService(consumerStatsService) .build(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index d0ac20ae06..b2af3671c6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -27,9 +27,9 @@ import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; @@ -433,10 +433,10 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueProducer> createTaskResultProducer() { - return TbKafkaProducerTemplate.>builder() - .clientId("task-result-producer-" + serviceInfoProvider.getServiceId()) - .defaultTopic(topicService.buildTopicName("tasks.results")) + public TbQueueProducer> createJobStatsProducer() { + return TbKafkaProducerTemplate.>builder() + .clientId("job-stats-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName("jobs.stats")) .settings(kafkaSettings) .admin(tasksAdmin) .build(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java index 10b84f0f65..571b14639c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java @@ -16,8 +16,8 @@ package org.thingsboard.server.queue.provider; import org.thingsboard.server.common.data.job.JobType; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; -import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; @@ -26,6 +26,6 @@ public interface TaskProcessorQueueFactory { TbQueueConsumer> createTaskConsumer(JobType jobType); - TbQueueProducer> createTaskResultProducer(); + TbQueueProducer> createJobStatsProducer(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java index 409348a12a..823ebea298 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java @@ -19,8 +19,8 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.gen.js.JsInvokeProtos; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; -import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -47,7 +47,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; * Responsible for initialization of various Producers and Consumers used by TB Core Node. * Implementation Depends on the queue queue.type from yml or TB_QUEUE_TYPE environment variable */ -public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory { +public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory, TaskProcessorQueueFactory { /** * Used to push messages to instances of TB Transport Service @@ -170,6 +170,6 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, Hous TbQueueProducer> createTaskProducer(JobType jobType); - TbQueueConsumer> createTaskResultConsumer(); + TbQueueConsumer> createJobStatsConsumer(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java index 98a3d78304..9900474a10 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueProducerProvider.java @@ -18,6 +18,7 @@ package org.thingsboard.server.queue.provider; import jakarta.annotation.PostConstruct; import org.springframework.stereotype.Service; import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -53,6 +54,7 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider { private TbQueueProducer> toHousekeeper; private TbQueueProducer> toCalculatedFields; private TbQueueProducer> toCalculatedFieldNotifications; + private TbQueueProducer> jobStatsProducer; public TbCoreQueueProducerProvider(TbCoreQueueFactory tbQueueProvider) { this.tbQueueProvider = tbQueueProvider; @@ -73,6 +75,7 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider { this.toEdgeEvents = tbQueueProvider.createEdgeEventMsgProducer(); this.toCalculatedFields = tbQueueProvider.createToCalculatedFieldMsgProducer(); this.toCalculatedFieldNotifications = tbQueueProvider.createToCalculatedFieldNotificationMsgProducer(); + this.jobStatsProducer = tbQueueProvider.createJobStatsProducer(); } @Override @@ -140,4 +143,9 @@ public class TbCoreQueueProducerProvider implements TbQueueProducerProvider { return toCalculatedFieldNotifications; } + @Override + public TbQueueProducer> getJobStatsProducer() { + return jobStatsProducer; + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java index 865637b2ff..428e673fa8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbQueueProducerProvider.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.queue.provider; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -97,4 +98,6 @@ public interface TbQueueProducerProvider { TbQueueProducer> getCalculatedFieldsNotificationsMsgProducer(); + TbQueueProducer> getJobStatsProducer(); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java index 8e1952fc14..9e77a2d4e7 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineProducerProvider.java @@ -18,6 +18,7 @@ package org.thingsboard.server.queue.provider; import jakarta.annotation.PostConstruct; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -51,6 +52,7 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider { private TbQueueProducer> toEdgeEvents; private TbQueueProducer> toCalculatedFields; private TbQueueProducer> toCalculatedFieldNotifications; + private TbQueueProducer> jobStatsProducer; public TbRuleEngineProducerProvider(TbRuleEngineQueueFactory tbQueueProvider) { this.tbQueueProvider = tbQueueProvider; @@ -70,6 +72,7 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider { this.toEdgeEvents = tbQueueProvider.createEdgeEventMsgProducer(); this.toCalculatedFields = tbQueueProvider.createToCalculatedFieldMsgProducer(); this.toCalculatedFieldNotifications = tbQueueProvider.createToCalculatedFieldNotificationMsgProducer(); + this.jobStatsProducer = tbQueueProvider.createJobStatsProducer(); } @Override @@ -137,4 +140,9 @@ public class TbRuleEngineProducerProvider implements TbQueueProducerProvider { return toCalculatedFieldNotifications; } + @Override + public TbQueueProducer> getJobStatsProducer() { + return jobStatsProducer; + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java index a7a34992cd..4472c6157e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbTransportQueueProducerProvider.java @@ -121,4 +121,10 @@ public class TbTransportQueueProducerProvider implements TbQueueProducerProvider public TbQueueProducer> getCalculatedFieldsNotificationsMsgProducer() { throw new RuntimeException("Not Implemented! Should not be used by Transport!"); } + + @Override + public TbQueueProducer> getJobStatsProducer() { + throw new RuntimeException("Not Implemented! Should not be used by Transport!"); + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java index 85c400d094..0370f8a4af 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbVersionControlProducerProvider.java @@ -118,4 +118,9 @@ public class TbVersionControlProducerProvider implements TbQueueProducerProvider throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); } + @Override + public TbQueueProducer> getJobStatsProducer() { + throw new RuntimeException("Not Implemented! Should not be used by Version Control Service!"); + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/task/JobStatsService.java b/common/queue/src/main/java/org/thingsboard/server/queue/task/JobStatsService.java new file mode 100644 index 0000000000..6c36c573ca --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/task/JobStatsService.java @@ -0,0 +1,63 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.queue.task; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Lazy; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.job.TaskResult; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.gen.transport.TransportProtos.JobStatsMsg; +import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto; +import org.thingsboard.server.queue.TbQueueCallback; +import org.thingsboard.server.queue.TbQueueProducer; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.provider.TbQueueProducerProvider; + +@Lazy +@Service +@Slf4j +@RequiredArgsConstructor +public class JobStatsService { + + private final TbQueueProducerProvider producerProvider; + + public void reportTaskResult(JobId jobId, TaskResult result) { + report(jobId, JobStatsMsg.newBuilder() + .setTaskResult(TaskResultProto.newBuilder() + .setValue(JacksonUtil.toString(result)) + .build())); + } + + public void reportAllTasksSubmitted(JobId jobId, int tasksCount) { + report(jobId, JobStatsMsg.newBuilder() + .setTotalTasksCount(tasksCount)); + } + + private void report(JobId jobId, JobStatsMsg.Builder statsMsg) { + log.info("[{}] Reporting: {}", jobId, statsMsg); + statsMsg.setJobIdMSB(jobId.getId().getMostSignificantBits()) + .setJobIdLSB(jobId.getId().getLeastSignificantBits()); + + TbProtoQueueMsg msg = new TbProtoQueueMsg<>(jobId.getId(), statsMsg.build()); + TbQueueProducer> producer = producerProvider.getJobStatsProducer(); + producer.send(TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(), msg, TbQueueCallback.EMPTY); + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java b/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java index 7123eaf1e0..52fd382c3d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java @@ -25,12 +25,8 @@ import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.job.Task; import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.data.job.TaskResult.TaskFailure; -import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos.TaskProto; -import org.thingsboard.server.gen.transport.TransportProtos.TaskResultProto; -import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueConsumer; -import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.provider.TaskProcessorQueueFactory; @@ -45,22 +41,22 @@ public abstract class TaskProcessor { @Autowired private TaskProcessorQueueFactory queueFactory; + @Autowired + private JobStatsService statsService; private QueueConsumerManager> taskConsumer; - private TbQueueProducer> taskResultProducer; private ExecutorService consumerExecutor; @PostConstruct public void init() { consumerExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName(getJobType().name().toLowerCase() + "-task-consumer")); taskConsumer = QueueConsumerManager.>builder() // fixme: should be consumer per partition - .name(getJobType().name() + "-tasks") + .name(getJobType().name().toLowerCase() + "-tasks") .msgPackProcessor(this::processMsgs) .pollInterval(125) .consumerCreator(() -> queueFactory.createTaskConsumer(getJobType())) .consumerExecutor(consumerExecutor) .build(); - taskResultProducer = queueFactory.createTaskResultProducer(); } @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) @@ -92,7 +88,7 @@ public abstract class TaskProcessor { reportSuccess(task); } catch (Exception e) { log.error("Failed to process task (attempt {}): {}", task.getAttempt(), task, e); - if (task.getAttempt() < 3) { + if (task.getAttempt() <= task.getRetries()) { processTask(task); } else { reportFailure(task, e); @@ -102,34 +98,19 @@ public abstract class TaskProcessor { private void reportSuccess(Task task) { TaskResult result = TaskResult.builder() - .tenantId(task.getTenantId()) - .jobId(task.getJobId()) .success(true) .build(); - reportResult(result); + statsService.reportTaskResult(task.getJobId(), result); } private void reportFailure(Task task, Throwable error) { TaskResult result = TaskResult.builder() - .tenantId(task.getTenantId()) - .jobId(task.getJobId()) .failure(TaskFailure.builder() .error(error.getMessage()) .task(task) .build()) .build(); - reportResult(result); - } - - private void reportResult(TaskResult result) { - log.info("Reporting result: {}", result); - TaskResultProto resultProto = TaskResultProto.newBuilder() - .setValue(JacksonUtil.toString(result)) - .build(); - TbProtoQueueMsg msg = new TbProtoQueueMsg<>(result.getJobId().getId(), resultProto); - taskResultProducer.send(TopicPartitionInfo.builder() - .topic(taskResultProducer.getDefaultTopic()) - .build(), msg, TbQueueCallback.EMPTY); + statsService.reportTaskResult(task.getJobId(), result); } protected abstract void process(T task) throws Exception; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 1d504b0ab2..707dcdc3a5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -745,6 +745,7 @@ public class ModelConstants { public static final String JOB_TABLE_NAME = "job"; public static final String JOB_TYPE_PROPERTY = "type"; public static final String JOB_KEY_PROPERTY = "key"; + public static final String JOB_DESCRIPTION_PROPERTY = "description"; public static final String JOB_STATUS_PROPERTY = "status"; public static final String JOB_CONFIGURATION_PROPERTY = "configuration"; public static final String JOB_RESULT_PROPERTY = "result"; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java index 6e8b3e7958..d13c4bbed3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java @@ -56,6 +56,9 @@ public class JobEntity extends BaseSqlEntity { @Column(name = ModelConstants.JOB_KEY_PROPERTY, nullable = false) private String key; + @Column(name = ModelConstants.JOB_DESCRIPTION_PROPERTY, nullable = false) + private String description; + @Enumerated(EnumType.STRING) @Column(name = ModelConstants.JOB_STATUS_PROPERTY, nullable = false) private JobStatus status; @@ -74,6 +77,7 @@ public class JobEntity extends BaseSqlEntity { this.tenantId = getTenantUuid(job.getTenantId()); this.type = job.getType(); this.key = job.getKey(); + this.description = job.getDescription(); this.status = job.getStatus(); this.configuration = toJson(job.getConfiguration()); this.result = toJson(job.getResult()); @@ -87,6 +91,7 @@ public class JobEntity extends BaseSqlEntity { job.setTenantId(getTenantId(tenantId)); job.setType(type); job.setKey(key); + job.setDescription(description); job.setStatus(status); job.setConfiguration(fromJson(configuration, JobConfiguration.class)); job.setResult(fromJson(result, JobResult.class)); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java index 273d8ef78a..9fff4d06b7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java @@ -23,14 +23,22 @@ import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.query.Param; import org.springframework.stereotype.Repository; +import org.thingsboard.server.common.data.job.JobStatus; +import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.dao.model.sql.JobEntity; +import java.util.List; import java.util.UUID; @Repository public interface JobRepository extends JpaRepository { - Page findByTenantId(UUID tenantId, Pageable pageable); + @Query("SELECT j FROM JobEntity j WHERE j.tenantId = :tenantId " + + "AND (:searchText IS NULL OR ilike(j.key, concat('%', :searchText, '%')) = true " + + "OR ilike(j.description, concat('%', :searchText, '%')) = true)") + Page findByTenantIdAndSearchText(@Param("tenantId") UUID tenantId, + @Param("searchText") String searchText, + Pageable pageable); @Modifying @Transactional @@ -45,7 +53,8 @@ public interface JobRepository extends JpaRepository { RETURNING ((result->>'successfulCount')::int + :count) + (result->>'failedCount')::int = (result->>'totalCount')::int """, nativeQuery = true) - boolean reportTaskSuccess(@Param("jobId") UUID jobId, @Param("count") int count); + boolean reportTaskSuccess(@Param("jobId") UUID jobId, + @Param("count") int count); @Modifying @Transactional @@ -64,6 +73,12 @@ public interface JobRepository extends JpaRepository { RETURNING ((result->>'failedCount')::int + 1) + (result->>'successfulCount')::int = (result->>'totalCount')::int """, nativeQuery = true) - boolean reportTaskFailure(@Param("jobId") UUID jobId, @Param("taskKey") String taskKey, @Param("error") String error); + boolean reportTaskFailure(@Param("jobId") UUID jobId, + @Param("taskKey") String taskKey, + @Param("error") String error); + + boolean existsByKeyAndStatusIn(String key, List statuses); + + boolean existsByTenantIdAndTypeAndStatusIn(UUID tenantId, JobType type, List statuses); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java index d6286e2d77..b8b36dbc23 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.sql.task; +import com.google.common.base.Strings; import lombok.RequiredArgsConstructor; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; @@ -22,6 +23,8 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobStatus; +import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.DaoUtil; @@ -30,6 +33,7 @@ import org.thingsboard.server.dao.sql.JpaAbstractDao; import org.thingsboard.server.dao.task.JobDao; import org.thingsboard.server.dao.util.SqlDao; +import java.util.Arrays; import java.util.UUID; @Component @@ -41,7 +45,7 @@ public class JpaJobDao extends JpaAbstractDao implements JobDao @Override public PageData findByTenantId(TenantId tenantId, PageLink pageLink) { - return DaoUtil.toPageData(jobRepository.findByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink))); + return DaoUtil.toPageData(jobRepository.findByTenantIdAndSearchText(tenantId.getId(), Strings.emptyToNull(pageLink.getTextSearch()), DaoUtil.toPageable(pageLink))); } @Override @@ -54,6 +58,16 @@ public class JpaJobDao extends JpaAbstractDao implements JobDao return jobRepository.reportTaskFailure(jobId.getId(), taskKey, error); } + @Override + public boolean existsByKeyAndStatusOneOf(String key, JobStatus... statuses) { + return jobRepository.existsByKeyAndStatusIn(key, Arrays.stream(statuses).toList()); + } + + @Override + public boolean existsByTenantIdAndTypeAndStatusOneOf(TenantId tenantId, JobType type, JobStatus... statuses) { + return jobRepository.existsByTenantIdAndTypeAndStatusIn(tenantId.getId(), type, Arrays.stream(statuses).toList()); + } + @Override public EntityType getEntityType() { return EntityType.JOB; diff --git a/dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java b/dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java index aa9e48600c..12d81c1c2c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java @@ -22,13 +22,14 @@ import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.job.Job; import org.thingsboard.server.common.data.job.JobResult; +import org.thingsboard.server.common.data.job.JobStats; import org.thingsboard.server.common.data.job.JobStatus; import org.thingsboard.server.common.data.job.TaskResult; import org.thingsboard.server.common.data.job.TaskResult.TaskFailure; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; - -import java.util.List; +import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.dao.service.DataValidator; @Service @RequiredArgsConstructor @@ -36,14 +37,21 @@ import java.util.List; public class DefaultJobService implements JobService { private final JobDao jobDao; + private final JobValidator validator = new JobValidator(); @Override public Job createJob(TenantId tenantId, Job job) { + validator.validate(job, Job::getTenantId); return jobDao.save(tenantId, job); } @Override - public void reportTaskResults(JobId jobId, List results) { + public Job findJobById(TenantId tenantId, JobId jobId) { + return jobDao.findById(tenantId, jobId.getId()); + } + + @Override + public void processStats(JobId jobId, JobStats jobStats) { Job job = jobDao.findById(TenantId.SYS_TENANT_ID, jobId.getId()); switch (job.getStatus()) { case PENDING -> { @@ -56,7 +64,11 @@ public class DefaultJobService implements JobService { } JobResult jobResult = job.getResult(); - for (TaskResult taskResult : results) { + if (jobStats.getTotalTasksCount() != null) { + jobResult.setTotalCount(jobStats.getTotalTasksCount()); + } + + for (TaskResult taskResult : jobStats.getTaskResults()) { if (taskResult.isSuccess()) { jobResult.setSuccessfulCount(jobResult.getSuccessfulCount() + 1); } else { @@ -67,7 +79,7 @@ public class DefaultJobService implements JobService { } } - if (jobResult.getSuccessfulCount() + jobResult.getFailedCount() >= jobResult.getTotalCount()) { + if (jobResult.getTotalCount() != null && jobResult.getSuccessfulCount() + jobResult.getFailedCount() >= jobResult.getTotalCount()) { if (jobResult.getFailures().isEmpty()) { job.setStatus(JobStatus.COMPLETED); } else { @@ -83,4 +95,22 @@ public class DefaultJobService implements JobService { return jobDao.findByTenantId(tenantId, pageLink); } + // todo: cancellation, reprocessing + + public class JobValidator extends DataValidator { + + @Override + protected void validateCreate(TenantId tenantId, Job job) { + if (jobDao.existsByTenantIdAndTypeAndStatusOneOf(tenantId, job.getType(), JobStatus.PENDING, JobStatus.RUNNING)) { + throw new DataValidationException("Job of this type is already running"); + } + } + + @Override + protected Job validateUpdate(TenantId tenantId, Job job) { + throw new IllegalArgumentException("Job can't be updated externally"); + } + + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java b/dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java index 6a5b002aea..5c3c5b977c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java @@ -18,6 +18,8 @@ package org.thingsboard.server.dao.task; import org.thingsboard.server.common.data.id.JobId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobStatus; +import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.Dao; @@ -30,4 +32,8 @@ public interface JobDao extends Dao { boolean reportTaskFailure(JobId jobId, String taskKey, String error); + boolean existsByKeyAndStatusOneOf(String key, JobStatus... statuses); + + boolean existsByTenantIdAndTypeAndStatusOneOf(TenantId tenantId, JobType type, JobStatus... statuses); + } diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 7fa31da5fb..5afc9398d2 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -955,6 +955,7 @@ CREATE TABLE IF NOT EXISTS job ( tenant_id uuid NOT NULL, type varchar NOT NULL, key varchar NOT NULL, + description varchar NOT NULL, status varchar NOT NULL, configuration varchar(1000) NOT NULL, result jsonb