Browse Source

Improvement after review

pull/11924/head
Andrii Landiak 2 years ago
parent
commit
732344c8d7
  1. 13
      application/src/main/java/org/thingsboard/server/service/edge/rpc/AbstractEdgeGrpcSession.java
  2. 36
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  3. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  4. 41
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeEventService.java
  5. 36
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java
  6. 9
      application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java
  7. 16
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java
  8. 41
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java
  9. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java
  10. 11
      application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java
  11. 2
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  12. 1
      application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java
  13. 36
      application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java
  14. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeEventService.java
  15. 59
      dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java
  16. 2
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java
  17. 32
      dao/src/main/java/org/thingsboard/server/dao/edge/PostgresEdgeEventService.java
  18. 30
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java
  19. 8
      dao/src/test/java/org/thingsboard/server/dao/service/EdgeEventServiceTest.java

13
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<T extends AbstractEdgeGrpcSession<T>> 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<T extends AbstractEdgeGrpcSession<
protected static final ConcurrentLinkedQueue<EdgeEvent> highPriorityQueue = new ConcurrentLinkedQueue<>();
protected UUID sessionId;
private BiConsumer<EdgeId, T> sessionOpenListener;
private BiConsumer<EdgeId, AbstractEdgeGrpcSession> sessionOpenListener;
private BiConsumer<Edge, UUID> sessionCloseListener;
private final EdgeSessionState sessionState = new EdgeSessionState();
@ -146,7 +146,7 @@ public abstract class AbstractEdgeGrpcSession<T extends AbstractEdgeGrpcSession<
private ScheduledExecutorService sendDownlinkExecutorService;
public AbstractEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver<ResponseMsg> outputStream,
BiConsumer<EdgeId, T> sessionOpenListener,
BiConsumer<EdgeId, AbstractEdgeGrpcSession> sessionOpenListener,
BiConsumer<Edge, UUID> sessionCloseListener,
ScheduledExecutorService sendDownlinkExecutorService,
int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) {
@ -299,11 +299,6 @@ public abstract class AbstractEdgeGrpcSession<T extends AbstractEdgeGrpcSession<
if (isConnected() && !pageData.getData().isEmpty()) {
log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, sessionId, pageData.getData().size());
List<DownlinkMsg> 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<T extends AbstractEdgeGrpcSession<
tenantId = edge.getTenantId();
try {
if (edge.getSecret().equals(request.getEdgeSecret())) {
sessionOpenListener.accept(edge.getId(), (T) this);
sessionOpenListener.accept(edge.getId(), this);
edgeVersion = request.getEdgeVersion();
processSaveEdgeVersionAsAttribute(request.getEdgeVersion().name());
return ConnectResponseMsg.newBuilder()

36
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -86,7 +86,7 @@ import java.util.function.Consumer;
@TbCoreComponent
public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService {
private final ConcurrentMap<EdgeId, AbstractEdgeGrpcSession<?>> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, AbstractEdgeGrpcSession> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, Lock> sessionNewEventsLocks = new ConcurrentHashMap<>();
private final Map<EdgeId, Boolean> sessionNewEvents = new HashMap<>();
private final ConcurrentMap<EdgeId, ScheduledFuture<?>> sessionEdgeEventChecks = new ConcurrentHashMap<>();
@ -210,11 +210,11 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) {
AbstractEdgeGrpcSession<?> session = createEdgeGrpcSession(outputStream);
AbstractEdgeGrpcSession session = createEdgeGrpcSession(outputStream);
return session.getInputStream();
}
private AbstractEdgeGrpcSession<?> createEdgeGrpcSession(StreamObserver<ResponseMsg> outputStream) {
private AbstractEdgeGrpcSession createEdgeGrpcSession(StreamObserver<ResponseMsg> 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();

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -26,7 +26,9 @@ public interface EdgeGrpcSession {
boolean isConnected();
ListenableFuture<Boolean> migrateEdgeEvents(boolean isMigrationProcessed) throws Exception;
void destroy();
ListenableFuture<Boolean> migrateEdgeEvents() throws Exception;
ListenableFuture<Boolean> processEdgeEvents() throws Exception;

41
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<EdgeEvent> edgeEventValidator;
@Lazy
private final TbQueueProducerProvider producerProvider;
@Lazy
private final TopicService topicService;
private final EdgeEventDao edgeEventDao;
@Lazy
private final TbQueueProducerProvider producerProvider;
@Override
public ListenableFuture<Void> 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<EdgeEvent> 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);
}
}

36
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<KafkaEdgeGrpcSession> {
public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession {
private final TbQueueConsumer<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> edgeEventsConsumer;
private volatile boolean isHighPriorityProcessing;
private volatile boolean isConsumerInit;
private QueueConsumerManager<TbProtoQueueMsg<ToEdgeEventNotificationMsg>> consumer;
private ExecutorService consumerExecutor;
public KafkaEdgeGrpcSession(EdgeContextComponent ctx, TbCoreQueueFactory tbCoreQueueFactory, StreamObserver<ResponseMsg> outputStream,
BiConsumer<EdgeId, KafkaEdgeGrpcSession> sessionOpenListener,
BiConsumer<EdgeId, AbstractEdgeGrpcSession> sessionOpenListener,
BiConsumer<Edge, UUID> sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService,
int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) {
super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession);
@ -82,32 +80,23 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession<KafkaEdgeGrpcS
log.error("[{}] Failed to process all downlink messages", sessionId, e);
}
} else {
try {
Thread.sleep(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval());
} catch (InterruptedException interruptedException) {
log.trace("Failed to wait until the server has capacity to handle new requests", interruptedException);
}
log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId);
}
}
@PreDestroy
private void destroy() {
consumer.stop();
consumerExecutor.shutdown();
}
public void stopConsumer() {
consumer.stop();
consumerExecutor.shutdown();
}
@Override
public ListenableFuture<Boolean> migrateEdgeEvents(boolean isMigrationProcessed) throws Exception {
if (isMigrationProcessed) {
return Futures.immediateFuture(Boolean.FALSE);
}
public ListenableFuture<Boolean> migrateEdgeEvents() throws Exception {
return super.processEdgeEvents();
}
@Override
public ListenableFuture<Boolean> processEdgeEvents() {
if (!isConsumerInit) {
if (consumer == null) {
this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer"));
this.consumer = QueueConsumerManager.<TbProtoQueueMsg<ToEdgeEventNotificationMsg>>builder()
.name("TB Edge events")
@ -117,7 +106,6 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession<KafkaEdgeGrpcS
.consumerExecutor(consumerExecutor)
.threadPrefix("edge-events")
.build();
isConsumerInit = true;
consumer.subscribe();
consumer.launch();
}
@ -130,4 +118,10 @@ public class KafkaEdgeGrpcSession extends AbstractEdgeGrpcSession<KafkaEdgeGrpcS
super.processHighPriorityEvents();
}
@Override
public void destroy() {
consumer.stop();
consumerExecutor.shutdown();
}
}

9
application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java

@ -29,10 +29,10 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.function.BiConsumer;
@Slf4j
public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession<PostgresEdgeGrpcSession> {
public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession {
PostgresEdgeGrpcSession(EdgeContextComponent ctx, StreamObserver<ResponseMsg> outputStream,
BiConsumer<EdgeId, PostgresEdgeGrpcSession> sessionOpenListener,
BiConsumer<EdgeId, AbstractEdgeGrpcSession> sessionOpenListener,
BiConsumer<Edge, UUID> sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService,
int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) {
super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession);
@ -40,7 +40,10 @@ public class PostgresEdgeGrpcSession extends AbstractEdgeGrpcSession<PostgresEdg
}
@Override
public ListenableFuture<Boolean> migrateEdgeEvents(boolean isMigrationProcessed) {
public void destroy() {}
@Override
public ListenableFuture<Boolean> migrateEdgeEvents() {
return Futures.immediateFuture(Boolean.FALSE);
}

16
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) {

41
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<Void> 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;
}
}

3
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);

11
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<TbKafkaSettings> kafkaSettings;
private final Optional<TbKafkaTopicConfigs> kafkaTopicConfigs;
public EntityStateSourcingListener(TopicService topicService, TbClusterService tbClusterService, TenantService tenantService,
Optional<TbKafkaSettings> kafkaSettings, Optional<TbKafkaTopicConfigs> kafkaTopicConfigs) {
this.topicService = topicService;
this.tbClusterService = tbClusterService;
this.tenantService = tenantService;
this.kafkaSettings = kafkaSettings;
this.kafkaTopicConfigs = kafkaTopicConfigs;
}
@PostConstruct
public void init() {
log.debug("EntityStateSourcingListener initiated");

2
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -236,7 +236,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
@Override
protected void startConsumers() {
super.startConsumers();
firmwareStatesConsumer. subscribe();
firmwareStatesConsumer.subscribe();
firmwareStatesConsumer.launch();
usageStatsConsumer.launch();
}

1
application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java

@ -20,6 +20,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.edge.EdgeEventDao;
import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import org.thingsboard.server.queue.discovery.PartitionService;

36
application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.service.ttl;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
@ -40,7 +41,6 @@ import org.thingsboard.server.service.state.DefaultDeviceStateService;
import java.time.Instant;
import java.util.Date;
import java.util.Optional;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
@ -52,8 +52,6 @@ import java.util.concurrent.TimeUnit;
@ConditionalOnExpression("'${queue.type:null}'=='kafka' && ${edges.enabled:true}")
public class KafkaEdgeTopicsCleanUpService {
private static final long ONE_MONTH_MILLIS = TimeUnit.DAYS.toChronoUnit().getDuration().multipliedBy(30).toMillis();
private final EdgeService edgeService;
private final TenantService tenantService;
private final AttributesService attributesService;
@ -64,6 +62,9 @@ public class KafkaEdgeTopicsCleanUpService {
private final TbKafkaSettings kafkaSettings;
private final TbKafkaTopicConfigs kafkaTopicConfigs;
@Value("${sql.ttl.edge_events.edge_events_ttl:2628000}")
private long ttlSeconds;
private final ExecutorService executorService = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("kafka-edge-topic-cleanup"));
@Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.edge_events.execution_interval_ms})}", fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}")
@ -87,25 +88,26 @@ public class KafkaEdgeTopicsCleanUpService {
PageDataIterable<EdgeId> 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<AttributeKvEntry> attributeOpt = attributesService.find(tenantId, edgeId, AttributeScope.SERVER_SCOPE, DefaultDeviceStateService.LAST_CONNECT_TIME).get();
if (attributeOpt.isPresent()) {
Optional<Long> 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;
}
}

4
common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeEventService.java

@ -28,10 +28,6 @@ public interface EdgeEventService {
PageData<EdgeEvent> 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);
}

59
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<EdgeEvent> edgeEventValidator;
@Override
public PageData<EdgeEvent> 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);
}
}

2
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java

@ -56,6 +56,4 @@ public interface EdgeEventDao extends Dao<EdgeEvent> {
*/
void cleanupEvents(long ttl);
void migrateEdgeEvents();
}

32
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<EdgeEvent> edgeEventValidator;
private final ApplicationEventPublisher eventPublisher;
private ExecutorService edgeEventExecutor;
@ -69,14 +58,7 @@ public class PostgresEdgeEventService implements EdgeEventService {
@Override
public ListenableFuture<Void> 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<Void> saveFuture = edgeEventDao.saveAsync(edgeEvent);
Futures.addCallback(saveFuture, new FutureCallback<>() {
@ -96,14 +78,4 @@ public class PostgresEdgeEventService implements EdgeEventService {
return saveFuture;
}
@Override
public PageData<EdgeEvent> 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);
}
}

30
dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java

@ -198,36 +198,6 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao<EdgeEventEnti
partitioningRepository.dropPartitionsBefore(TABLE_NAME, ttl, TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
@Override
public void migrateEdgeEvents() {
long startTime = edgeEventsTtl > 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));

8
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<Void> saveEdgeEventWithProvidedTime(long time, EdgeId edgeId, EntityId entityId, TenantId tenantId) throws Exception {

Loading…
Cancel
Save