Browse Source
* Improve edge notification for entities' CRUD operations. Use service layer to notify instead of TbService * Improve queue service, delete unused class for edge event updates * Improve alarm delete and add handle fox delete dao event notification * Refactoring: provide notification for relations and alarms. Improve logic and bad edge event type using * Add entity type to SaveEvent to process correct message type to edge * Improve relation service publish event * Introduce EdgeEventSourcing service instead of saving edge events on controller/service layers. Part #2 * Improved stability of device edge test * Push credential updated event only in case update * Add tenantId to saveUser signature to send correct notification for listener * Fix tests to send correct notification msg to edge * Fix tests with correct action type * Add delete msg to edge for customer * Refactor ActionEntityEvent to use lombok builder * Remove unnecessary comments * Added edgeSynchronizationManager into BaseAlarmProcessor and BaseRelationProcessor * Fixed license header * Remove notification to edge from Version Control Service * Fixed alarm del processing - find related edges inside edge processor * Fix controller test for publish event to listener if entity was deleted * Added check for edge imitator messages during login as tenant admin * Refactoring: Added filtering of relation on EdgeEventSourcingListener * Refactored to be in sync with PE * Refactored edge test to be in sync with PE edge test changes * EdgeControllerTest - moved await block into separate method to reuse it * Fixed EdgeControllerTest * Fixed testAssignEdgeToCustomerFromDifferentTenant test * testSyncEdge - make stable * Refacroting - update utils method name to pop* in EdgeControllerTest * testSyncEdge - fixed order and nubmer of edge events * testGetEdgeEvents - check by pop items, and not by index to improve stability on slow machines * testGetEdgeEvents - added check that list is empty * Removed test debug output * EntityServiceTest - Fixed compilation error after merge * Improve service layer event publisher to process each notification and validate in listener * Improve BaseAlarmService to send notification to listener * Fix asset-device notification action to send delete to all edges * Delete unnecessary usage of sendMsgToEdge * Improve processEntityNotification to be in sync with changed needed for PE * Pull request review - minor refactoring * Fix tests after review-refactoring * Refactor tests to be in sync with PE * Fixed repeated update - added check for old_edge_event table existance before migration * DeviceEdgeProcessor - do edgeSynchronizationManager as soon as possible to avoid unnecessary downlinks * BaseEdgeProcessor - refactoring and remove duplicate methods. Introduce EdgeEventType.isAllEdgesRelated * Organize imports * Improve Edge test: add sync completed message to await * Minor refactoring for EdgeProcessor notification: asset and device * EdgeEventSourcingListener - updated logging to avoid null pointer exception * BaseAlarmService - added check for alarm to avoid NPE. EdgeEventSourcingListener - added try/catch blocks * EdgeEventSourcingListener - fixed error message log level --------- Co-authored-by: Volodymyr Babak <volodymyr.babak@gmail.com>pull/9052/head
committed by
GitHub
122 changed files with 1524 additions and 1171 deletions
@ -0,0 +1,158 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.edge; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Component; |
|||
import org.springframework.transaction.event.TransactionalEventListener; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.cluster.TbClusterService; |
|||
import org.thingsboard.server.common.data.OtaPackageInfo; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.alarm.AlarmApiCallResult; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.relation.EntityRelation; |
|||
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|||
import org.thingsboard.server.common.data.rule.RuleChain; |
|||
import org.thingsboard.server.common.data.rule.RuleChainType; |
|||
import org.thingsboard.server.common.data.security.Authority; |
|||
import org.thingsboard.server.dao.edge.EdgeSynchronizationManager; |
|||
import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; |
|||
import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; |
|||
import org.thingsboard.server.dao.eventsourcing.RelationActionEvent; |
|||
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
|
|||
import static org.thingsboard.server.service.entitiy.DefaultTbNotificationEntityService.edgeTypeByActionType; |
|||
|
|||
|
|||
/** |
|||
* This event listener does not support async event processing because relay on ThreadLocal |
|||
* Another possible approach is to implement a special annotation and a bunch of classes similar to TransactionalApplicationListener |
|||
* This class is the simplest approach to maintain edge synchronization within the single class. |
|||
* <p> |
|||
* For async event publishers, you have to decide whether publish event on creating async task in the same thread where dao method called |
|||
* @Autowired |
|||
* EdgeEventSynchronizationManager edgeSynchronizationManager |
|||
* ... |
|||
* //some async write action make future
|
|||
* if (!edgeSynchronizationManager.isSync()) { |
|||
* future.addCallback(eventPublisher.publishEvent(...)) |
|||
* } |
|||
* */ |
|||
@Component |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class EdgeEventSourcingListener { |
|||
|
|||
private final TbClusterService tbClusterService; |
|||
private final EdgeSynchronizationManager edgeSynchronizationManager; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
log.info("EdgeEventSourcingListener initiated"); |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(SaveEntityEvent<?> event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
if (!isValidEdgeEventEntity(event.getEntity())) { |
|||
return; |
|||
} |
|||
log.trace("[{}] SaveEntityEvent called: {}", event.getTenantId(), event); |
|||
EdgeEventActionType action = Boolean.TRUE.equals(event.getAdded()) ? EdgeEventActionType.ADDED : EdgeEventActionType.UPDATED; |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, event.getEntityId(), |
|||
null, null, action); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process SaveEntityEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(DeleteEntityEvent<?> event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(), |
|||
JacksonUtil.toString(event.getEntity()), null, EdgeEventActionType.DELETED); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process DeleteEntityEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(ActionEntityEvent event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event); |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(), |
|||
event.getBody(), null, edgeTypeByActionType(event.getActionType())); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process ActionEntityEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(RelationActionEvent event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
EntityRelation relation = event.getRelation(); |
|||
if (relation == null) { |
|||
log.trace("[{}] skipping RelationActionEvent event in case relation is null: {}", event.getTenantId(), event); |
|||
return; |
|||
} |
|||
if (!RelationTypeGroup.COMMON.equals(relation.getTypeGroup())) { |
|||
log.trace("[{}] skipping RelationActionEvent event in case NOT COMMON relation type group: {}", event.getTenantId(), event); |
|||
return; |
|||
} |
|||
log.trace("[{}] RelationActionEvent called: {}", event.getTenantId(), event); |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, null, |
|||
JacksonUtil.toString(relation), EdgeEventType.RELATION, edgeTypeByActionType(event.getActionType())); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process RelationActionEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
private boolean isValidEdgeEventEntity(Object entity) { |
|||
if (entity instanceof OtaPackageInfo) { |
|||
OtaPackageInfo otaPackageInfo = (OtaPackageInfo) entity; |
|||
return otaPackageInfo.hasUrl() || otaPackageInfo.isHasData(); |
|||
} else if (entity instanceof RuleChain) { |
|||
RuleChain ruleChain = (RuleChain) entity; |
|||
return RuleChainType.EDGE.equals(ruleChain.getType()); |
|||
} else if (entity instanceof User) { |
|||
User user = (User) entity; |
|||
return !Authority.SYS_ADMIN.equals(user.getAuthority()); |
|||
} else if (entity instanceof AlarmApiCallResult) { |
|||
AlarmApiCallResult alarmApiCallResult = (AlarmApiCallResult) entity; |
|||
return alarmApiCallResult.isModified(); |
|||
} |
|||
// Default: If the entity doesn't match any of the conditions, consider it as valid.
|
|||
return true; |
|||
} |
|||
} |
|||
@ -0,0 +1,23 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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; |
|||
|
|||
public interface EdgeSynchronizationManager { |
|||
|
|||
ThreadLocal<Boolean> getSync(); |
|||
|
|||
boolean isSync(); |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue