Browse Source

Save attributes strategies: ensure device state service is notified when inactivity timeout is updated

pull/12764/head
Dmytro Skarzhynets 2 years ago
parent
commit
4495a2fa4b
No known key found for this signature in database GPG Key ID: 2B51652F224037DF
  1. 61
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java
  2. 83
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  3. 299
      application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java

61
application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java

@ -48,8 +48,6 @@ import org.thingsboard.server.queue.discovery.event.OtherServiceShutdownEvent;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.state.DefaultDeviceStateService;
import org.thingsboard.server.service.state.DeviceStateService;
import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate;
import org.thingsboard.server.service.ws.notification.sub.NotificationsSubscriptionUpdate;
@ -75,7 +73,6 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
private final TbServiceInfoProvider serviceInfoProvider;
private final TbQueueProducerProvider producerProvider;
private final TbLocalSubscriptionService localSubscriptionService;
private final DeviceStateService deviceStateService;
private final TbClusterService clusterService;
private final SubscriptionSchedulerComponent scheduler;
@ -170,7 +167,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
callback.onSuccess();
}
public void onTimeSeriesUpdate(EntityId entityId, List<TsKvEntry> update) {
private void onTimeSeriesUpdate(EntityId entityId, List<TsKvEntry> update) {
getEntityUpdatesInfo(entityId).timeSeriesUpdateTs = System.currentTimeMillis();
TbEntityRemoteSubsInfo subInfo = entitySubscriptions.get(entityId);
if (subInfo != null) {
@ -202,11 +199,6 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
public void onAttributesUpdate(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, TbCallback callback) {
getEntityUpdatesInfo(entityId).attributesUpdateTs = System.currentTimeMillis();
processAttributesUpdate(entityId, scope, attributes);
if (entityId.getEntityType() == EntityType.DEVICE) {
if (TbAttributeSubscriptionScope.SERVER_SCOPE.name().equalsIgnoreCase(scope)) {
updateDeviceInactivityTimeout(tenantId, entityId, attributes);
}
}
callback.onSuccess();
}
@ -219,19 +211,13 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
public void onAttributesDelete(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, TbCallback callback) {
processAttributesUpdate(entityId, scope,
keys.stream().map(key -> new BaseAttributeKvEntry(0, new StringDataEntry(key, ""))).collect(Collectors.toList()));
if (entityId.getEntityType() == EntityType.DEVICE) {
if (TbAttributeSubscriptionScope.SERVER_SCOPE.name().equalsIgnoreCase(scope)
|| TbAttributeSubscriptionScope.ANY_SCOPE.name().equalsIgnoreCase(scope)) {
deleteDeviceInactivityTimeout(tenantId, entityId, keys);
} else if (TbAttributeSubscriptionScope.SHARED_SCOPE.name().equalsIgnoreCase(scope) && notifyDevice) {
clusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete(tenantId,
new DeviceId(entityId.getId()), scope, keys), null);
}
if (entityId.getEntityType() == EntityType.DEVICE && TbAttributeSubscriptionScope.SHARED_SCOPE.name().equalsIgnoreCase(scope) && notifyDevice) {
clusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete(tenantId, new DeviceId(entityId.getId()), scope, keys), null);
}
callback.onSuccess();
}
public void processAttributesUpdate(EntityId entityId, String scope, List<AttributeKvEntry> update) {
private void processAttributesUpdate(EntityId entityId, String scope, List<AttributeKvEntry> update) {
TbEntityRemoteSubsInfo subInfo = entitySubscriptions.get(entityId);
if (subInfo != null) {
log.trace("[{}] Handling attributes update: {}", entityId, update);
@ -259,22 +245,6 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
}
}
private void updateDeviceInactivityTimeout(TenantId tenantId, EntityId entityId, List<? extends KvEntry> kvEntries) {
for (KvEntry kvEntry : kvEntries) {
if (kvEntry.getKey().equals(DefaultDeviceStateService.INACTIVITY_TIMEOUT)) {
deviceStateService.onDeviceInactivityTimeoutUpdate(tenantId, new DeviceId(entityId.getId()), getLongValue(kvEntry));
}
}
}
private void deleteDeviceInactivityTimeout(TenantId tenantId, EntityId entityId, List<String> keys) {
for (String key : keys) {
if (key.equals(DefaultDeviceStateService.INACTIVITY_TIMEOUT)) {
deviceStateService.onDeviceInactivityTimeoutUpdate(tenantId, new DeviceId(entityId.getId()), 0);
}
}
}
@Override
public void onAlarmUpdate(TenantId tenantId, EntityId entityId, AlarmInfo alarm, TbCallback callback) {
onAlarmSubUpdate(tenantId, entityId, alarm, false, callback);
@ -344,29 +314,6 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene
}
}
private static long getLongValue(KvEntry kve) {
switch (kve.getDataType()) {
case LONG:
return kve.getLongValue().orElse(0L);
case DOUBLE:
return kve.getDoubleValue().orElse(0.0).longValue();
case STRING:
try {
return Long.parseLong(kve.getStrValue().orElse("0"));
} catch (NumberFormatException e) {
return 0L;
}
case JSON:
try {
return Long.parseLong(kve.getJsonValue().orElse("0"));
} catch (NumberFormatException e) {
return 0L;
}
default:
return 0L;
}
}
private static <T extends KvEntry> List<T> getSubList(List<T> ts, Set<String> keys) {
List<T> update = null;
for (T entry : ts) {

83
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -31,10 +31,12 @@ import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.rule.engine.api.AttributesDeleteRequest;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.TimeseriesDeleteRequest;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView;
@ -43,6 +45,7 @@ import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult;
@ -55,12 +58,11 @@ import org.thingsboard.server.dao.util.KvUtils;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.cf.CalculatedFieldQueueService;
import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService;
import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope;
import org.thingsboard.server.service.state.DefaultDeviceStateService;
import org.thingsboard.server.service.subscription.TbSubscriptionUtils;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@ -70,6 +72,11 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.function.Consumer;
import static java.util.Comparator.comparing;
import static java.util.Comparator.comparingLong;
import static java.util.Comparator.naturalOrder;
import static java.util.Comparator.nullsFirst;
/**
* Created by ashvayka on 27.03.18.
*/
@ -83,6 +90,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
private final TbApiUsageReportClient apiUsageClient;
private final TbApiUsageStateService apiUsageStateService;
private final CalculatedFieldQueueService calculatedFieldQueueService;
private final DeviceStateManager deviceStateManager;
private ExecutorService tsCallBackExecutor;
@ -94,13 +102,15 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
@Lazy TbEntityViewService tbEntityViewService,
TbApiUsageReportClient apiUsageClient,
TbApiUsageStateService apiUsageStateService,
CalculatedFieldQueueService calculatedFieldQueueService) {
CalculatedFieldQueueService calculatedFieldQueueService,
DeviceStateManager deviceStateManager) {
this.attrService = attrService;
this.tsService = tsService;
this.tbEntityViewService = tbEntityViewService;
this.apiUsageClient = apiUsageClient;
this.apiUsageStateService = apiUsageStateService;
this.calculatedFieldQueueService = calculatedFieldQueueService;
this.deviceStateManager = deviceStateManager;
}
@PostConstruct
@ -201,20 +211,51 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
}
}, t -> request.getCallback().onFailure(t));
if (strategy.saveAttributes()
&& entityId.getEntityType() == EntityType.DEVICE
&& TbAttributeSubscriptionScope.SHARED_SCOPE.name().equalsIgnoreCase(request.getScope().name())
&& request.isNotifyDevice()) {
if (shouldSendSharedAttributesUpdatedNotification(request)) {
addMainCallback(resultFuture, success -> clusterService.pushMsgToCore(
DeviceAttributesEventNotificationMsg.onUpdate(tenantId, new DeviceId(entityId.getId()), DataConstants.SHARED_SCOPE, request.getEntries()), null
));
}
if (shouldCheckForInactivityTimeoutUpdates(request)) {
findNewInactivityTimeout(request.getEntries()).ifPresent(newInactivityTimeout ->
addMainCallback(resultFuture, success -> deviceStateManager.onDeviceInactivityTimeoutUpdate(
tenantId, new DeviceId(entityId.getId()), newInactivityTimeout, TbCallback.EMPTY)
)
);
}
if (strategy.sendWsUpdate()) {
addWsCallback(resultFuture, success -> onAttributesUpdate(tenantId, entityId, request.getScope().name(), request.getEntries()));
}
}
private static boolean shouldSendSharedAttributesUpdatedNotification(AttributesSaveRequest request) {
return request.getStrategy().saveAttributes() && shouldSendSharedAttributesNotification(request.getEntityId(), request.getScope(), request.isNotifyDevice());
}
private static boolean shouldCheckForInactivityTimeoutUpdates(AttributesSaveRequest request) {
return request.getStrategy().saveAttributes()
&& request.getEntityId().getEntityType() == EntityType.DEVICE
&& request.getScope() == AttributeScope.SERVER_SCOPE;
}
private static Optional<Long> findNewInactivityTimeout(List<AttributeKvEntry> entries) {
return entries.stream()
.filter(entry -> Objects.equals(DefaultDeviceStateService.INACTIVITY_TIMEOUT, entry.getKey()))
// Select the entry with the highest version, or if the versions are equal, the one with the most recent update timestamp
.max(comparing(AttributeKvEntry::getVersion, nullsFirst(naturalOrder())).thenComparingLong(AttributeKvEntry::getLastUpdateTs))
.map(DefaultTelemetrySubscriptionService::parseAsLong);
}
private static long parseAsLong(KvEntry kve) {
try {
return Long.parseLong(kve.getValueAsString());
} catch (NumberFormatException e) {
return 0L;
}
}
@Override
public void deleteAttributes(AttributesDeleteRequest request) {
checkInternalEntity(request.getEntityId());
@ -233,17 +274,37 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
t -> request.getCallback().onFailure(t)
);
if (entityId.getEntityType() == EntityType.DEVICE
&& TbAttributeSubscriptionScope.SHARED_SCOPE.name().equalsIgnoreCase(request.getScope().name())
&& request.isNotifyDevice()) {
if (shouldSendSharedAttributesDeletedNotification(request)) {
addMainCallback(deleteFuture, success -> clusterService.pushMsgToCore(
DeviceAttributesEventNotificationMsg.onDelete(tenantId, new DeviceId(entityId.getId()), DataConstants.SHARED_SCOPE, request.getKeys()), null
));
}
if (inactivityTimeoutDeleted(request)) {
addMainCallback(deleteFuture, success -> deviceStateManager.onDeviceInactivityTimeoutUpdate(
tenantId, new DeviceId(entityId.getId()), 0L, TbCallback.EMPTY)
);
}
addWsCallback(deleteFuture, success -> onAttributesDelete(tenantId, entityId, request.getScope().name(), request.getKeys()));
}
private static boolean shouldSendSharedAttributesDeletedNotification(AttributesDeleteRequest request) {
return shouldSendSharedAttributesNotification(request.getEntityId(), request.getScope(), request.isNotifyDevice());
}
private static boolean shouldSendSharedAttributesNotification(EntityId entityId, AttributeScope scope, boolean notifyDevice) {
return entityId.getEntityType() == EntityType.DEVICE
&& scope == AttributeScope.SHARED_SCOPE
&& notifyDevice;
}
private static boolean inactivityTimeoutDeleted(AttributesDeleteRequest request) {
return request.getEntityId().getEntityType() == EntityType.DEVICE
&& request.getScope() == AttributeScope.SERVER_SCOPE
&& request.getKeys().stream().anyMatch(key -> Objects.equals(DefaultDeviceStateService.INACTIVITY_TIMEOUT, key));
}
@Override
public void deleteTimeseries(TimeseriesDeleteRequest request) {
checkInternalEntity(request.getEntityId());
@ -293,7 +354,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
if (entries != null) {
Optional<TsKvEntry> tsKvEntry = entries.stream()
.filter(entry -> entry.getTs() > startTs && entry.getTs() <= endTs)
.max(Comparator.comparingLong(TsKvEntry::getTs));
.max(comparingLong(TsKvEntry::getTs));
tsKvEntry.ifPresent(entityViewLatest::add);
}
}

299
application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java

@ -31,6 +31,7 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.rule.engine.api.AttributesDeleteRequest;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.DeviceStateManager;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
@ -51,6 +52,7 @@ import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.DoubleDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult;
import org.thingsboard.server.common.data.kv.TsKvEntry;
@ -137,12 +139,14 @@ class DefaultTelemetrySubscriptionServiceTest {
TbApiUsageStateService apiUsageStateService;
@Mock
CalculatedFieldQueueService calculatedFieldQueueService;
@Mock
DeviceStateManager deviceStateManager;
DefaultTelemetrySubscriptionService telemetryService;
@BeforeEach
void setup() {
telemetryService = new DefaultTelemetrySubscriptionService(attrService, tsService, tbEntityViewService, apiUsageClient, apiUsageStateService, calculatedFieldQueueService);
telemetryService = new DefaultTelemetrySubscriptionService(attrService, tsService, tbEntityViewService, apiUsageClient, apiUsageStateService, calculatedFieldQueueService, deviceStateManager);
ReflectionTestUtils.setField(telemetryService, "clusterService", clusterService);
ReflectionTestUtils.setField(telemetryService, "partitionService", partitionService);
ReflectionTestUtils.setField(telemetryService, "subscriptionManagerService", Optional.of(subscriptionManagerService));
@ -658,6 +662,184 @@ class DefaultTelemetrySubscriptionServiceTest {
then(clusterService).should(never()).pushMsgToCore(any(), any());
}
@Test
void shouldNotifyDeviceStateManagerWhenDeviceInactivityTimeoutWasUpdated() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var inactivityTimeout = new BaseAttributeKvEntry(123L, new LongDataEntry("inactivityTimeout", 5000L));
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.entry(inactivityTimeout)
.strategy(new AttributesSaveRequest.Strategy(true, false, false))
.build();
given(attrService.save(tenantId, deviceId, request.getScope(), request.getEntries())).willReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).should().onDeviceInactivityTimeoutUpdate(tenantId, deviceId, 5000L, TbCallback.EMPTY);
}
@Test
void shouldNotNotifyDeviceStateManagerWhenDeviceInactivityTimeoutSaveWasSkipped() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var inactivityTimeout = new BaseAttributeKvEntry(123L, new LongDataEntry("inactivityTimeout", 5000L));
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.entry(inactivityTimeout)
.strategy(new AttributesSaveRequest.Strategy(false, true, true))
.build();
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@ParameterizedTest
@EnumSource(
value = EntityType.class,
names = {"DEVICE", "API_USAGE_STATE"}, // API usage state excluded due to coverage in another test
mode = EnumSource.Mode.EXCLUDE
)
void shouldNotNotifyDeviceStateManagerWhenInactivityTimeoutWasUpdatedButEntityTypeIsNotDevice(EntityType entityType) {
// GIVEN
var nonDeviceId = EntityIdFactory.getByTypeAndUuid(entityType, "cc51e450-53e1-11ee-883e-e56b48fd2088");
var inactivityTimeout = new BaseAttributeKvEntry(123L, new LongDataEntry("inactivityTimeout", 5000L));
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(nonDeviceId)
.scope(AttributeScope.SERVER_SCOPE)
.entry(inactivityTimeout)
.strategy(new AttributesSaveRequest.Strategy(true, false, false))
.build();
given(attrService.save(tenantId, nonDeviceId, request.getScope(), request.getEntries())).willReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@ParameterizedTest
@EnumSource(
value = AttributeScope.class,
names = {"SERVER_SCOPE"},
mode = EnumSource.Mode.EXCLUDE
)
void shouldNotNotifyDeviceStateManagerWhenInactivityTimeoutWasUpdatedButAttributeScopeIsNotServer(AttributeScope nonServerScope) {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var inactivityTimeout = new BaseAttributeKvEntry(123L, new LongDataEntry("inactivityTimeout", 5000L));
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(nonServerScope)
.entry(inactivityTimeout)
.strategy(new AttributesSaveRequest.Strategy(true, false, false))
.build();
given(attrService.save(tenantId, deviceId, request.getScope(), request.getEntries())).willReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@Test
void shouldNotNotifyDeviceStateManagerWhenUpdatedAttributesDoNotContainInactivityTimeout() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var inactivityTimeout = new BaseAttributeKvEntry(123L, new LongDataEntry("notInactivityTimeout", 5000L));
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.entry(inactivityTimeout)
.strategy(new AttributesSaveRequest.Strategy(true, false, false))
.build();
given(attrService.save(tenantId, deviceId, request.getScope(), request.getEntries())).willReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@Test
void shouldUseInactivityTimeoutEntryWithTheGreatestVersion() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
List<AttributeKvEntry> entries = List.of(
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 0L), 0L, null),
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 1000L), 3L, 1L),
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 2000L), 2L, 2L),
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 3000L), 1L, 3L)
);
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.entries(entries)
.strategy(new AttributesSaveRequest.Strategy(true, false, false))
.build();
given(attrService.save(tenantId, deviceId, request.getScope(), request.getEntries())).willReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).should().onDeviceInactivityTimeoutUpdate(tenantId, deviceId, 3000L, TbCallback.EMPTY);
}
@Test
void shouldUseInactivityTimeoutEntryWithTheGreatestLastUpdateTsWhenVersionsAreTheSame() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
List<AttributeKvEntry> entries = List.of(
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 1000L), 1L, 1L),
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 2000L), 2L, 1L),
new BaseAttributeKvEntry(new LongDataEntry("inactivityTimeout", 3000L), 3L, 1L)
);
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.entries(entries)
.strategy(new AttributesSaveRequest.Strategy(true, false, false))
.build();
given(attrService.save(tenantId, deviceId, request.getScope(), request.getEntries())).willReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
then(deviceStateManager).should().onDeviceInactivityTimeoutUpdate(tenantId, deviceId, 3000L, TbCallback.EMPTY);
}
/* --- Delete attributes API --- */
@Test
@ -807,6 +989,121 @@ class DefaultTelemetrySubscriptionServiceTest {
then(clusterService).should(never()).pushMsgToCore(any(), any());
}
@Test
void shouldNotifyDeviceStateManagerWhenDeviceInactivityTimeoutWasDeleted() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var request = AttributesDeleteRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.keys(List.of("inactivityTimeout", "someOtherDeletedAttribute"))
.build();
given(attrService.removeAll(tenantId, deviceId, request.getScope(), request.getKeys())).willReturn(immediateFuture(request.getKeys()));
// WHEN
telemetryService.deleteAttributes(request);
// THEN
then(deviceStateManager).should().onDeviceInactivityTimeoutUpdate(tenantId, deviceId, 0L, TbCallback.EMPTY);
}
@ParameterizedTest
@EnumSource(
value = EntityType.class,
names = {"DEVICE", "API_USAGE_STATE"}, // API usage state excluded due to coverage in another test
mode = EnumSource.Mode.EXCLUDE
)
void shouldNotNotifyDeviceStateManagerWhenInactivityTimeoutWasDeletedButEntityTypeIsNotDevice(EntityType entityType) {
// GIVEN
var nonDeviceId = EntityIdFactory.getByTypeAndUuid(entityType, "cc51e450-53e1-11ee-883e-e56b48fd2088");
var request = AttributesDeleteRequest.builder()
.tenantId(tenantId)
.entityId(nonDeviceId)
.scope(AttributeScope.SERVER_SCOPE)
.keys(List.of("inactivityTimeout", "someOtherDeletedAttribute"))
.build();
given(attrService.removeAll(tenantId, nonDeviceId, request.getScope(), request.getKeys())).willReturn(immediateFuture(request.getKeys()));
// WHEN
telemetryService.deleteAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@ParameterizedTest
@EnumSource(
value = AttributeScope.class,
names = {"SERVER_SCOPE"},
mode = EnumSource.Mode.EXCLUDE
)
void shouldNotNotifyDeviceStateManagerWhenInactivityTimeoutWasDeletedButAttributeScopeIsNotServer(AttributeScope nonServerScope) {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var request = AttributesDeleteRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(nonServerScope)
.keys(List.of("inactivityTimeout", "someOtherDeletedAttribute"))
.build();
given(attrService.removeAll(tenantId, deviceId, request.getScope(), request.getKeys())).willReturn(immediateFuture(request.getKeys()));
// WHEN
telemetryService.deleteAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@Test
void shouldNotNotifyDeviceStateManagerWhenInactivityTimeoutWasNotDeleted() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var request = AttributesDeleteRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.keys(List.of("someOtherDeletedAttribute"))
.build();
given(attrService.removeAll(tenantId, deviceId, request.getScope(), request.getKeys())).willReturn(immediateFuture(request.getKeys()));
// WHEN
telemetryService.deleteAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
@Test
void shouldNotNotifyDeviceStateManagerWhenDeviceInactivityTimeoutDeleteFailed() {
// GIVEN
var deviceId = DeviceId.fromString("cc51e450-53e1-11ee-883e-e56b48fd2088");
var request = AttributesDeleteRequest.builder()
.tenantId(tenantId)
.entityId(deviceId)
.scope(AttributeScope.SERVER_SCOPE)
.keys(List.of("inactivityTimeout", "someOtherDeletedAttribute"))
.build();
given(attrService.removeAll(tenantId, deviceId, request.getScope(), request.getKeys())).willReturn(immediateFailedFuture(new RuntimeException("failed to delete")));
// WHEN
telemetryService.deleteAttributes(request);
// THEN
then(deviceStateManager).shouldHaveNoInteractions();
}
// used to emulate versions returned by save APIs
private static List<Long> listOfNNumbers(int N) {
return LongStream.range(0, N).boxed().toList();

Loading…
Cancel
Save