diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java index 2f06b4d6e4..9fd94b4988 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java @@ -105,7 +105,7 @@ import java.util.function.BiConsumer; @Slf4j @Data -public abstract class AbstractEdgeGrpcSession> implements EdgeGrpcSession, Closeable { +public abstract class AbstractEdgeGrpcSession implements EdgeGrpcSession, Closeable { private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; private static final String QUEUE_START_SEQ_ID_ATTR_KEY = "queueStartSeqId"; @@ -116,7 +116,7 @@ public abstract class AbstractEdgeGrpcSession highPriorityQueue = new ConcurrentLinkedQueue<>(); protected UUID sessionId; - private BiConsumer sessionOpenListener; + private BiConsumer sessionOpenListener; private BiConsumer sessionCloseListener; private final EdgeSessionState sessionState = new EdgeSessionState(); @@ -146,7 +146,7 @@ public abstract class AbstractEdgeGrpcSession outputStream, - BiConsumer sessionOpenListener, + BiConsumer sessionOpenListener, BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { @@ -299,11 +299,6 @@ public abstract class AbstractEdgeGrpcSession downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); - for (DownlinkMsg downlinkMsg : downlinkMsgsPack) { - if (downlinkMsg.getEntityDataCount() > 0) { - System.out.println("downlink = " + downlinkMsg); - } - } Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { @Override public void onSuccess(@Nullable Boolean isInterrupted) { @@ -351,7 +346,7 @@ public abstract class AbstractEdgeGrpcSession> sessions = new ConcurrentHashMap<>(); + private final ConcurrentMap sessions = new ConcurrentHashMap<>(); private final ConcurrentMap sessionNewEventsLocks = new ConcurrentHashMap<>(); private final Map sessionNewEvents = new HashMap<>(); private final ConcurrentMap> sessionEdgeEventChecks = new ConcurrentHashMap<>(); @@ -210,11 +210,11 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public StreamObserver handleMsgs(StreamObserver outputStream) { - AbstractEdgeGrpcSession session = createEdgeGrpcSession(outputStream); + AbstractEdgeGrpcSession session = createEdgeGrpcSession(outputStream); return session.getInputStream(); } - private AbstractEdgeGrpcSession createEdgeGrpcSession(StreamObserver outputStream) { + private AbstractEdgeGrpcSession createEdgeGrpcSession(StreamObserver outputStream) { return isKafkaSupported ? new KafkaEdgeGrpcSession(ctx, tbCoreQueueFactory, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession) @@ -250,7 +250,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void updateEdge(TenantId tenantId, Edge edge) { - AbstractEdgeGrpcSession session = sessions.get(edge.getId()); + AbstractEdgeGrpcSession session = sessions.get(edge.getId()); if (session != null && session.isConnected()) { log.debug("[{}] Updating configuration for edge [{}] [{}]", tenantId, edge.getName(), edge.getId()); session.onConfigurationUpdate(edge); @@ -261,7 +261,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i @Override public void deleteEdge(TenantId tenantId, EdgeId edgeId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); session.close(); @@ -278,7 +278,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEventUpdate [{}]", tenantId, edgeId.getId()); updateSessionEventsFlag(tenantId, edgeId); @@ -289,7 +289,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i TenantId tenantId = msg.getTenantId(); EdgeEvent edgeEvent = msg.getEdgeEvent(); EdgeId edgeId = edgeEvent.getEdgeId(); - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null && session.isConnected()) { log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId); session.addEventToHighPriorityQueue(edgeEvent); @@ -310,7 +310,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void onEdgeConnect(EdgeId edgeId, AbstractEdgeGrpcSession edgeGrpcSession) { + private void onEdgeConnect(EdgeId edgeId, AbstractEdgeGrpcSession edgeGrpcSession) { Edge edge = edgeGrpcSession.getEdge(); TenantId tenantId = edge.getTenantId(); log.info("[{}][{}] edge [{}] connected successfully.", tenantId, edgeGrpcSession.getSessionId(), edgeId); @@ -333,7 +333,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } private void startSyncProcess(TenantId tenantId, EdgeId edgeId, UUID requestId, String requestServiceId) { - AbstractEdgeGrpcSession session = sessions.get(edgeId); + AbstractEdgeGrpcSession session = sessions.get(edgeId); if (session != null) { if (!session.isSyncCompleted()) { clusterService.pushEdgeSyncResponseToCore(new FromEdgeSyncResponse(requestId, tenantId, edgeId, false, "Sync process is active at the moment"), requestServiceId); @@ -353,7 +353,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i ToEdgeSyncRequest request = new ToEdgeSyncRequest(UUID.randomUUID(), tenantId, edgeId, serviceInfoProvider.getServiceId()); UUID requestId = request.getId(); - AbstractEdgeGrpcSession session = sessions.get(request.getEdgeId()); + AbstractEdgeGrpcSession session = sessions.get(request.getEdgeId()); if (session != null && !session.isSyncCompleted()) { responseConsumer.accept(new FromEdgeSyncResponse(requestId, request.getTenantId(), request.getEdgeId(), false, "Sync process is active at the moment")); } else { @@ -387,7 +387,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void scheduleEdgeEventsCheck(AbstractEdgeGrpcSession session) { + private void scheduleEdgeEventsCheck(AbstractEdgeGrpcSession session) { EdgeId edgeId = session.getEdge().getId(); TenantId tenantId = session.getEdge().getTenantId(); if (sessions.containsKey(edgeId)) { @@ -434,14 +434,14 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } } - private void processEdgeEventMigrationIfNeeded(AbstractEdgeGrpcSession session, EdgeId edgeId) throws Exception { + private void processEdgeEventMigrationIfNeeded(AbstractEdgeGrpcSession session, EdgeId edgeId) throws Exception { boolean isMigrationProcessed = edgeEventsProcessed.getOrDefault(edgeId, Boolean.FALSE); if (!isMigrationProcessed) { - Boolean migrated = session.migrateEdgeEvents(false).get(); - if (Boolean.TRUE.equals(migrated)) { + Boolean eventsExist = session.migrateEdgeEvents().get(); + if (Boolean.TRUE.equals(eventsExist)) { sessionNewEvents.put(edgeId, true); scheduleEdgeEventsCheck(session); - } else if (Boolean.FALSE.equals(migrated)) { + } else if (Boolean.FALSE.equals(eventsExist)) { edgeEventsProcessed.put(edgeId, true); } else { scheduleEdgeEventsCheck(session); @@ -463,7 +463,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i private void onEdgeDisconnect(Edge edge, UUID sessionId) { EdgeId edgeId = edge.getId(); log.info("[{}][{}] edge disconnected!", edgeId, sessionId); - AbstractEdgeGrpcSession toRemove = sessions.get(edgeId); + AbstractEdgeGrpcSession toRemove = sessions.get(edgeId); if (toRemove.getSessionId().equals(sessionId)) { toRemove = sessions.remove(edgeId); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); @@ -473,9 +473,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i } finally { newEventLock.unlock(); } - if (isKafkaSupported) { - ((KafkaEdgeGrpcSession) toRemove).stopConsumer(); - } + toRemove.destroy(); TenantId tenantId = toRemove.getEdge().getTenantId(); save(tenantId, edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); long lastDisconnectTs = System.currentTimeMillis(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 3db3a72288..880712d6b2 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -26,7 +26,9 @@ public interface EdgeGrpcSession { boolean isConnected(); - ListenableFuture migrateEdgeEvents(boolean isMigrationProcessed) throws Exception; + void destroy(); + + ListenableFuture migrateEdgeEvents() throws Exception; ListenableFuture processEdgeEvents() throws Exception; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java index 1126214e60..9b93238400 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java @@ -22,20 +22,10 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; -import org.thingsboard.server.cache.limits.RateLimitService; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.EdgeEvent; -import org.thingsboard.server.common.data.id.EdgeId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.limit.LimitedApi; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.util.ProtoUtils; -import org.thingsboard.server.dao.edge.EdgeEventDao; -import org.thingsboard.server.dao.edge.EdgeEventService; -import org.thingsboard.server.dao.service.DataValidator; +import org.thingsboard.server.dao.edge.BaseEdgeEventService; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.TopicService; @@ -47,25 +37,16 @@ import java.util.UUID; @Service @RequiredArgsConstructor @ConditionalOnExpression("'${queue.type:null}'=='kafka'") -public class KafkaEdgeEventService implements EdgeEventService { +public class KafkaEdgeEventService extends BaseEdgeEventService { - private final RateLimitService rateLimitService; - private final DataValidator edgeEventValidator; - @Lazy - private final TbQueueProducerProvider producerProvider; @Lazy private final TopicService topicService; - private final EdgeEventDao edgeEventDao; + @Lazy + private final TbQueueProducerProvider producerProvider; @Override public ListenableFuture saveAsync(EdgeEvent edgeEvent) { - if (!rateLimitService.checkRateLimit(LimitedApi.EDGE_EVENTS, edgeEvent.getTenantId())) { - throw new TbRateLimitsException(EntityType.TENANT); - } - if (!rateLimitService.checkRateLimit(LimitedApi.EDGE_EVENTS_PER_EDGE, edgeEvent.getTenantId(), edgeEvent.getEdgeId())) { - throw new TbRateLimitsException(EntityType.EDGE); - } - edgeEventValidator.validate(edgeEvent, EdgeEvent::getTenantId); + validateEdgeEvent(edgeEvent); TopicPartitionInfo tpi = topicService.getEdgeEventNotificationsTopic(edgeEvent.getTenantId(), edgeEvent.getEdgeId()); ToEdgeEventNotificationMsg msg = ToEdgeEventNotificationMsg.newBuilder().setEdgeEventMsg(ProtoUtils.toProto(edgeEvent)).build(); @@ -74,16 +55,4 @@ public class KafkaEdgeEventService implements EdgeEventService { return Futures.immediateFuture(null); } - @Override - public PageData findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) { - // To support fetching edge events on connect from postgres if there are any: - return edgeEventDao.findEdgeEvents(tenantId.getId(), edgeId, seqIdStart, seqIdEnd, pageLink); - } - - @Override - public void cleanupEvents(long ttl) { - // To delete deprecated events by ttl - edgeEventDao.cleanupEvents(ttl); - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java index 4fabfd5945..15a87ec3f1 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java @@ -18,7 +18,6 @@ package org.thingsboard.server.service.edge.rpc; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import io.grpc.stub.StreamObserver; -import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.edge.Edge; @@ -43,19 +42,18 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.function.BiConsumer; @Slf4j -public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { +public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { private final TbQueueConsumer> edgeEventsConsumer; private volatile boolean isHighPriorityProcessing; - private volatile boolean isConsumerInit; private QueueConsumerManager> consumer; private ExecutorService consumerExecutor; public KafkaEdgeGrpcSession(EdgeContextComponent ctx, TbCoreQueueFactory tbCoreQueueFactory, StreamObserver outputStream, - BiConsumer sessionOpenListener, + BiConsumer sessionOpenListener, BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); @@ -82,32 +80,23 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession migrateEdgeEvents(boolean isMigrationProcessed) throws Exception { - if (isMigrationProcessed) { - return Futures.immediateFuture(Boolean.FALSE); - } + public ListenableFuture migrateEdgeEvents() throws Exception { return super.processEdgeEvents(); } @Override public ListenableFuture processEdgeEvents() { - if (!isConsumerInit) { + if (consumer == null) { this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer")); this.consumer = QueueConsumerManager.>builder() .name("TB Edge events") @@ -117,7 +106,6 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession { +public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession { PostgresEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, - BiConsumer sessionOpenListener, + BiConsumer sessionOpenListener, BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); @@ -40,7 +40,10 @@ public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession migrateEdgeEvents(boolean isMigrationProcessed) { + public void destroy() {} + + @Override + public ListenableFuture migrateEdgeEvents() { return Futures.immediateFuture(Boolean.FALSE); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index e8e2ca7c3a..e448fbd937 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -51,6 +51,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.dao.edge.EdgeSynchronizationManager; +import org.thingsboard.server.dao.entity.EntityDaoRegistry; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.TbQueueCallback; @@ -77,6 +78,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected EdgeContextComponent edgeCtx; + @Autowired + protected EntityDaoRegistry entityDaoRegistry; + @Autowired protected EdgeSynchronizationManager edgeSynchronizationManager; @@ -315,17 +319,7 @@ public abstract class BaseEdgeProcessor { } protected boolean isEntityExists(TenantId tenantId, EntityId entityId) { - return switch (entityId.getEntityType()) { - case TENANT -> edgeCtx.getTenantService().findTenantById(tenantId) != null; - case DEVICE -> edgeCtx.getDeviceService().findDeviceById(tenantId, new DeviceId(entityId.getId())) != null; - case ASSET -> edgeCtx.getAssetService().findAssetById(tenantId, new AssetId(entityId.getId())) != null; - case ENTITY_VIEW -> edgeCtx.getEntityViewService().findEntityViewById(tenantId, new EntityViewId(entityId.getId())) != null; - case CUSTOMER -> edgeCtx.getCustomerService().findCustomerById(tenantId, new CustomerId(entityId.getId())) != null; - case USER -> edgeCtx.getUserService().findUserById(tenantId, new UserId(entityId.getId())) != null; - case DASHBOARD -> edgeCtx.getDashboardService().findDashboardById(tenantId, new DashboardId(entityId.getId())) != null; - case EDGE -> edgeCtx.getEdgeService().findEdgeById(tenantId, new EdgeId(entityId.getId())) != null; - default -> false; - }; + return entityDaoRegistry.getDao(entityId.getEntityType()).existsById(tenantId, entityId.getId()); } protected void createRelationFromEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java index 650aafb8cc..f0a8012660 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java @@ -21,23 +21,18 @@ import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EdgeUtils; -import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmComment; -import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.id.AlarmId; -import org.thingsboard.server.common.data.id.AssetId; -import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageDataIterableByTenantIdEntityId; +import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; @@ -59,6 +54,9 @@ public abstract class AlarmEdgeProcessor extends BaseAlarmProcessor implements A @Autowired private AlarmMsgConstructorFactory alarmMsgConstructorFactory; + @Autowired + private EntityService entityService; + @Override public ListenableFuture processAlarmMsgFromEdge(TenantId tenantId, EdgeId edgeId, AlarmUpdateMsg alarmUpdateMsg) { log.trace("[{}] processAlarmMsgFromEdge [{}]", tenantId, alarmUpdateMsg); @@ -178,42 +176,19 @@ public abstract class AlarmEdgeProcessor extends BaseAlarmProcessor implements A case ADDED, UPDATED, ALARM_ACK, ALARM_CLEAR -> { Alarm alarm = edgeCtx.getAlarmService().findAlarmById(tenantId, alarmId); if (alarm != null) { - return msgConstructor.constructAlarmUpdatedMsg(msgType, alarm, findOriginatorEntityName(tenantId, alarm)); + return msgConstructor.constructAlarmUpdatedMsg(msgType, alarm, + entityService.fetchEntityName(tenantId, alarm.getOriginator()).orElse(null)); } } case ALARM_DELETE, DELETED -> { Alarm deletedAlarm = JacksonUtil.convertValue(body, Alarm.class); if (deletedAlarm != null) { - return msgConstructor.constructAlarmUpdatedMsg(msgType, deletedAlarm, findOriginatorEntityName(tenantId, deletedAlarm)); + return msgConstructor.constructAlarmUpdatedMsg(msgType, deletedAlarm, + entityService.fetchEntityName(tenantId, deletedAlarm.getOriginator()).orElse(null)); } } } return null; } - private String findOriginatorEntityName(TenantId tenantId, Alarm alarm) { - String entityName = null; - switch (alarm.getOriginator().getEntityType()) { - case DEVICE -> { - Device deviceById = edgeCtx.getDeviceService().findDeviceById(tenantId, new DeviceId(alarm.getOriginator().getId())); - if (deviceById != null) { - entityName = deviceById.getName(); - } - } - case ASSET -> { - Asset assetById = edgeCtx.getAssetService().findAssetById(tenantId, new AssetId(alarm.getOriginator().getId())); - if (assetById != null) { - entityName = assetById.getName(); - } - } - case ENTITY_VIEW -> { - EntityView entityViewById = edgeCtx.getEntityViewService().findEntityViewById(tenantId, new EntityViewId(alarm.getOriginator().getId())); - if (entityViewById != null) { - entityName = entityViewById.getName(); - } - } - } - return entityName; - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java index 9ba863f62a..0c6da73a25 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java @@ -36,8 +36,7 @@ public abstract class BaseRelationProcessor extends BaseEdgeProcessor { switch (relationUpdateMsg.getMsgType()) { case ENTITY_CREATED_RPC_MESSAGE: case ENTITY_UPDATED_RPC_MESSAGE: - if (isEntityExists(tenantId, entityRelation.getTo()) - && isEntityExists(tenantId, entityRelation.getFrom())) { + if (isEntityExists(tenantId, entityRelation.getTo()) && isEntityExists(tenantId, entityRelation.getFrom())) { edgeCtx.getRelationService().saveRelation(tenantId, entityRelation); } else { log.warn("[{}] Skipping relating update msg because from/to entity doesn't exists on edge, {}", tenantId, relationUpdateMsg); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index a613a2b5a1..38cebbb509 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.entitiy; import com.fasterxml.jackson.core.type.TypeReference; import jakarta.annotation.PostConstruct; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionalEventListener; @@ -64,6 +65,7 @@ import java.util.Set; @Slf4j @Component +@RequiredArgsConstructor public class EntityStateSourcingListener { private final TopicService topicService; @@ -73,15 +75,6 @@ public class EntityStateSourcingListener { private final Optional kafkaSettings; private final Optional kafkaTopicConfigs; - public EntityStateSourcingListener(TopicService topicService, TbClusterService tbClusterService, TenantService tenantService, - Optional kafkaSettings, Optional kafkaTopicConfigs) { - this.topicService = topicService; - this.tbClusterService = tbClusterService; - this.tenantService = tenantService; - this.kafkaSettings = kafkaSettings; - this.kafkaTopicConfigs = kafkaTopicConfigs; - } - @PostConstruct public void init() { log.debug("EntityStateSourcingListener initiated"); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 8e4af4fa28..2a64dd3388 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -236,7 +236,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService edgeIds = new PageDataIterable<>(link -> edgeService.findEdgeIdsByTenantId(tenantId, link), 1024); long currentTimeMillis = System.currentTimeMillis(); + long ttlMillis = TimeUnit.SECONDS.toChronoUnit().getDuration().multipliedBy(ttlSeconds).toMillis(); for (EdgeId edgeId : edgeIds) { - Optional attributeOpt = attributesService.find(tenantId, edgeId, AttributeScope.SERVER_SCOPE, DefaultDeviceStateService.LAST_CONNECT_TIME).get(); - if (attributeOpt.isPresent()) { - Optional lastConnectTimeOpt = attributeOpt.get().getLongValue(); - if (lastConnectTimeOpt.isPresent() && isTopicExpired(lastConnectTimeOpt.get(), currentTimeMillis)) { - String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); - TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); - if (kafkaAdmin.isTopicEmpty(topic)) { - kafkaAdmin.deleteTopic(topic); - log.info("Removed outdated topic for tenant {} and edge with id {} older than {}", tenantId, edgeId, Date.from(Instant.ofEpochMilli(currentTimeMillis - ONE_MONTH_MILLIS))); - } - } - } + attributesService.find(tenantId, edgeId, AttributeScope.SERVER_SCOPE, DefaultDeviceStateService.LAST_CONNECT_TIME).get() + .flatMap(AttributeKvEntry::getLongValue) + .filter(lastConnectTime -> isTopicExpired(lastConnectTime, ttlMillis, currentTimeMillis)) + .ifPresent(lastConnectTime -> { + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); + if (kafkaAdmin.isTopicEmpty(topic)) { + kafkaAdmin.deleteTopic(topic); + log.info("Removed outdated topic for tenant {} and edge with id {} older than {}", + tenantId, edgeId, Date.from(Instant.ofEpochMilli(currentTimeMillis - ttlMillis))); + } + }); } } - private boolean isTopicExpired(long lastConnectTime, long currentTimeMillis) { - return lastConnectTime + ONE_MONTH_MILLIS < currentTimeMillis; + private boolean isTopicExpired(long lastConnectTime, long ttlMillis, long currentTimeMillis) { + return lastConnectTime + ttlMillis < currentTimeMillis; } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeEventService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeEventService.java index 9a7fc854d1..b7735228eb 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeEventService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeEventService.java @@ -28,10 +28,6 @@ public interface EdgeEventService { PageData findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink); - /** - * Executes stored procedure to cleanup old edge events. - * @param ttl the ttl for edge events in seconds - */ void cleanupEvents(long ttl); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java new file mode 100644 index 0000000000..b8736e913a --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java @@ -0,0 +1,59 @@ +/** + * Copyright © 2016-2024 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.edge; + +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.cache.limits.RateLimitService; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.limit.LimitedApi; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.msg.tools.TbRateLimitsException; +import org.thingsboard.server.dao.service.DataValidator; + +public abstract class BaseEdgeEventService implements EdgeEventService { + + @Autowired + private EdgeEventDao edgeEventDao; + @Autowired + private RateLimitService rateLimitService; + @Autowired + private DataValidator edgeEventValidator; + + @Override + public PageData findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) { + return edgeEventDao.findEdgeEvents(tenantId.getId(), edgeId, seqIdStart, seqIdEnd, pageLink); + } + + @Override + public void cleanupEvents(long ttl) { + edgeEventDao.cleanupEvents(ttl); + } + + protected void validateEdgeEvent(EdgeEvent edgeEvent) { + if (!rateLimitService.checkRateLimit(LimitedApi.EDGE_EVENTS, edgeEvent.getTenantId())) { + throw new TbRateLimitsException(EntityType.TENANT); + } + if (!rateLimitService.checkRateLimit(LimitedApi.EDGE_EVENTS_PER_EDGE, edgeEvent.getTenantId(), edgeEvent.getEdgeId())) { + throw new TbRateLimitsException(EntityType.EDGE); + } + edgeEventValidator.validate(edgeEvent, EdgeEvent::getTenantId); + } + +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java index 1f9562724d..572733d75b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java @@ -56,6 +56,4 @@ public interface EdgeEventDao extends Dao { */ void cleanupEvents(long ttl); - void migrateEdgeEvents(); - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java index 8d23c8ee5d..b021ac90df 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java @@ -27,17 +27,8 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.server.cache.limits.RateLimitService; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.EdgeEvent; -import org.thingsboard.server.common.data.id.EdgeId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.limit.LimitedApi; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.TimePageLink; -import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; -import org.thingsboard.server.dao.service.DataValidator; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -46,11 +37,9 @@ import java.util.concurrent.Executors; @Service @RequiredArgsConstructor @ConditionalOnExpression("'${queue.type:null}'!='kafka'") -public class PostgresEdgeEventService implements EdgeEventService { +public class PostgresEdgeEventService extends BaseEdgeEventService { private final EdgeEventDao edgeEventDao; - private final RateLimitService rateLimitService; - private final DataValidator edgeEventValidator; private final ApplicationEventPublisher eventPublisher; private ExecutorService edgeEventExecutor; @@ -69,14 +58,7 @@ public class PostgresEdgeEventService implements EdgeEventService { @Override public ListenableFuture saveAsync(EdgeEvent edgeEvent) { - if (!rateLimitService.checkRateLimit(LimitedApi.EDGE_EVENTS, edgeEvent.getTenantId())) { - throw new TbRateLimitsException(EntityType.TENANT); - } - if (!rateLimitService.checkRateLimit(LimitedApi.EDGE_EVENTS_PER_EDGE, edgeEvent.getTenantId(), edgeEvent.getEdgeId())) { - throw new TbRateLimitsException(EntityType.EDGE); - } - edgeEventValidator.validate(edgeEvent, EdgeEvent::getTenantId); - + validateEdgeEvent(edgeEvent); ListenableFuture saveFuture = edgeEventDao.saveAsync(edgeEvent); Futures.addCallback(saveFuture, new FutureCallback<>() { @@ -96,14 +78,4 @@ public class PostgresEdgeEventService implements EdgeEventService { return saveFuture; } - @Override - public PageData findEdgeEvents(TenantId tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) { - return edgeEventDao.findEdgeEvents(tenantId.getId(), edgeId, seqIdStart, seqIdEnd, pageLink); - } - - @Override - public void cleanupEvents(long ttl) { - edgeEventDao.cleanupEvents(ttl); - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java index d9d2282fca..065a057160 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java @@ -198,36 +198,6 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao 0 ? System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(edgeEventsTtl) : 1629158400000L; - - long currentTime = System.currentTimeMillis(); - var partitionStepInMs = TimeUnit.HOURS.toMillis(partitionSizeInHours); - long numberOfPartitions = (currentTime - startTime) / partitionStepInMs; - - if (numberOfPartitions > 1000) { - String error = "Please adjust your edge event partitioning configuration. Configuration with partition size " + - "of " + partitionSizeInHours + " hours and corresponding TTL will use " + numberOfPartitions + " " + - "(> 1000) partitions which is not recommended!"; - log.error(error); - throw new RuntimeException(error); - } - - while (startTime < currentTime) { - var endTime = startTime + partitionStepInMs; - log.info("Migrating edge event for time period: {} - {}", startTime, endTime); - callMigrationFunction(startTime, endTime, partitionStepInMs); - startTime = endTime; - } - log.info("Event edge migration finished"); - jdbcTemplate.execute("DROP TABLE IF EXISTS old_edge_event"); - } - - private void callMigrationFunction(long startTime, long endTime, long partitionSIzeInMs) { - jdbcTemplate.update("CALL migrate_edge_event(?, ?, ?)", startTime, endTime, partitionSIzeInMs); - } - @Override public void createPartition(EdgeEventEntity entity) { partitioningRepository.createPartitionIfNotExists(TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours)); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java index 2cc179fae2..be8cb5d749 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java @@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.dao.edge.EdgeEventDao; import org.thingsboard.server.dao.edge.EdgeEventService; import java.io.IOException; @@ -48,6 +49,9 @@ public class EdgeEventServiceTest extends AbstractServiceTest { @Autowired EdgeEventService edgeEventService; + @Autowired + EdgeEventDao edgeEventDao; + long timeBeforeStartTime; long startTime; long eventTime; @@ -81,7 +85,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest { Assert.assertEquals(saved.getAction(), edgeEvent.getAction()); Assert.assertEquals(saved.getBody(), edgeEvent.getBody()); - edgeEventService.cleanupEvents(1); + edgeEventDao.cleanupEvents(1); } protected EdgeEvent generateEdgeEvent(TenantId tenantId, EdgeId edgeId, EntityId entityId) throws IOException { @@ -129,7 +133,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest { Assert.assertEquals(Uuids.startOf(eventTime), edgeEvents.getData().get(0).getUuidId()); Assert.assertFalse(edgeEvents.hasNext()); - edgeEventService.cleanupEvents(1); + edgeEventDao.cleanupEvents(1); } private ListenableFuture saveEdgeEventWithProvidedTime(long time, EdgeId edgeId, EntityId entityId, TenantId tenantId) throws Exception {