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 new file mode 100644 index 0000000000..3b63d3736f --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/CfReprocessingJobProcessor.java @@ -0,0 +1,83 @@ +/** + * 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.EntityType; +import org.thingsboard.server.common.data.ProfileEntityIdInfo; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.job.CfReprocessingJobConfiguration; +import org.thingsboard.server.common.data.job.CfReprocessingTask; +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 org.thingsboard.server.common.data.page.PageDataIterable; +import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.device.DeviceService; + +import java.util.function.Consumer; + +@Component +@RequiredArgsConstructor +public class CfReprocessingJobProcessor extends JobProcessor { + + private final DeviceService deviceService; + private final AssetService assetService; + + @Override + public void process(Job job, Consumer taskConsumer) { + CfReprocessingJobConfiguration configuration = job.getConfiguration(); + + CalculatedField calculatedField = configuration.getCalculatedField(); + EntityId entityId = calculatedField.getEntityId(); + + if (entityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) { + taskConsumer.accept(createTask(job, configuration, entityId)); + } 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); + } else { + throw new IllegalArgumentException("Unsupported CF entity type " + entityId.getEntityType()); + } + entities.forEach(device -> { + taskConsumer.accept(createTask(job, configuration, device.getEntityId())); + }); + } + } + + private Task createTask(Job job, CfReprocessingJobConfiguration configuration, EntityId entityId) { + return CfReprocessingTask.builder() + .tenantId(job.getTenantId()) + .jobId(job.getId()) + .key(entityId.getEntityType().getNormalName() + " " + entityId.getId()) + .calculatedField(configuration.getCalculatedField()) + .entityId(entityId) + .startTs(configuration.getStartTs()) + .endTs(configuration.getEndTs()) + .build(); + } + + @Override + public JobType getType() { + return JobType.CF_REPROCESSING; + } + +} 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 new file mode 100644 index 0000000000..73367e2af5 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/DefaultJobManager.java @@ -0,0 +1,136 @@ +/** + * 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 jakarta.annotation.PreDestroy; +import lombok.SneakyThrows; +import lombok.extern.slf4j.Slf4j; +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.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.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; +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.util.AfterStartUp; +import org.thingsboard.server.queue.util.TbCoreComponent; + +import java.util.Arrays; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.function.Function; +import java.util.stream.Collectors; + +@TbCoreComponent +@Component +@Slf4j +public class DefaultJobManager implements JobManager { + + private final JobService jobService; + private final TbCoreQueueFactory queueFactory; + private final Map jobProcessors; + private final Map>> taskProducers; + private final QueueConsumerManager> taskResultConsumer; + private final ExecutorService consumerExecutor; + + public DefaultJobManager(JobService jobService, TbCoreQueueFactory queueFactory, List jobProcessors) { + this.jobService = jobService; + this.queueFactory = queueFactory; + 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) + .pollInterval(125) + .consumerCreator(queueFactory::createTaskResultConsumer) + .consumerExecutor(consumerExecutor) + .build(); + } + + @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) + public void afterStartUp() { + taskResultConsumer.subscribe(); + taskResultConsumer.launch(); + } + + @Override + public void submitJob(Job job) { + job = jobService.createJob(job.getTenantId(), job); + log.info("Submitting job: {}", job); + jobProcessors.get(job.getType()).process(job, this::submitTask); + } + + private void submitTask(Task task) { + log.info("Submitting task: {}", task); + TaskProto taskProto = TaskProto.newBuilder() + .setValue(JacksonUtil.toString(task)) + .build(); + + TbQueueProducer> producer = taskProducers.get(task.getJobType()); + TbProtoQueueMsg msg = new TbProtoQueueMsg<>(task.getTenantId().getId(), taskProto); // one job at a time for a given tenant + producer.send(TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(), msg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + log.trace("Submitted task: {}", task); + } + + @Override + public void onFailure(Throwable t) { + log.warn("Failed to submit task: {}", task, t); + } + }); + } + + @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) -> { + try { + log.info("[{}] Processing task results: {}", jobId, taskResults); + jobService.reportTaskResults(jobId, taskResults); + } catch (Exception e) { + log.warn("Failed to report task results for job {}: {}", jobId, taskResults, e); + } + }); + consumer.commit(); + + Thread.sleep(5000); + } + + @PreDestroy + private void destroy() { + taskResultConsumer.stop(); + consumerExecutor.shutdownNow(); + } + +} 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 new file mode 100644 index 0000000000..5db52ef448 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/JobManager.java @@ -0,0 +1,24 @@ +/** + * 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 org.thingsboard.server.common.data.job.Job; + +public interface JobManager { + + void 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 new file mode 100644 index 0000000000..cf15209616 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/JobProcessor.java @@ -0,0 +1,30 @@ +/** + * 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 org.thingsboard.server.common.data.job.Task; +import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobType; + +import java.util.function.Consumer; + +public abstract class JobProcessor { + + public abstract void process(Job job, Consumer taskConsumer); + + public abstract JobType getType(); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/job/task/CfReprocessingTaskProcessor.java b/application/src/main/java/org/thingsboard/server/service/job/task/CfReprocessingTaskProcessor.java new file mode 100644 index 0000000000..36899516f9 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/job/task/CfReprocessingTaskProcessor.java @@ -0,0 +1,57 @@ +/** + * 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 com.google.common.util.concurrent.SettableFuture; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Component; +import org.thingsboard.server.actors.calculatedField.CalculatedFieldReprocessingService; +import org.thingsboard.server.common.data.job.CfReprocessingTask; +import org.thingsboard.server.common.data.job.JobType; +import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.queue.task.TaskProcessor; + +import java.util.concurrent.TimeUnit; + +@Component +@RequiredArgsConstructor +public class CfReprocessingTaskProcessor extends TaskProcessor { + + private final CalculatedFieldReprocessingService cfReprocessingService; + + @Override + protected void process(CfReprocessingTask task) throws Exception { + SettableFuture future = SettableFuture.create(); + cfReprocessingService.reprocess(task, new TbCallback() { + @Override + public void onSuccess() { + future.set(null); + } + + @Override + public void onFailure(Throwable t) { + future.setException(t); + } + }); + future.get(1, TimeUnit.MINUTES); + } + + @Override + public JobType getJobType() { + return JobType.CF_REPROCESSING; + } + +} diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 418263076b..6670140b30 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1667,6 +1667,8 @@ queue: edqs-requests: "${TB_QUEUE_KAFKA_EDQS_REQUESTS_TOPIC_PROPERTIES:retention.ms:180000;segment.bytes:52428800;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}" # Kafka properties for EDQS state topic (infinite retention, compaction) edqs-state: "${TB_QUEUE_KAFKA_EDQS_STATE_TOPIC_PROPERTIES:retention.ms:-1;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1;cleanup.policy:compact}" + # Kafka properties for tasks topics + tasks: "${TB_QUEUE_KAFKA_TASKS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:52428800;retention.bytes:104857600;partitions:100;min.insync.replicas:1}" consumer-stats: # Prints lag between consumer group offset and last messages offset in Kafka topics enabled: "${TB_QUEUE_KAFKA_CONSUMER_STATS_ENABLED:true}" diff --git a/application/src/test/resources/logback-test.xml b/application/src/test/resources/logback-test.xml index 13c93da411..a0efcf52c1 100644 --- a/application/src/test/resources/logback-test.xml +++ b/application/src/test/resources/logback-test.xml @@ -9,7 +9,7 @@ - + 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 new file mode 100644 index 0000000000..17e7645605 --- /dev/null +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/task/JobService.java @@ -0,0 +1,35 @@ +/** + * 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.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.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); + + PageData findJobsByTenantId(TenantId tenantId, PageLink pageLink); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java index 93e754eb2c..31cc983168 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java @@ -63,7 +63,8 @@ public enum EntityType { MOBILE_APP(37), MOBILE_APP_BUNDLE(38), CALCULATED_FIELD(39), - CALCULATED_FIELD_LINK(40); + CALCULATED_FIELD_LINK(40), + JOB(41); @Getter private final int protoNumber; // Corresponds to EntityTypeProto diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java index f5dd4b12a0..dcf59a4ea4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java @@ -117,6 +117,8 @@ public class EntityIdFactory { return new CalculatedFieldId(uuid); case CALCULATED_FIELD_LINK: return new CalculatedFieldLinkId(uuid); + case JOB: + return new JobId(uuid); } throw new IllegalArgumentException("EntityType " + type + " is not supported!"); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java new file mode 100644 index 0000000000..e6688f0eb0 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/id/JobId.java @@ -0,0 +1,38 @@ +/** + * 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.id; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.swagger.v3.oas.annotations.media.Schema; +import org.thingsboard.server.common.data.EntityType; + +import java.util.UUID; + +public class JobId extends UUIDBased implements EntityId { + + @JsonCreator + public JobId(@JsonProperty("id") UUID id) { + super(id); + } + + @Schema(requiredMode = Schema.RequiredMode.REQUIRED, description = "string", example = "TASK", allowableValues = "TASK") + @Override + public EntityType getEntityType() { + return EntityType.JOB; + } + +} 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 new file mode 100644 index 0000000000..4f5a33dacd --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobConfiguration.java @@ -0,0 +1,39 @@ +/** + * 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 org.thingsboard.server.common.data.cf.CalculatedField; + +@Data +@AllArgsConstructor +@NoArgsConstructor +@Builder +public class CfReprocessingJobConfiguration implements JobConfiguration { + + private CalculatedField calculatedField; + private long startTs; + private long endTs; + + @Override + public JobType getType() { + return JobType.CF_REPROCESSING; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobResult.java new file mode 100644 index 0000000000..2d756f6d53 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingJobResult.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 CfReprocessingJobResult extends JobResult { + + @Override + public JobType getJobType() { + return JobType.CF_REPROCESSING; + } + +} 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 new file mode 100644 index 0000000000..0e15b9473f --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/CfReprocessingTask.java @@ -0,0 +1,53 @@ +/** + * 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.EqualsAndHashCode; +import lombok.NoArgsConstructor; +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) +public class CfReprocessingTask extends Task { + + private CalculatedField calculatedField; + private EntityId entityId; + 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/Job.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java new file mode 100644 index 0000000000..e0161cf98e --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/Job.java @@ -0,0 +1,56 @@ +/** + * 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.Builder; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.BaseData; +import org.thingsboard.server.common.data.HasTenantId; +import org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.id.TenantId; + +@Data +@NoArgsConstructor +@EqualsAndHashCode(callSuper = true) +public class Job extends BaseData implements HasTenantId { + + private TenantId tenantId; + private JobType type; + private String key; + private JobStatus status; + private JobConfiguration configuration; + private JobResult result; + + @Builder + public Job(TenantId tenantId, JobType type, String key, JobConfiguration configuration) { + this.tenantId = tenantId; + this.type = type; + this.key = key; + this.configuration = configuration; + this.status = JobStatus.PENDING; + this.result = switch (type) { + case CF_REPROCESSING -> new CfReprocessingJobResult(); + }; + } + + @SuppressWarnings("unchecked") + public C getConfiguration() { + return (C) configuration; + } + +} 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 new file mode 100644 index 0000000000..8899206c49 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobConfiguration.java @@ -0,0 +1,32 @@ +/** + * 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 com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonSubTypes.Type; +import com.fasterxml.jackson.annotation.JsonTypeInfo; + +@JsonIgnoreProperties(ignoreUnknown = true) +@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type") +@JsonSubTypes({ + @Type(name = "CF_REPROCESSING", value = CfReprocessingJobConfiguration.class), +}) +public interface JobConfiguration { + + JobType getType(); + +} 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 new file mode 100644 index 0000000000..406eb0f65e --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobResult.java @@ -0,0 +1,41 @@ +/** + * 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 com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; +import lombok.Data; + +import java.util.HashMap; +import java.util.Map; + +@JsonIgnoreProperties(ignoreUnknown = true) +@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType") +@JsonSubTypes({ + @JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingJobResult.class), +}) +@Data +public abstract class JobResult { + + private int successfulCount; + private int failedCount; + private int totalCount; + private Map failures = new HashMap<>(); + + public abstract JobType getJobType(); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/job/JobStatus.java b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobStatus.java new file mode 100644 index 0000000000..026e19c5b2 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobStatus.java @@ -0,0 +1,24 @@ +/** + * 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 enum JobStatus { + PENDING, + RUNNING, + COMPLETED, + FAILED, + CANCELLED +} 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 new file mode 100644 index 0000000000..60cac8173f --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/JobType.java @@ -0,0 +1,26 @@ +/** + * 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 enum JobType { + + CF_REPROCESSING; + + 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 new file mode 100644 index 0000000000..afeaeba393 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/Task.java @@ -0,0 +1,50 @@ +/** + * 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 com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; +import lombok.Data; +import org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.id.TenantId; + +@Data +@JsonIgnoreProperties(ignoreUnknown = true) +@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "jobType") +@JsonSubTypes({ + @JsonSubTypes.Type(name = "CF_REPROCESSING", value = CfReprocessingTask.class), +}) +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; + } + + public Task() { + } + + private int attempt = 0; + + public abstract JobType getJobType(); + +} 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 new file mode 100644 index 0000000000..bfbef46180 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/job/TaskResult.java @@ -0,0 +1,45 @@ +/** + * 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 org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.id.TenantId; + +@Data +@AllArgsConstructor +@NoArgsConstructor +@Builder +public class TaskResult { + + private TenantId tenantId; + private JobId jobId; + private boolean success; + private TaskFailure failure; + + @Data + @AllArgsConstructor + @NoArgsConstructor + @Builder + public static class TaskFailure { + private String error; + private Task task; + } + +} diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 938a1692ae..de03b11b6a 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -1846,3 +1846,11 @@ message EdqsRequestMsg { message EdqsResponseMsg { string value = 1; } + +message TaskProto { + string value = 1; // fixme: TMP, make more efficient +} + +message TaskResultProto { + string value = 1; // fixme: TMP, make more efficient +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaTopicConfigs.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaTopicConfigs.java index aebda5a5bc..5d5834d20a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaTopicConfigs.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaTopicConfigs.java @@ -62,6 +62,8 @@ public class TbKafkaTopicConfigs { private String edqsRequestsProperties; @Value("${queue.kafka.topic-properties.edqs-state:}") private String edqsStateProperties; + @Value("${queue.kafka.topic-properties.tasks:}") + private String tasksProperties; @Getter private Map coreConfigs; @@ -99,6 +101,8 @@ public class TbKafkaTopicConfigs { private Map edqsRequestsConfigs; @Getter private Map edqsStateConfigs; + @Getter + private Map tasksConfigs; @PostConstruct private void init() { @@ -122,6 +126,7 @@ public class TbKafkaTopicConfigs { edqsEventsConfigs = PropertyUtils.getProps(edqsEventsProperties); edqsRequestsConfigs = PropertyUtils.getProps(edqsRequestsProperties); edqsStateConfigs = PropertyUtils.getProps(edqsStateProperties); + tasksConfigs = PropertyUtils.getProps(tasksProperties); } } 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 085d04f28c..c8866e52b4 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 @@ -20,6 +20,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.job.JobType; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -258,9 +259,28 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE .build(); } + @Override + public TbQueueProducer> createTaskProducer(JobType jobType) { + return new InMemoryTbQueueProducer<>(storage, jobType.getTasksTopic()); + } + + @Override + public TbQueueConsumer> createTaskConsumer(JobType jobType) { + return new InMemoryTbQueueConsumer<>(storage, jobType.getTasksTopic()); + } + + @Override + public TbQueueProducer> createTaskResultProducer() { + return new InMemoryTbQueueProducer<>(storage, "tasks.results"); + } + + @Override + public TbQueueConsumer> createTaskResultConsumer() { + return new InMemoryTbQueueConsumer<>(storage, "tasks.results"); + } + @Scheduled(fixedRateString = "${queue.in_memory.stats.print-interval-ms:60000}") private void printInMemoryStats() { storage.printStats(); } - } 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 2269004f90..809099aa8b 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 @@ -23,12 +23,15 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.DataConstants; 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.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.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; +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; @@ -110,6 +113,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi private final TbQueueAdmin cfStateAdmin; private final TbQueueAdmin edqsEventsAdmin; private final TbKafkaAdmin edqsRequestsAdmin; + private final TbQueueAdmin tasksAdmin; private final AtomicLong consumerCount = new AtomicLong(); private final AtomicLong edgeConsumerCount = new AtomicLong(); @@ -158,6 +162,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi this.cfStateAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldStateConfigs()); this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs()); this.edqsRequestsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsRequestsConfigs()); + this.tasksAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTasksConfigs()); } @Override @@ -641,6 +646,52 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi .build(); } + @Override + public TbQueueProducer> createTaskProducer(JobType jobType) { + return TbKafkaProducerTemplate.>builder() + .clientId(jobType.name().toLowerCase() + "-task-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName(jobType.getTasksTopic())) + .settings(kafkaSettings) + .admin(tasksAdmin) + .build(); + } + + @Override + public TbQueueConsumer> createTaskConsumer(JobType jobType) { + return TbKafkaConsumerTemplate.>builder() + .settings(kafkaSettings) + .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> createTaskResultProducer() { + return TbKafkaProducerTemplate.>builder() + .clientId("task-result-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName("tasks.results")) + .settings(kafkaSettings) + .admin(tasksAdmin) + .build(); + } + + @Override + public TbQueueConsumer> createTaskResultConsumer() { + 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())) + .admin(tasksAdmin) + .statsService(consumerStatsService) + .build(); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { 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 ea7c56f0aa..71a9669ba4 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 @@ -22,9 +22,12 @@ import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; 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.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.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -105,6 +108,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { private final TbQueueAdmin cfAdmin; private final TbQueueAdmin edqsEventsAdmin; private final TbKafkaAdmin edqsRequestsAdmin; + private final TbQueueAdmin tasksAdmin; private final AtomicLong consumerCount = new AtomicLong(); private final AtomicLong edgeConsumerCount = new AtomicLong(); @@ -153,6 +157,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs()); this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs()); this.edqsRequestsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsRequestsConfigs()); + this.tasksAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTasksConfigs()); } @Override @@ -520,6 +525,29 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory { .build(); } + @Override + public TbQueueProducer> createTaskProducer(JobType jobType) { + return TbKafkaProducerTemplate.>builder() + .clientId(jobType.name().toLowerCase() + "-task-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName(jobType.getTasksTopic())) + .settings(kafkaSettings) + .admin(tasksAdmin) + .build(); + } + + @Override + public TbQueueConsumer> createTaskResultConsumer() { + 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())) + .admin(tasksAdmin) + .statsService(consumerStatsService) + .build(); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { 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 b4884ae72c..d0ac20ae06 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 @@ -22,12 +22,15 @@ import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.job.JobType; 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.TaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; @@ -96,6 +99,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { private final TbQueueAdmin cfAdmin; private final TbQueueAdmin cfStateAdmin; private final TbQueueAdmin edqsEventsAdmin; + private final TbQueueAdmin tasksAdmin; private final AtomicLong consumerCount = new AtomicLong(); public KafkaTbRuleEngineQueueFactory(TopicService topicService, TbKafkaSettings kafkaSettings, @@ -133,6 +137,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs()); this.cfStateAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldStateConfigs()); this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdqsEventsConfigs()); + this.tasksAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getTasksConfigs()); } @Override @@ -414,6 +419,29 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { throw new UnsupportedOperationException(); } + @Override + public TbQueueConsumer> createTaskConsumer(JobType jobType) { + return TbKafkaConsumerTemplate.>builder() + .settings(kafkaSettings) + .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> createTaskResultProducer() { + return TbKafkaProducerTemplate.>builder() + .clientId("task-result-producer-" + serviceInfoProvider.getServiceId()) + .defaultTopic(topicService.buildTopicName("tasks.results")) + .settings(kafkaSettings) + .admin(tasksAdmin) + .build(); + } + @PreDestroy private void destroy() { if (coreAdmin != null) { @@ -438,4 +466,5 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { cfAdmin.destroy(); } } + } 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 new file mode 100644 index 0000000000..10b84f0f65 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TaskProcessorQueueFactory.java @@ -0,0 +1,31 @@ +/** + * 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.provider; + +import org.thingsboard.server.common.data.job.JobType; +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; + +public interface TaskProcessorQueueFactory { + + TbQueueConsumer> createTaskConsumer(JobType jobType); + + TbQueueProducer> createTaskResultProducer(); + +} 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 037d1f2087..409348a12a 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 @@ -17,7 +17,10 @@ package org.thingsboard.server.queue.provider; 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.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; @@ -165,4 +168,8 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, Hous TbQueueProducer> createToCalculatedFieldNotificationMsgProducer(); + TbQueueProducer> createTaskProducer(JobType jobType); + + TbQueueConsumer> createTaskResultConsumer(); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java index 18bb6db14a..83c467c992 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java @@ -41,7 +41,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 TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory { +public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory, HousekeeperClientQueueFactory, EdqsClientQueueFactory, TaskProcessorQueueFactory { /** * Used to push messages to instances of TB Transport Service 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 new file mode 100644 index 0000000000..7123eaf1e0 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/task/TaskProcessor.java @@ -0,0 +1,139 @@ +/** + * 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 jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.common.util.ThingsBoardThreadFactory; +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; +import org.thingsboard.server.queue.util.AfterStartUp; + +import java.util.List; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +@Slf4j +public abstract class TaskProcessor { + + @Autowired + private TaskProcessorQueueFactory queueFactory; + + 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") + .msgPackProcessor(this::processMsgs) + .pollInterval(125) + .consumerCreator(() -> queueFactory.createTaskConsumer(getJobType())) + .consumerExecutor(consumerExecutor) + .build(); + taskResultProducer = queueFactory.createTaskResultProducer(); + } + + @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) + public void afterStartUp() { + taskConsumer.subscribe(); + taskConsumer.launch(); + } + + @PreDestroy + public void destroy() { + taskConsumer.stop(); + consumerExecutor.shutdownNow(); + } + + private void processMsgs(List> msgs, TbQueueConsumer> consumer) { + for (TbProtoQueueMsg msg : msgs) { + TaskProto taskProto = msg.getValue(); + Task task = JacksonUtil.fromString(taskProto.getValue(), Task.class); + processTask((T) task); + } + consumer.commit(); + } + + private void processTask(T task) { + task.setAttempt(task.getAttempt() + 1); + log.info("Processing task: {}", task); + try { + process(task); + reportSuccess(task); + } catch (Exception e) { + log.error("Failed to process task (attempt {}): {}", task.getAttempt(), task, e); + if (task.getAttempt() < 3) { + processTask(task); + } else { + reportFailure(task, e); + } + } + } + + private void reportSuccess(Task task) { + TaskResult result = TaskResult.builder() + .tenantId(task.getTenantId()) + .jobId(task.getJobId()) + .success(true) + .build(); + reportResult(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); + } + + protected abstract void process(T task) throws Exception; + + public abstract JobType getJobType(); + +} 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 148908d063..1d504b0ab2 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 @@ -739,6 +739,16 @@ public class ModelConstants { public static final String CALCULATED_FIELD_LINK_ENTITY_ID = ENTITY_ID_COLUMN; public static final String CALCULATED_FIELD_LINK_CALCULATED_FIELD_ID = "calculated_field_id"; + /** + * Tasks constants. + */ + 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_STATUS_PROPERTY = "status"; + public static final String JOB_CONFIGURATION_PROPERTY = "configuration"; + public static final String JOB_RESULT_PROPERTY = "result"; + protected static final String[] NONE_AGGREGATION_COLUMNS = new String[]{LONG_VALUE_COLUMN, DOUBLE_VALUE_COLUMN, BOOLEAN_VALUE_COLUMN, STRING_VALUE_COLUMN, JSON_VALUE_COLUMN, KEY_COLUMN, TS_COLUMN}; protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN), count(JSON_VALUE_COLUMN), max(TS_COLUMN)}; 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 new file mode 100644 index 0000000000..6e8b3e7958 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/JobEntity.java @@ -0,0 +1,96 @@ +/** + * 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.dao.model.sql; + +import com.fasterxml.jackson.databind.JsonNode; +import jakarta.persistence.Column; +import jakarta.persistence.Convert; +import jakarta.persistence.Entity; +import jakarta.persistence.EnumType; +import jakarta.persistence.Enumerated; +import jakarta.persistence.Table; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.NoArgsConstructor; +import org.hibernate.annotations.JdbcType; +import org.hibernate.dialect.PostgreSQLJsonPGObjectJsonbType; +import org.thingsboard.server.common.data.id.JobId; +import org.thingsboard.server.common.data.job.Job; +import org.thingsboard.server.common.data.job.JobConfiguration; +import org.thingsboard.server.common.data.job.JobResult; +import org.thingsboard.server.common.data.job.JobStatus; +import org.thingsboard.server.common.data.job.JobType; +import org.thingsboard.server.dao.model.BaseSqlEntity; +import org.thingsboard.server.dao.model.ModelConstants; +import org.thingsboard.server.dao.util.mapping.JsonConverter; + +import java.util.UUID; + +@Data +@EqualsAndHashCode(callSuper = true) +@NoArgsConstructor +@Entity +@Table(name = ModelConstants.JOB_TABLE_NAME) +public class JobEntity extends BaseSqlEntity { + + @Column(name = ModelConstants.TENANT_ID_PROPERTY, nullable = false) + private UUID tenantId; + + @Enumerated(EnumType.STRING) + @Column(name = ModelConstants.JOB_TYPE_PROPERTY, nullable = false) + private JobType type; + + @Column(name = ModelConstants.JOB_KEY_PROPERTY, nullable = false) + private String key; + + @Enumerated(EnumType.STRING) + @Column(name = ModelConstants.JOB_STATUS_PROPERTY, nullable = false) + private JobStatus status; + + @Convert(converter = JsonConverter.class) + @Column(name = ModelConstants.JOB_CONFIGURATION_PROPERTY, nullable = false) + private JsonNode configuration; + + @Convert(converter = JsonConverter.class) + @JdbcType(PostgreSQLJsonPGObjectJsonbType.class) + @Column(name = ModelConstants.JOB_RESULT_PROPERTY) + private JsonNode result; + + public JobEntity(Job job) { + super(job); + this.tenantId = getTenantUuid(job.getTenantId()); + this.type = job.getType(); + this.key = job.getKey(); + this.status = job.getStatus(); + this.configuration = toJson(job.getConfiguration()); + this.result = toJson(job.getResult()); + } + + @Override + public Job toData() { + Job job = new Job(); + job.setId(new JobId(id)); + job.setCreatedTime(createdTime); + job.setTenantId(getTenantId(tenantId)); + job.setType(type); + job.setKey(key); + job.setStatus(status); + job.setConfiguration(fromJson(configuration, JobConfiguration.class)); + job.setResult(fromJson(result, JobResult.class)); + return job; + } + +} 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 new file mode 100644 index 0000000000..273d8ef78a --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/task/JobRepository.java @@ -0,0 +1,69 @@ +/** + * 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.dao.sql.task; + +import jakarta.transaction.Transactional; +import org.springframework.data.domain.Page; +import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.JpaRepository; +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.dao.model.sql.JobEntity; + +import java.util.UUID; + +@Repository +public interface JobRepository extends JpaRepository { + + Page findByTenantId(UUID tenantId, Pageable pageable); + + @Modifying + @Transactional + @Query(value = """ + UPDATE job + SET result = jsonb_set( + result, + '{successfulCount}', + to_jsonb((result->>'successfulCount')::int + :count) + ) + WHERE id = :jobId + RETURNING ((result->>'successfulCount')::int + :count) + + (result->>'failedCount')::int = (result->>'totalCount')::int + """, nativeQuery = true) + boolean reportTaskSuccess(@Param("jobId") UUID jobId, @Param("count") int count); + + @Modifying + @Transactional + @Query(value = """ + UPDATE job + SET result = jsonb_set( + jsonb_set( + result, + '{failedCount}', + to_jsonb((result->>'failedCount')::int + 1) + ), + ARRAY['failures', :taskKey], + to_jsonb(:error) + ) + WHERE id = :jobId + 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); + +} 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 new file mode 100644 index 0000000000..d6286e2d77 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/task/JpaJobDao.java @@ -0,0 +1,72 @@ +/** + * 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.dao.sql.task; + +import lombok.RequiredArgsConstructor; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.stereotype.Component; +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.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.dao.DaoUtil; +import org.thingsboard.server.dao.model.sql.JobEntity; +import org.thingsboard.server.dao.sql.JpaAbstractDao; +import org.thingsboard.server.dao.task.JobDao; +import org.thingsboard.server.dao.util.SqlDao; + +import java.util.UUID; + +@Component +@SqlDao +@RequiredArgsConstructor +public class JpaJobDao extends JpaAbstractDao implements JobDao { + + private final JobRepository jobRepository; + + @Override + public PageData findByTenantId(TenantId tenantId, PageLink pageLink) { + return DaoUtil.toPageData(jobRepository.findByTenantId(tenantId.getId(), DaoUtil.toPageable(pageLink))); + } + + @Override + public boolean reportTaskSuccess(JobId jobId, int tasksCount) { + return jobRepository.reportTaskSuccess(jobId.getId(), tasksCount); + } + + @Override + public boolean reportTaskFailure(JobId jobId, String taskKey, String error) { + return jobRepository.reportTaskFailure(jobId.getId(), taskKey, error); + } + + @Override + public EntityType getEntityType() { + return EntityType.JOB; + } + + @Override + protected Class getEntityClass() { + return JobEntity.class; + } + + @Override + protected JpaRepository getRepository() { + return jobRepository; + } + +} 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 new file mode 100644 index 0000000000..aa9e48600c --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/task/DefaultJobService.java @@ -0,0 +1,86 @@ +/** + * 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.dao.task; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +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.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; + +@Service +@RequiredArgsConstructor +@Slf4j +public class DefaultJobService implements JobService { + + private final JobDao jobDao; + + @Override + public Job createJob(TenantId tenantId, Job job) { + return jobDao.save(tenantId, job); + } + + @Override + public void reportTaskResults(JobId jobId, List results) { + Job job = jobDao.findById(TenantId.SYS_TENANT_ID, jobId.getId()); + switch (job.getStatus()) { + case PENDING -> { + job.setStatus(JobStatus.RUNNING); + } + case CANCELLED, COMPLETED, FAILED -> { + // got some stale stats + return; + } + } + + JobResult jobResult = job.getResult(); + for (TaskResult taskResult : results) { + if (taskResult.isSuccess()) { + jobResult.setSuccessfulCount(jobResult.getSuccessfulCount() + 1); + } else { + TaskFailure failure = taskResult.getFailure(); + String key = failure.getTask().getKey(); + jobResult.setFailedCount(jobResult.getFailedCount() + 1); + jobResult.getFailures().put(key, failure.getError()); + } + } + + if (jobResult.getSuccessfulCount() + jobResult.getFailedCount() >= jobResult.getTotalCount()) { + if (jobResult.getFailures().isEmpty()) { + job.setStatus(JobStatus.COMPLETED); + } else { + job.setStatus(JobStatus.FAILED); + } + } + log.info("Saving job {}", job); + jobDao.save(TenantId.SYS_TENANT_ID, job); + } + + @Override + public PageData findJobsByTenantId(TenantId tenantId, PageLink pageLink) { + return jobDao.findByTenantId(tenantId, pageLink); + } + +} 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 new file mode 100644 index 0000000000..6a5b002aea --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/task/JobDao.java @@ -0,0 +1,33 @@ +/** + * 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.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.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.dao.Dao; + +public interface JobDao extends Dao { + + PageData findByTenantId(TenantId tenantId, PageLink pageLink); + + boolean reportTaskSuccess(JobId jobId, int tasksCount); + + boolean reportTaskFailure(JobId jobId, String taskKey, String error); + +} diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index b425550e7e..7fa31da5fb 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -948,3 +948,14 @@ CREATE TABLE IF NOT EXISTS cf_debug_event ( e_result varchar, e_error varchar ) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS job ( + id uuid NOT NULL CONSTRAINT job_pkey PRIMARY KEY, + created_time bigint NOT NULL, + tenant_id uuid NOT NULL, + type varchar NOT NULL, + key varchar NOT NULL, + status varchar NOT NULL, + configuration varchar(1000) NOT NULL, + result jsonb +);