Browse Source

Edge: impl new approach with oauth

pull/11231/head
Andrii Landiak 2 years ago
parent
commit
0ab35b72d2
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  3. 13
      application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java
  4. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  6. 34
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/oauth2/OAuth2MsgConstructor.java
  7. 29
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java
  8. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java
  9. 19
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java
  10. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/profile/DeviceProfileEdgeProcessor.java
  11. 94
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/oauth2/OAuth2EdgeProcessor.java
  12. 8
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  13. 8
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  14. 6
      application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java
  15. 19
      application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
  16. 144
      application/src/test/java/org/thingsboard/server/edge/OAuth2EdgeTest.java
  17. 14
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  18. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientService.java
  19. 3
      common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java
  20. 4
      common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java
  21. 2
      common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java
  22. 19
      common/edge-api/src/main/proto/edge.proto
  23. 4
      dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java
  24. 3
      dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientDao.java
  25. 8
      dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java
  26. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/domain/JpaDomainDao.java
  27. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/oauth2/JpaOAuth2ClientDao.java
  28. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/oauth2/OAuth2ClientRepository.java

2
application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java

@ -201,7 +201,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
case NOTIFICATION_RULE, NOTIFICATION_TARGET, NOTIFICATION_TEMPLATE ->
notificationEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
case TB_RESOURCE -> resourceEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
case OAUTH2_CLIENT -> oAuth2EdgeProcessor.processOAuth2Notification(tenantId, edgeNotificationMsg);
case DOMAIN, OAUTH2_CLIENT -> oAuth2EdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
default -> log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type);
}
} catch (Exception e) {

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

@ -29,13 +29,13 @@ import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.dashboard.DashboardService;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.domain.DomainService;
import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.edge.EdgeService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.notification.NotificationRuleService;
import org.thingsboard.server.dao.notification.NotificationTargetService;
import org.thingsboard.server.dao.notification.NotificationTemplateService;
import org.thingsboard.server.dao.oauth2.OAuth2ClientService;
import org.thingsboard.server.dao.ota.OtaPackageService;
import org.thingsboard.server.dao.queue.QueueService;
import org.thingsboard.server.dao.resource.ResourceService;
@ -167,7 +167,7 @@ public class EdgeContextComponent {
private NotificationTemplateService notificationTemplateService;
@Autowired
private OAuth2ClientService oAuth2ClientService;
private DomainService domainService;
@Autowired
private RateLimitService rateLimitService;

13
application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java

@ -32,10 +32,10 @@ import org.thingsboard.server.common.data.alarm.AlarmApiCallResult;
import org.thingsboard.server.common.data.alarm.AlarmComment;
import org.thingsboard.server.common.data.alarm.EntityAlarm;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.domain.Domain;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.rule.RuleChain;
@ -204,11 +204,12 @@ public class EdgeEventSourcingListener {
return !event.getCreated();
case API_USAGE_STATE, EDGE:
return false;
case DOMAIN:
if (entity instanceof Domain domain) {
return domain.isPropagateToEdge();
}
}
}
// if (entity instanceof OAuth2Info oAuth2Info) {
// return oAuth2Info.isEdgeEnabled();
// }
// Default: If the entity doesn't match any of the conditions, consider it as valid.
return true;
}
@ -233,8 +234,6 @@ public class EdgeEventSourcingListener {
private EdgeEventType getEdgeEventTypeForEntityEvent(Object entity) {
if (entity instanceof AlarmComment) {
return EdgeEventType.ALARM_COMMENT;
} else if (entity instanceof OAuth2Client) {
return EdgeEventType.OAUTH2_CLIENT;
}
return null;
}
@ -242,8 +241,6 @@ public class EdgeEventSourcingListener {
private String getBodyMsgForEntityEvent(Object entity) {
if (entity instanceof AlarmComment) {
return JacksonUtil.toString(entity);
} else if (entity instanceof OAuth2Client) {
return JacksonUtil.toString(entity);
}
return null;
}

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

@ -699,7 +699,9 @@ public final class EdgeGrpcSession implements Closeable {
case NOTIFICATION_TEMPLATE:
return ctx.getNotificationEdgeProcessor().convertNotificationTemplateToDownlink(edgeEvent);
case OAUTH2_CLIENT:
return ctx.getOAuth2EdgeProcessor().convertOAuth2ProviderEventToDownlink(edgeEvent);
return ctx.getOAuth2EdgeProcessor().convertOAuth2ClientEventToDownlink(edgeEvent, this.edgeVersion);
case DOMAIN:
return ctx.getOAuth2EdgeProcessor().convertOAuth2DomainEventToDownlink(edgeEvent, this.edgeVersion);
default:
log.warn("[{}] Unsupported edge event type [{}]", this.tenantId, edgeEvent);
return null;

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java

@ -88,7 +88,7 @@ public class EdgeSyncCursor {
fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService()));
fetchers.add(new TenantResourcesEdgeEventFetcher(ctx.getResourceService()));
fetchers.add(new OAuth2EdgeEventFetcher(ctx.getOAuth2ClientService()));
fetchers.add(new OAuth2EdgeEventFetcher(ctx.getDomainService()));
}
}

34
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/oauth2/OAuth2MsgConstructor.java

@ -17,16 +17,44 @@ package org.thingsboard.server.service.edge.rpc.constructor.oauth2;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.domain.DomainInfo;
import org.thingsboard.server.common.data.id.DomainId;
import org.thingsboard.server.common.data.id.OAuth2ClientId;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.gen.edge.v1.OAuth2UpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.queue.util.TbCoreComponent;
@Component
@TbCoreComponent
public class OAuth2MsgConstructor {
public OAuth2UpdateMsg constructOAuth2UpdateMsg(OAuth2Client oAuth2Client) {
return OAuth2UpdateMsg.newBuilder().setEntity(JacksonUtil.toString(oAuth2Client)).build();
public OAuth2ClientUpdateMsg constructOAuth2ClientUpdateMsg(UpdateMsgType msgType, OAuth2Client oAuth2Client) {
return OAuth2ClientUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(oAuth2Client))
.setIdMSB(oAuth2Client.getId().getId().getMostSignificantBits())
.setIdLSB(oAuth2Client.getId().getId().getLeastSignificantBits()).build();
}
public OAuth2ClientUpdateMsg constructOAuth2ClientDeleteMsg(OAuth2ClientId oAuth2ClientId) {
return OAuth2ClientUpdateMsg.newBuilder()
.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE)
.setIdMSB(oAuth2ClientId.getId().getMostSignificantBits())
.setIdLSB(oAuth2ClientId.getId().getLeastSignificantBits()).build();
}
public OAuth2DomainUpdateMsg constructOAuth2DomainUpdateMsg(UpdateMsgType msgType, DomainInfo domainInfo) {
return OAuth2DomainUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(domainInfo))
.setIdMSB(domainInfo.getId().getId().getMostSignificantBits())
.setIdLSB(domainInfo.getId().getId().getLeastSignificantBits()).build();
}
public OAuth2DomainUpdateMsg constructOAuth2DomainDeleteMsg(DomainId domainId) {
return OAuth2DomainUpdateMsg.newBuilder()
.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE)
.setIdMSB(domainId.getId().getMostSignificantBits())
.setIdLSB(domainId.getId().getLeastSignificantBits())
.build();
}
}

29
application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/OAuth2EdgeEventFetcher.java

@ -17,43 +17,32 @@ package org.thingsboard.server.service.edge.rpc.fetch;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.domain.DomainInfo;
import org.thingsboard.server.common.data.edge.Edge;
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.TenantId;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.oauth2.OAuth2ClientService;
import java.util.ArrayList;
import java.util.List;
import org.thingsboard.server.dao.domain.DomainService;
@AllArgsConstructor
@Slf4j
public class OAuth2EdgeEventFetcher implements EdgeEventFetcher {
public class OAuth2EdgeEventFetcher extends BasePageableEdgeEventFetcher<DomainInfo> {
private final OAuth2ClientService oAuth2ClientService;
private final DomainService domainService;
@Override
public PageLink getPageLink(int pageSize) {
return null;
PageData<DomainInfo> fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) {
return domainService.findDomainInfosByTenantId(TenantId.SYS_TENANT_ID, pageLink);
}
@Override
public PageData<EdgeEvent> fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) {
// OAuth2Info oAuth2Info = oAuth2Service.findOAuth2Info();
// if (!oAuth2Info.isEdgeEnabled()) {
// return new PageData<>();
// }
List<EdgeEvent> result = new ArrayList<>();
// result.add(EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.OAUTH2,
// EdgeEventActionType.ADDED, null, JacksonUtil.valueToTree(oAuth2Info)));
// returns PageData object to be in sync with other fetchers
return new PageData<>(result, 1, result.size(), false);
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, DomainInfo domainInfo) {
return EdgeUtils.constructEdgeEvent(TenantId.SYS_TENANT_ID, edge.getId(), EdgeEventType.DOMAIN,
EdgeEventActionType.ADDED, domainInfo.getId(), null);
}
}

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

@ -69,6 +69,7 @@ import org.thingsboard.server.dao.dashboard.DashboardService;
import org.thingsboard.server.dao.device.DeviceCredentialsService;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.domain.DomainService;
import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.edge.EdgeService;
import org.thingsboard.server.dao.edge.EdgeSynchronizationManager;
@ -242,6 +243,9 @@ public abstract class BaseEdgeProcessor {
@Autowired
protected OAuth2ClientService oAuth2ClientService;
@Autowired
protected DomainService domainService;
@Autowired
@Lazy
protected TbQueueProducerProvider producerProvider;

19
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java

@ -216,6 +216,7 @@ public abstract class DeviceEdgeProcessor extends BaseDeviceProcessor implements
public DownlinkMsg convertDeviceEventToDownlink(EdgeEvent edgeEvent, EdgeId edgeId, EdgeVersion edgeVersion) {
DeviceId deviceId = new DeviceId(edgeEvent.getEntityId());
DownlinkMsg downlinkMsg = null;
var msgConstructor = (DeviceMsgConstructor) deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion);
switch (edgeEvent.getAction()) {
case ADDED:
case UPDATED:
@ -233,24 +234,20 @@ public abstract class DeviceEdgeProcessor extends BaseDeviceProcessor implements
.addDeviceUpdateMsg(deviceUpdateMsg);
DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(edgeEvent.getTenantId(), deviceId);
if (deviceCredentials != null) {
DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = ((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion)).constructDeviceCredentialsUpdatedMsg(deviceCredentials);
DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = msgConstructor.constructDeviceCredentialsUpdatedMsg(deviceCredentials);
builder.addDeviceCredentialsUpdateMsg(deviceCredentialsUpdateMsg).build();
}
if (UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE.equals(msgType)) {
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(edgeEvent.getTenantId(), device.getDeviceProfileId());
deviceProfile = checkIfDeviceProfileDefaultFieldsAssignedToEdge(edgeEvent.getTenantId(), edgeId, deviceProfile, edgeVersion);
builder.addDeviceProfileUpdateMsg(((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion))
.constructDeviceProfileUpdatedMsg(msgType, deviceProfile));
builder.addDeviceProfileUpdateMsg(msgConstructor.constructDeviceProfileUpdatedMsg(msgType, deviceProfile));
}
downlinkMsg = builder.build();
}
break;
case DELETED:
case UNASSIGNED_FROM_EDGE:
DeviceUpdateMsg deviceUpdateMsg = ((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion)).constructDeviceDeleteMsg(deviceId);
DeviceUpdateMsg deviceUpdateMsg = msgConstructor.constructDeviceDeleteMsg(deviceId);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceUpdateMsg(deviceUpdateMsg)
@ -259,8 +256,7 @@ public abstract class DeviceEdgeProcessor extends BaseDeviceProcessor implements
case CREDENTIALS_UPDATED:
DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(edgeEvent.getTenantId(), deviceId);
if (deviceCredentials != null) {
DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = ((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion)).constructDeviceCredentialsUpdatedMsg(deviceCredentials);
DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = msgConstructor.constructDeviceCredentialsUpdatedMsg(deviceCredentials);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceCredentialsUpdateMsg(deviceCredentialsUpdateMsg)
@ -270,11 +266,10 @@ public abstract class DeviceEdgeProcessor extends BaseDeviceProcessor implements
case RPC_CALL:
return DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceRpcCallMsg(((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion))
.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()))
.addDeviceRpcCallMsg(msgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()))
.build();
}
return downlinkMsg;
}
}

7
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/profile/DeviceProfileEdgeProcessor.java

@ -96,14 +96,14 @@ public abstract class DeviceProfileEdgeProcessor extends BaseDeviceProfileProces
public DownlinkMsg convertDeviceProfileEventToDownlink(EdgeEvent edgeEvent, EdgeId edgeId, EdgeVersion edgeVersion) {
DeviceProfileId deviceProfileId = new DeviceProfileId(edgeEvent.getEntityId());
DownlinkMsg downlinkMsg = null;
var msgConstructor = (DeviceMsgConstructor) deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion);
switch (edgeEvent.getAction()) {
case ADDED, UPDATED -> {
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(edgeEvent.getTenantId(), deviceProfileId);
if (deviceProfile != null) {
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction());
deviceProfile = checkIfDeviceProfileDefaultFieldsAssignedToEdge(edgeEvent.getTenantId(), edgeId, deviceProfile, edgeVersion);
DeviceProfileUpdateMsg deviceProfileUpdateMsg = ((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion)).constructDeviceProfileUpdatedMsg(msgType, deviceProfile);
DeviceProfileUpdateMsg deviceProfileUpdateMsg = msgConstructor.constructDeviceProfileUpdatedMsg(msgType, deviceProfile);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceProfileUpdateMsg(deviceProfileUpdateMsg)
@ -111,8 +111,7 @@ public abstract class DeviceProfileEdgeProcessor extends BaseDeviceProfileProces
}
}
case DELETED -> {
DeviceProfileUpdateMsg deviceProfileUpdateMsg = ((DeviceMsgConstructor)
deviceMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion)).constructDeviceProfileDeleteMsg(deviceProfileId);
DeviceProfileUpdateMsg deviceProfileUpdateMsg = msgConstructor.constructDeviceProfileDeleteMsg(deviceProfileId);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceProfileUpdateMsg(deviceProfileUpdateMsg)

94
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/oauth2/OAuth2EdgeProcessor.java

@ -15,49 +15,95 @@
*/
package org.thingsboard.server.service.edge.rpc.processor.oauth2;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.domain.DomainInfo;
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.TenantId;
import org.thingsboard.server.common.data.id.DomainId;
import org.thingsboard.server.common.data.id.OAuth2ClientId;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.gen.edge.v1.DownlinkMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2UpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.edge.v1.EdgeVersion;
import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor;
import org.thingsboard.server.service.edge.rpc.utils.EdgeVersionUtils;
@Slf4j
@Component
@TbCoreComponent
public class OAuth2EdgeProcessor extends BaseEdgeProcessor {
public DownlinkMsg convertOAuth2ProviderEventToDownlink(EdgeEvent edgeEvent) {
public DownlinkMsg convertOAuth2DomainEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) {
if (EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_7_1)) {
return null;
}
DomainId domainId = new DomainId(edgeEvent.getEntityId());
DownlinkMsg downlinkMsg = null;
// OAuth2Info oAuth2Info = JacksonUtil.convertValue(edgeEvent.getBody(), OAuth2Info.class);
// if (oAuth2Info != null && oAuth2Info.isEdgeEnabled()) {
// OAuth2UpdateMsg oAuth2UpdateMsg = oAuth2MsgConstructor.constructOAuth2UpdateMsg(oAuth2Info);
// downlinkMsg = DownlinkMsg.newBuilder()
// .setDownlinkMsgId(EdgeUtils.nextPositiveInt())
// .addOAuth2UpdateMsg(oAuth2ProviderUpdateMsg)
// .build();
// }
switch (edgeEvent.getAction()) {
case ADDED, UPDATED -> {
DomainInfo domainInfo = domainService.findDomainInfoById(edgeEvent.getTenantId(), domainId);
if (domainInfo != null && domainInfo.isPropagateToEdge()) {
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction());
OAuth2DomainUpdateMsg oAuth2DomainUpdateMsg = oAuth2MsgConstructor.constructOAuth2DomainUpdateMsg(msgType, domainInfo);
DownlinkMsg.Builder builder = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addOAuth2DomainUpdateMsg(oAuth2DomainUpdateMsg);
domainInfo.getOauth2ClientInfos().forEach(clientInfo -> {
OAuth2Client oauth2Client = oAuth2ClientService.findOAuth2ClientById(edgeEvent.getTenantId(), clientInfo.getId());
OAuth2ClientUpdateMsg oAuth2ClientUpdateMsg = oAuth2MsgConstructor.constructOAuth2ClientUpdateMsg(msgType, oauth2Client);
builder.addOAuth2ClientUpdateMsg(oAuth2ClientUpdateMsg);
});
downlinkMsg = builder.build();
}
}
case DELETED -> {
OAuth2DomainUpdateMsg oAuth2DomainUpdateMsg = oAuth2MsgConstructor.constructOAuth2DomainDeleteMsg(domainId);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addOAuth2DomainUpdateMsg(oAuth2DomainUpdateMsg)
.build();
}
}
return downlinkMsg;
}
public ListenableFuture<Void> processOAuth2Notification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) {
OAuth2Client oAuth2Info = JacksonUtil.fromString(edgeNotificationMsg.getBody(), OAuth2Client.class);
if (oAuth2Info == null) {
return Futures.immediateFuture(null);
public DownlinkMsg convertOAuth2ClientEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) {
if (EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_7_1)) {
return null;
}
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType());
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction());
return processActionForAllEdges(tenantId, type, actionType, null, JacksonUtil.toJsonNode(edgeNotificationMsg.getBody()), null);
OAuth2ClientId oAuth2ClientId = new OAuth2ClientId(edgeEvent.getEntityId());
DownlinkMsg downlinkMsg = null;
switch (edgeEvent.getAction()) {
case ADDED, UPDATED -> {
boolean isPropagateToEdge = oAuth2ClientService.isPropagateOAuth2ClientToEdge(edgeEvent.getTenantId(), oAuth2ClientId);
if (!isPropagateToEdge) {
return null;
}
OAuth2Client oAuth2Client = oAuth2ClientService.findOAuth2ClientById(edgeEvent.getTenantId(), oAuth2ClientId);
if (oAuth2Client != null) {
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction());
OAuth2ClientUpdateMsg oAuth2ClientUpdateMsg = oAuth2MsgConstructor.constructOAuth2ClientUpdateMsg(msgType, oAuth2Client);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addOAuth2ClientUpdateMsg(oAuth2ClientUpdateMsg)
.build();
}
}
case DELETED -> {
OAuth2ClientUpdateMsg oAuth2ClientDeleteMsg = oAuth2MsgConstructor.constructOAuth2ClientDeleteMsg(oAuth2ClientId);
downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addOAuth2ClientUpdateMsg(oAuth2ClientDeleteMsg)
.build();
}
}
return downlinkMsg;
}
}

8
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -614,13 +614,11 @@ public class DefaultTbClusterService implements TbClusterService {
private void pushDeviceUpdateMessage(TenantId tenantId, EdgeId edgeId, EntityId entityId, EdgeEventActionType action) {
log.trace("{} Going to send edge update notification for device actor, device id {}, edge id {}", tenantId, entityId, edgeId);
switch (action) {
case ASSIGNED_TO_EDGE:
pushMsgToCore(new DeviceEdgeUpdateMsg(tenantId, new DeviceId(entityId.getId()), edgeId), null);
break;
case UNASSIGNED_FROM_EDGE:
case ASSIGNED_TO_EDGE -> pushMsgToCore(new DeviceEdgeUpdateMsg(tenantId, new DeviceId(entityId.getId()), edgeId), null);
case UNASSIGNED_FROM_EDGE -> {
EdgeId relatedEdgeId = findRelatedEdgeIdIfAny(tenantId, entityId);
pushMsgToCore(new DeviceEdgeUpdateMsg(tenantId, new DeviceId(entityId.getId()), relatedEdgeId), null);
break;
}
}
}

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

@ -108,7 +108,6 @@ import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpd
import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate;
import java.util.List;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@ -701,13 +700,6 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
actorContext.tell(new TransportToDeviceActorMsgWrapper(toDeviceActorMsg, callback));
}
private void forwardToAppActor(UUID id, Optional<TbActorMsg> actorMsg, TbCallback callback) {
if (actorMsg.isPresent()) {
forwardToAppActor(id, actorMsg.get());
}
callback.onSuccess();
}
private void forwardToAppActor(UUID id, TbActorMsg actorMsg) {
log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg);
actorContext.tell(actorMsg);

6
application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java

@ -81,7 +81,8 @@ import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DeviceProfileUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.v1.EdgeVersion;
import org.thingsboard.server.gen.edge.v1.OAuth2UpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg;
import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg;
import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg;
import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg;
@ -892,7 +893,8 @@ public class EdgeControllerTest extends AbstractControllerTest {
EdgeImitator edgeImitator = new EdgeImitator(EDGE_HOST, EDGE_PORT, edge.getRoutingKey(), edge.getSecret());
edgeImitator.ignoreType(UserCredentialsUpdateMsg.class);
edgeImitator.ignoreType(OAuth2UpdateMsg.class);
edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class);
edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class);
edgeImitator.expectMessageAmount(27);
edgeImitator.connect();

19
application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java

@ -63,7 +63,6 @@ import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.ota.ChecksumAlgorithm;
import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.data.page.PageData;
@ -87,7 +86,8 @@ import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DeviceProfileUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.v1.EdgeConfiguration;
import org.thingsboard.server.gen.edge.v1.OAuth2UpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg;
import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg;
import org.thingsboard.server.gen.edge.v1.RuleChainMetadataRequestMsg;
import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg;
@ -144,7 +144,8 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
installation();
edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret());
edgeImitator.ignoreType(OAuth2UpdateMsg.class);
edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class);
edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class);
edgeImitator.expectMessageAmount(24);
edgeImitator.connect();
@ -547,18 +548,6 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
Assert.assertTrue(customer.isPublic());
}
private void validateOAuth2() throws Exception {
Optional<OAuth2UpdateMsg> oAuth2UpdateMsgOpt = edgeImitator.findMessageByType(OAuth2UpdateMsg.class);
Assert.assertTrue(oAuth2UpdateMsgOpt.isPresent());
OAuth2UpdateMsg oAuth2ProviderUpdateMsg = oAuth2UpdateMsgOpt.get();
OAuth2Client oAuth2Client = JacksonUtil.fromString(oAuth2ProviderUpdateMsg.getEntity(), OAuth2Client.class, true);
Assert.assertNotNull(oAuth2Client);
OAuth2Client auth2Info = doGet("/api/oauth2/config", OAuth2Client.class);
Assert.assertNotNull(auth2Info);
Assert.assertEquals(oAuth2Client, auth2Info);
testAutoGeneratedCodeByProtobuf(oAuth2ProviderUpdateMsg);
}
private void validateSyncCompleted() {
Optional<SyncCompletedMsg> syncCompletedMsgOpt = edgeImitator.findMessageByType(SyncCompletedMsg.class);
Assert.assertTrue(syncCompletedMsgOpt.isPresent());

144
application/src/test/java/org/thingsboard/server/edge/OAuth2EdgeTest.java

@ -20,55 +20,128 @@ import org.junit.Assert;
import org.junit.Ignore;
import org.junit.Test;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.domain.Domain;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.oauth2.MapperType;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.oauth2.OAuth2CustomMapperConfig;
import org.thingsboard.server.common.data.oauth2.OAuth2MapperConfig;
import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.edge.v1.OAuth2UpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import java.util.Arrays;
import java.util.Collections;
import java.util.Optional;
import java.util.UUID;
@DaoSqlTest
@Ignore("Ignored till fixed")
public class OAuth2EdgeTest extends AbstractEdgeTest {
@Test
public void testOAuth2Support() throws Exception {
public void testOAuth2DomainSupport() throws Exception {
loginSysAdmin();
// enable oauth
// enable oauth and save domain
edgeImitator.allowIgnoredTypes();
edgeImitator.expectMessageAmount(1);
OAuth2Client oAuth2Client = createDefaultOAuth2Info();
oAuth2Client = doPost("/api/oauth2/client", oAuth2Client, OAuth2Client.class);
Domain savedDomain = doPost("/api/domain", constructDomain(), Domain.class);
Assert.assertTrue(edgeImitator.waitForMessages());
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof OAuth2UpdateMsg);
OAuth2UpdateMsg oAuth2ProviderUpdateMsg = (OAuth2UpdateMsg) latestMessage;
OAuth2Client result = JacksonUtil.fromString(oAuth2ProviderUpdateMsg.getEntity(), OAuth2Client.class, true);
Assert.assertEquals(oAuth2Client, result);
Assert.assertTrue(latestMessage instanceof OAuth2DomainUpdateMsg);
OAuth2DomainUpdateMsg oAuth2DomainUpdateMsg = (OAuth2DomainUpdateMsg) latestMessage;
Domain result = JacksonUtil.fromString(oAuth2DomainUpdateMsg.getEntity(), Domain.class, true);
Assert.assertEquals(savedDomain, result);
// disable oauth support: no update of domain events is sending to Edge
edgeImitator.expectMessageAmount(1);
savedDomain.setPropagateToEdge(false);
doPost("/api/domain", savedDomain, Domain.class);
Assert.assertFalse(edgeImitator.waitForMessages(5));
// disable oauth support
// delete domain
edgeImitator.expectMessageAmount(1);
// oAuth2Info.setEnabled(false);
// oAuth2Info.setEdgeEnabled(false);
// doPost("/api/oauth2/config", oAuth2Info, OAuth2Info.class);
// Assert.assertFalse(edgeImitator.waitForMessages(5));
doDelete("/api/domain/" + savedDomain.getId().getId());
Assert.assertTrue(edgeImitator.waitForMessages());
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof OAuth2DomainUpdateMsg);
oAuth2DomainUpdateMsg = (OAuth2DomainUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, oAuth2DomainUpdateMsg.getMsgType());
Assert.assertEquals(savedDomain.getUuidId().getMostSignificantBits(), oAuth2DomainUpdateMsg.getIdMSB());
Assert.assertEquals(savedDomain.getUuidId().getLeastSignificantBits(), oAuth2DomainUpdateMsg.getIdLSB());
edgeImitator.ignoreType(OAuth2UpdateMsg.class);
edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class);
edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class);
loginTenantAdmin();
}
private OAuth2Client createDefaultOAuth2Info() {
return validRegistrationInfo();
@Test
public void testOAuth2ClientSupport() throws Exception {
loginSysAdmin();
// enable oauth and save domain
edgeImitator.allowIgnoredTypes();
edgeImitator.expectMessageAmount(2);
OAuth2Client savedOAuth2Client = validClientInfo(TenantId.SYS_TENANT_ID, "test edge google client");
savedOAuth2Client = doPost("/api/oauth2/client", savedOAuth2Client, OAuth2Client.class);
Domain savedDomain = doPost("/api/domain?oauth2ClientIds=" + savedOAuth2Client.getId().getId(), constructDomain(), Domain.class);
Assert.assertTrue(edgeImitator.waitForMessages());
Optional<OAuth2DomainUpdateMsg> oAuth2DomainUpdateMsgOpt = edgeImitator.findMessageByType(OAuth2DomainUpdateMsg.class);
Assert.assertTrue(oAuth2DomainUpdateMsgOpt.isPresent());
Domain result = JacksonUtil.fromString(oAuth2DomainUpdateMsgOpt.get().getEntity(), Domain.class, true);
Assert.assertEquals(savedDomain, result);
Optional<OAuth2ClientUpdateMsg> oAuth2ClientUpdateMsgOpt = edgeImitator.findMessageByType(OAuth2ClientUpdateMsg.class);
Assert.assertTrue(oAuth2ClientUpdateMsgOpt.isPresent());
OAuth2Client clientResult = JacksonUtil.fromString(oAuth2ClientUpdateMsgOpt.get().getEntity(), OAuth2Client.class, true);
Assert.assertEquals(savedOAuth2Client, clientResult);
// disable oauth support: no update of domain events and client events are sending to Edge
edgeImitator.expectMessageAmount(1);
savedDomain.setPropagateToEdge(false);
doPost("/api/domain", savedDomain, Domain.class);
Assert.assertFalse(edgeImitator.waitForMessages(5));
edgeImitator.expectMessageAmount(1);
savedOAuth2Client.setTitle("Updated title");
doPost("/api/oauth2/client", savedOAuth2Client, OAuth2Client.class);
Assert.assertFalse(edgeImitator.waitForMessages(5));
// delete oauth2Client
edgeImitator.expectMessageAmount(1);
doDelete("/api/oauth2/client/" + savedOAuth2Client.getId().getId());
Assert.assertTrue(edgeImitator.waitForMessages());
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof OAuth2ClientUpdateMsg);
OAuth2ClientUpdateMsg oAuth2ClientUpdateMsg = (OAuth2ClientUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, oAuth2ClientUpdateMsg.getMsgType());
Assert.assertEquals(savedOAuth2Client.getUuidId().getMostSignificantBits(), oAuth2ClientUpdateMsg.getIdMSB());
Assert.assertEquals(savedOAuth2Client.getUuidId().getLeastSignificantBits(), oAuth2ClientUpdateMsg.getIdLSB());
// delete domain
edgeImitator.expectMessageAmount(1);
doDelete("/api/domain/" + savedDomain.getId().getId());
Assert.assertTrue(edgeImitator.waitForMessages());
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof OAuth2DomainUpdateMsg);
OAuth2DomainUpdateMsg oAuth2DomainUpdateMsg = (OAuth2DomainUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, oAuth2DomainUpdateMsg.getMsgType());
Assert.assertEquals(savedDomain.getUuidId().getMostSignificantBits(), oAuth2DomainUpdateMsg.getIdMSB());
Assert.assertEquals(savedDomain.getUuidId().getLeastSignificantBits(), oAuth2DomainUpdateMsg.getIdLSB());
edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class);
edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class);
loginTenantAdmin();
}
private OAuth2Client validRegistrationInfo() {
private OAuth2Client validClientInfo(TenantId tenantId, String title) {
OAuth2Client oAuth2Client = new OAuth2Client();
oAuth2Client.setTitle(UUID.randomUUID().toString());
oAuth2Client.setTenantId(tenantId);
oAuth2Client.setTitle(title);
oAuth2Client.setClientId(UUID.randomUUID().toString());
oAuth2Client.setClientSecret(UUID.randomUUID().toString());
oAuth2Client.setAuthorizationUri(UUID.randomUUID().toString());
@ -80,16 +153,29 @@ public class OAuth2EdgeTest extends AbstractEdgeTest {
oAuth2Client.setJwkSetUri(UUID.randomUUID().toString());
oAuth2Client.setClientAuthenticationMethod(UUID.randomUUID().toString());
oAuth2Client.setLoginButtonLabel(UUID.randomUUID().toString());
oAuth2Client.setLoginButtonIcon(UUID.randomUUID().toString());
oAuth2Client.setAdditionalInfo(JacksonUtil.newObjectNode().put(UUID.randomUUID().toString(), UUID.randomUUID().toString()));
oAuth2Client.setMapperConfig(
OAuth2MapperConfig.builder()
.type(MapperType.CUSTOM)
.custom(
OAuth2CustomMapperConfig.builder()
.url(UUID.randomUUID().toString())
.build()
)
.build());
OAuth2MapperConfig.builder()
.allowUserCreation(true)
.activateUser(true)
.type(MapperType.CUSTOM)
.custom(
OAuth2CustomMapperConfig.builder()
.url(UUID.randomUUID().toString())
.build()
)
.build());
return oAuth2Client;
}
private Domain constructDomain() {
Domain domain = new Domain();
domain.setTenantId(TenantId.SYS_TENANT_ID);
domain.setName("my.edge.domain");
domain.setOauth2Enabled(true);
domain.setPropagateToEdge(true);
return domain;
}
}

14
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -47,7 +47,8 @@ import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg;
import org.thingsboard.server.gen.edge.v1.NotificationRuleUpdateMsg;
import org.thingsboard.server.gen.edge.v1.NotificationTargetUpdateMsg;
import org.thingsboard.server.gen.edge.v1.NotificationTemplateUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2UpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg;
import org.thingsboard.server.gen.edge.v1.OtaPackageUpdateMsg;
import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg;
import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg;
@ -318,9 +319,14 @@ public class EdgeImitator {
result.add(saveDownlinkMsg(resourceUpdateMsg));
}
}
if (downlinkMsg.getOAuth2UpdateMsgCount() > 0) {
for (OAuth2UpdateMsg oAuth2UpdateMsg : downlinkMsg.getOAuth2UpdateMsgList()) {
result.add(saveDownlinkMsg(oAuth2UpdateMsg));
if (downlinkMsg.getOAuth2ClientUpdateMsgCount() > 0) {
for (OAuth2ClientUpdateMsg oAuth2ClientUpdateMsg : downlinkMsg.getOAuth2ClientUpdateMsgList()) {
result.add(saveDownlinkMsg(oAuth2ClientUpdateMsg));
}
}
if (downlinkMsg.getOAuth2DomainUpdateMsgCount() > 0) {
for (OAuth2DomainUpdateMsg oAuth2DomainUpdateMsg : downlinkMsg.getOAuth2DomainUpdateMsgList()) {
result.add(saveDownlinkMsg(oAuth2DomainUpdateMsg));
}
}
if (downlinkMsg.getNotificationTemplateUpdateMsgCount() > 0) {

3
common/dao-api/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientService.java

@ -49,4 +49,7 @@ public interface OAuth2ClientService extends EntityDaoService {
PageData<OAuth2ClientInfo> findOAuth2ClientInfosByTenantId(TenantId tenantId, PageLink pageLink);
List<OAuth2ClientInfo> findOAuth2ClientInfosByIds(TenantId tenantId, List<OAuth2ClientId> oAuth2ClientIds);
boolean isPropagateOAuth2ClientToEdge(TenantId tenantId, OAuth2ClientId oAuth2ClientId);
}

3
common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java

@ -45,7 +45,8 @@ public enum EdgeEventType {
NOTIFICATION_TARGET (true, EntityType.NOTIFICATION_TARGET),
NOTIFICATION_TEMPLATE (true, EntityType.NOTIFICATION_TEMPLATE),
TB_RESOURCE(true, EntityType.TB_RESOURCE),
OAUTH2_CLIENT(true, EntityType.OAUTH2_CLIENT);
OAUTH2_CLIENT(true, EntityType.OAUTH2_CLIENT),
DOMAIN(true, EntityType.DOMAIN);
private final boolean allEdgesRelated;

4
common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java

@ -159,6 +159,10 @@ public class EntityIdFactory {
return new NotificationTargetId(uuid);
case NOTIFICATION_TEMPLATE:
return new NotificationTemplateId(uuid);
case OAUTH2_CLIENT:
return new OAuth2ClientId(uuid);
case DOMAIN:
return new DomainId(uuid);
}
throw new IllegalArgumentException("EdgeEventType " + edgeEventType + " is not supported!");
}

2
common/edge-api/src/main/java/org/thingsboard/edge/rpc/EdgeGrpcClient.java

@ -136,7 +136,7 @@ public class EdgeGrpcClient implements EdgeRpcClient {
.setConnectRequestMsg(ConnectRequestMsg.newBuilder()
.setEdgeRoutingKey(edgeKey)
.setEdgeSecret(edgeSecret)
.setEdgeVersion(EdgeVersion.V_3_7_0)
.setEdgeVersion(EdgeVersion.V_3_7_1)
.setMaxInboundMessageSize(maxInboundMessageSize)
.build())
.build());

19
common/edge-api/src/main/proto/edge.proto

@ -39,6 +39,7 @@ enum EdgeVersion {
V_3_6_2 = 5;
V_3_6_4 = 6;
V_3_7_0 = 7;
V_3_7_1 = 8;
}
/**
@ -469,8 +470,18 @@ message ResourceUpdateMsg {
string entity = 11;
}
message OAuth2UpdateMsg {
string entity = 1;
message OAuth2ClientUpdateMsg {
UpdateMsgType msgType = 1;
optional int64 idMSB = 2;
optional int64 idLSB = 3;
optional string entity = 4;
}
message OAuth2DomainUpdateMsg {
UpdateMsgType msgType = 1;
optional int64 idMSB = 2;
optional int64 idLSB = 3;
optional string entity = 4;
}
message NotificationRuleUpdateMsg {
@ -695,9 +706,9 @@ message DownlinkMsg {
repeated TenantProfileUpdateMsg tenantProfileUpdateMsg = 27;
repeated ResourceUpdateMsg resourceUpdateMsg = 28;
repeated AlarmCommentUpdateMsg alarmCommentUpdateMsg = 29;
repeated OAuth2UpdateMsg oAuth2UpdateMsg = 30;
repeated OAuth2ClientUpdateMsg oAuth2ClientUpdateMsg = 30;
repeated NotificationRuleUpdateMsg notificationRuleUpdateMsg = 31;
repeated NotificationTargetUpdateMsg notificationTargetUpdateMsg = 32;
repeated NotificationTemplateUpdateMsg notificationTemplateUpdateMsg = 33;
repeated OAuth2DomainUpdateMsg oAuth2DomainUpdateMsg = 34;
}

4
dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java

@ -58,7 +58,7 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe
log.trace("Executing saveDomain [{}]", domain);
try {
Domain savedDomain = domainDao.save(tenantId, domain);
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entity(savedDomain).build());
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entityId(savedDomain.getId()).entity(savedDomain).build());
return savedDomain;
} catch (Exception e) {
checkConstraintViolation(e,
@ -86,8 +86,6 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe
for (DomainOauth2Client client : newClientList) {
domainDao.addOauth2Client(client);
}
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId)
.entityId(domainId).created(false).build());
}
@Override

3
dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientDao.java

@ -43,4 +43,7 @@ public interface OAuth2ClientDao extends Dao<OAuth2Client> {
void deleteByTenantId(UUID tenantId);
List<OAuth2Client> findByIds(UUID tenantId, List<OAuth2ClientId> oAuth2ClientIds);
boolean isPropagateToEdge(TenantId tenantId, UUID oAuth2ClientId);
}

8
dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java

@ -73,7 +73,7 @@ public class OAuth2ClientServiceImpl extends AbstractEntityService implements OA
log.trace("Executing saveOAuth2Client [{}]", oAuth2Client);
oAuth2ClientDataValidator.validate(oAuth2Client, OAuth2Client::getTenantId);
OAuth2Client savedOauth2Client = oauth2ClientDao.save(tenantId, oAuth2Client);
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(TenantId.SYS_TENANT_ID).entity(oAuth2Client).build());
eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entityId(savedOauth2Client.getId()).entity(savedOauth2Client).build());
return savedOauth2Client;
}
@ -129,6 +129,12 @@ public class OAuth2ClientServiceImpl extends AbstractEntityService implements OA
.collect(Collectors.toList());
}
@Override
public boolean isPropagateOAuth2ClientToEdge(TenantId tenantId, OAuth2ClientId oAuth2ClientId) {
log.trace("Executing isPropagateOAuth2ClientToEdge, tenantId [{}], oAuth2ClientId [{}]", tenantId, oAuth2ClientId);
return oauth2ClientDao.isPropagateToEdge(tenantId, oAuth2ClientId.getId());
}
@Override
public void deleteByTenantId(TenantId tenantId) {
deleteOauth2ClientsByTenantId(tenantId);

2
dao/src/main/java/org/thingsboard/server/dao/sql/domain/JpaDomainDao.java

@ -67,7 +67,7 @@ public class JpaDomainDao extends JpaAbstractDao<DomainEntity, Domain> implement
@Override
public List<DomainOauth2Client> findOauth2ClientsByDomainId(TenantId tenantId, DomainId domainId) {
return DaoUtil.convertDataList(domainOauth2ClientRepository.findAllByDomainId(domainId.getId()));
return DaoUtil.convertDataList(domainOauth2ClientRepository.findAllByDomainId(domainId.getId()));
}
@Override

5
dao/src/main/java/org/thingsboard/server/dao/sql/oauth2/JpaOAuth2ClientDao.java

@ -94,6 +94,11 @@ public class JpaOAuth2ClientDao extends JpaAbstractDao<OAuth2ClientEntity, OAuth
return DaoUtil.convertDataList(repository.findByTenantIdAndIdIn(tenantId, toUUIDs(oAuth2ClientIds)));
}
@Override
public boolean isPropagateToEdge(TenantId tenantId, UUID oAuth2ClientId) {
return repository.isPropagateToEdge(tenantId.getId(), oAuth2ClientId);
}
@Override
public EntityType getEntityType() {
return EntityType.OAUTH2_CLIENT;

5
dao/src/main/java/org/thingsboard/server/dao/sql/oauth2/OAuth2ClientRepository.java

@ -81,4 +81,9 @@ public interface OAuth2ClientRepository extends JpaRepository<OAuth2ClientEntity
List<OAuth2ClientEntity> findByTenantIdAndIdIn(UUID tenantId, List<UUID> uuids);
@Query("SELECT COUNT(d) > 0 FROM DomainEntity d " +
"JOIN DomainOauth2ClientEntity doc ON d.id = doc.domainId " +
"WHERE d.tenantId = :tenantId AND doc.oauth2ClientId = :oAuth2ClientId AND d.propagateToEdge = true")
boolean isPropagateToEdge(@Param("tenantId") UUID tenantId, @Param("oAuth2ClientId") UUID oAuth2ClientId);
}

Loading…
Cancel
Save