committed by
GitHub
198 changed files with 2856 additions and 1046 deletions
@ -0,0 +1,315 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.subscription; |
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import lombok.Data; |
||||
|
import lombok.Getter; |
||||
|
import lombok.Setter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.id.UserId; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.query.ComplexFilterPredicate; |
||||
|
import org.thingsboard.server.common.data.query.DynamicValue; |
||||
|
import org.thingsboard.server.common.data.query.DynamicValueSourceType; |
||||
|
import org.thingsboard.server.common.data.query.EntityCountQuery; |
||||
|
import org.thingsboard.server.common.data.query.EntityKeyType; |
||||
|
import org.thingsboard.server.common.data.query.FilterPredicateType; |
||||
|
import org.thingsboard.server.common.data.query.KeyFilter; |
||||
|
import org.thingsboard.server.common.data.query.KeyFilterPredicate; |
||||
|
import org.thingsboard.server.common.data.query.SimpleKeyFilterPredicate; |
||||
|
import org.thingsboard.server.common.data.query.TsValue; |
||||
|
import org.thingsboard.server.dao.attributes.AttributesService; |
||||
|
import org.thingsboard.server.dao.entity.EntityService; |
||||
|
import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; |
||||
|
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; |
||||
|
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.Optional; |
||||
|
import java.util.Set; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
import java.util.concurrent.ScheduledFuture; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Data |
||||
|
public abstract class TbAbstractSubCtx<T extends EntityCountQuery> { |
||||
|
|
||||
|
protected final String serviceId; |
||||
|
protected final SubscriptionServiceStatistics stats; |
||||
|
protected final TelemetryWebSocketService wsService; |
||||
|
protected final EntityService entityService; |
||||
|
protected final TbLocalSubscriptionService localSubscriptionService; |
||||
|
protected final AttributesService attributesService; |
||||
|
protected final TelemetryWebSocketSessionRef sessionRef; |
||||
|
protected final int cmdId; |
||||
|
protected final Set<Integer> subToDynamicValueKeySet; |
||||
|
@Getter |
||||
|
protected final Map<DynamicValueKey, List<DynamicValue>> dynamicValues; |
||||
|
@Getter |
||||
|
@Setter |
||||
|
protected T query; |
||||
|
@Setter |
||||
|
protected volatile ScheduledFuture<?> refreshTask; |
||||
|
|
||||
|
public TbAbstractSubCtx(String serviceId, TelemetryWebSocketService wsService, |
||||
|
EntityService entityService, TbLocalSubscriptionService localSubscriptionService, |
||||
|
AttributesService attributesService, SubscriptionServiceStatistics stats, |
||||
|
TelemetryWebSocketSessionRef sessionRef, int cmdId) { |
||||
|
this.serviceId = serviceId; |
||||
|
this.wsService = wsService; |
||||
|
this.entityService = entityService; |
||||
|
this.localSubscriptionService = localSubscriptionService; |
||||
|
this.attributesService = attributesService; |
||||
|
this.stats = stats; |
||||
|
this.sessionRef = sessionRef; |
||||
|
this.cmdId = cmdId; |
||||
|
this.subToDynamicValueKeySet = ConcurrentHashMap.newKeySet(); |
||||
|
this.dynamicValues = new ConcurrentHashMap<>(); |
||||
|
} |
||||
|
|
||||
|
public void setAndResolveQuery(T query) { |
||||
|
dynamicValues.clear(); |
||||
|
this.query = query; |
||||
|
if (query != null && query.getKeyFilters() != null) { |
||||
|
for (KeyFilter filter : query.getKeyFilters()) { |
||||
|
registerDynamicValues(filter.getPredicate()); |
||||
|
} |
||||
|
} |
||||
|
resolve(getTenantId(), getCustomerId(), getUserId()); |
||||
|
} |
||||
|
|
||||
|
public void resolve(TenantId tenantId, CustomerId customerId, UserId userId) { |
||||
|
List<ListenableFuture<DynamicValueKeySub>> futures = new ArrayList<>(); |
||||
|
for (DynamicValueKey key : dynamicValues.keySet()) { |
||||
|
switch (key.getSourceType()) { |
||||
|
case CURRENT_TENANT: |
||||
|
futures.add(resolveEntityValue(tenantId, tenantId, key)); |
||||
|
break; |
||||
|
case CURRENT_CUSTOMER: |
||||
|
if (customerId != null && !customerId.isNullUid()) { |
||||
|
futures.add(resolveEntityValue(tenantId, customerId, key)); |
||||
|
} |
||||
|
break; |
||||
|
case CURRENT_USER: |
||||
|
if (userId != null && !userId.isNullUid()) { |
||||
|
futures.add(resolveEntityValue(tenantId, userId, key)); |
||||
|
} |
||||
|
break; |
||||
|
} |
||||
|
} |
||||
|
try { |
||||
|
Map<EntityId, Map<String, DynamicValueKeySub>> tmpSubMap = new HashMap<>(); |
||||
|
for (DynamicValueKeySub sub : Futures.successfulAsList(futures).get()) { |
||||
|
tmpSubMap.computeIfAbsent(sub.getEntityId(), tmp -> new HashMap<>()).put(sub.getKey().getSourceAttribute(), sub); |
||||
|
} |
||||
|
for (EntityId entityId : tmpSubMap.keySet()) { |
||||
|
Map<String, Long> keyStates = new HashMap<>(); |
||||
|
Map<String, DynamicValueKeySub> dynamicValueKeySubMap = tmpSubMap.get(entityId); |
||||
|
dynamicValueKeySubMap.forEach((k, v) -> keyStates.put(k, v.getLastUpdateTs())); |
||||
|
int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); |
||||
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
||||
|
.serviceId(serviceId) |
||||
|
.sessionId(sessionRef.getSessionId()) |
||||
|
.subscriptionId(subIdx) |
||||
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
||||
|
.entityId(entityId) |
||||
|
.updateConsumer((s, subscriptionUpdate) -> dynamicValueSubUpdate(s, subscriptionUpdate, dynamicValueKeySubMap)) |
||||
|
.allKeys(false) |
||||
|
.keyStates(keyStates) |
||||
|
.scope(TbAttributeSubscriptionScope.SERVER_SCOPE) |
||||
|
.build(); |
||||
|
subToDynamicValueKeySet.add(subIdx); |
||||
|
localSubscriptionService.addSubscription(sub); |
||||
|
} |
||||
|
} catch (InterruptedException | ExecutionException e) { |
||||
|
log.info("[{}][{}][{}] Failed to resolve dynamic values: {}", tenantId, customerId, userId, dynamicValues.keySet()); |
||||
|
} |
||||
|
|
||||
|
} |
||||
|
|
||||
|
private void dynamicValueSubUpdate(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, |
||||
|
Map<String, DynamicValueKeySub> dynamicValueKeySubMap) { |
||||
|
Map<String, TsValue> latestUpdate = new HashMap<>(); |
||||
|
subscriptionUpdate.getData().forEach((k, v) -> { |
||||
|
Object[] data = (Object[]) v.get(0); |
||||
|
latestUpdate.put(k, new TsValue((Long) data[0], (String) data[1])); |
||||
|
}); |
||||
|
|
||||
|
boolean invalidateFilter = false; |
||||
|
for (Map.Entry<String, TsValue> entry : latestUpdate.entrySet()) { |
||||
|
String k = entry.getKey(); |
||||
|
TsValue tsValue = entry.getValue(); |
||||
|
DynamicValueKeySub sub = dynamicValueKeySubMap.get(k); |
||||
|
if (sub.updateValue(tsValue)) { |
||||
|
invalidateFilter = true; |
||||
|
updateDynamicValuesByKey(sub, tsValue); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
if (invalidateFilter) { |
||||
|
update(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public abstract void fetchData(); |
||||
|
|
||||
|
protected abstract void update(); |
||||
|
|
||||
|
public void clearSubscriptions() { |
||||
|
clearDynamicValueSubscriptions(); |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
private static class DynamicValueKeySub { |
||||
|
private final DynamicValueKey key; |
||||
|
private final EntityId entityId; |
||||
|
private long lastUpdateTs; |
||||
|
private String lastUpdateValue; |
||||
|
|
||||
|
boolean updateValue(TsValue value) { |
||||
|
if (value.getTs() > lastUpdateTs && (lastUpdateValue == null || !lastUpdateValue.equals(value.getValue()))) { |
||||
|
this.lastUpdateTs = value.getTs(); |
||||
|
this.lastUpdateValue = value.getValue(); |
||||
|
return true; |
||||
|
} else { |
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private ListenableFuture<DynamicValueKeySub> resolveEntityValue(TenantId tenantId, EntityId entityId, DynamicValueKey key) { |
||||
|
ListenableFuture<Optional<AttributeKvEntry>> entry = attributesService.find(tenantId, entityId, |
||||
|
TbAttributeSubscriptionScope.SERVER_SCOPE.name(), key.getSourceAttribute()); |
||||
|
return Futures.transform(entry, attributeOpt -> { |
||||
|
DynamicValueKeySub sub = new DynamicValueKeySub(key, entityId); |
||||
|
if (attributeOpt.isPresent()) { |
||||
|
AttributeKvEntry attribute = attributeOpt.get(); |
||||
|
sub.setLastUpdateTs(attribute.getLastUpdateTs()); |
||||
|
sub.setLastUpdateValue(attribute.getValueAsString()); |
||||
|
updateDynamicValuesByKey(sub, new TsValue(attribute.getLastUpdateTs(), attribute.getValueAsString())); |
||||
|
} |
||||
|
return sub; |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
protected void updateDynamicValuesByKey(DynamicValueKeySub sub, TsValue tsValue) { |
||||
|
DynamicValueKey dvk = sub.getKey(); |
||||
|
switch (dvk.getPredicateType()) { |
||||
|
case STRING: |
||||
|
dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(tsValue.getValue())); |
||||
|
break; |
||||
|
case NUMERIC: |
||||
|
try { |
||||
|
Double dValue = Double.parseDouble(tsValue.getValue()); |
||||
|
dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(dValue)); |
||||
|
} catch (NumberFormatException e) { |
||||
|
dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(null)); |
||||
|
} |
||||
|
break; |
||||
|
case BOOLEAN: |
||||
|
Boolean bValue = Boolean.parseBoolean(tsValue.getValue()); |
||||
|
dynamicValues.get(dvk).forEach(dynamicValue -> dynamicValue.setResolvedValue(bValue)); |
||||
|
break; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@SuppressWarnings("unchecked") |
||||
|
private void registerDynamicValues(KeyFilterPredicate predicate) { |
||||
|
switch (predicate.getType()) { |
||||
|
case STRING: |
||||
|
case NUMERIC: |
||||
|
case BOOLEAN: |
||||
|
Optional<DynamicValue> value = getDynamicValueFromSimplePredicate((SimpleKeyFilterPredicate) predicate); |
||||
|
if (value.isPresent()) { |
||||
|
DynamicValue dynamicValue = value.get(); |
||||
|
DynamicValueKey key = new DynamicValueKey( |
||||
|
predicate.getType(), |
||||
|
dynamicValue.getSourceType(), |
||||
|
dynamicValue.getSourceAttribute()); |
||||
|
dynamicValues.computeIfAbsent(key, tmp -> new ArrayList<>()).add(dynamicValue); |
||||
|
} |
||||
|
break; |
||||
|
case COMPLEX: |
||||
|
((ComplexFilterPredicate) predicate).getPredicates().forEach(this::registerDynamicValues); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private Optional<DynamicValue<T>> getDynamicValueFromSimplePredicate(SimpleKeyFilterPredicate<T> predicate) { |
||||
|
if (predicate.getValue().getUserValue() == null) { |
||||
|
return Optional.ofNullable(predicate.getValue().getDynamicValue()); |
||||
|
} else { |
||||
|
return Optional.empty(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public String getSessionId() { |
||||
|
return sessionRef.getSessionId(); |
||||
|
} |
||||
|
|
||||
|
public TenantId getTenantId() { |
||||
|
return sessionRef.getSecurityCtx().getTenantId(); |
||||
|
} |
||||
|
|
||||
|
public CustomerId getCustomerId() { |
||||
|
return sessionRef.getSecurityCtx().getCustomerId(); |
||||
|
} |
||||
|
|
||||
|
public UserId getUserId() { |
||||
|
return sessionRef.getSecurityCtx().getId(); |
||||
|
} |
||||
|
|
||||
|
protected void clearDynamicValueSubscriptions() { |
||||
|
if (subToDynamicValueKeySet != null) { |
||||
|
for (Integer subId : subToDynamicValueKeySet) { |
||||
|
localSubscriptionService.cancelSubscription(sessionRef.getSessionId(), subId); |
||||
|
} |
||||
|
subToDynamicValueKeySet.clear(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void setRefreshTask(ScheduledFuture<?> task) { |
||||
|
this.refreshTask = task; |
||||
|
} |
||||
|
|
||||
|
public void cancelTasks() { |
||||
|
if (this.refreshTask != null) { |
||||
|
log.trace("[{}][{}] Canceling old refresh task", sessionRef.getSessionId(), cmdId); |
||||
|
this.refreshTask.cancel(true); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
public static class DynamicValueKey { |
||||
|
@Getter |
||||
|
private final FilterPredicateType predicateType; |
||||
|
@Getter |
||||
|
private final DynamicValueSourceType sourceType; |
||||
|
@Getter |
||||
|
private final String sourceAttribute; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,55 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.subscription; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.common.data.query.EntityCountQuery; |
||||
|
import org.thingsboard.server.common.data.query.EntityKeyType; |
||||
|
import org.thingsboard.server.dao.attributes.AttributesService; |
||||
|
import org.thingsboard.server.dao.entity.EntityService; |
||||
|
import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; |
||||
|
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; |
||||
|
import org.thingsboard.server.service.telemetry.cmd.v2.EntityCountUpdate; |
||||
|
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; |
||||
|
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; |
||||
|
|
||||
|
@Slf4j |
||||
|
public class TbEntityCountSubCtx extends TbAbstractSubCtx<EntityCountQuery> { |
||||
|
|
||||
|
private volatile int result; |
||||
|
|
||||
|
public TbEntityCountSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService, |
||||
|
TbLocalSubscriptionService localSubscriptionService, AttributesService attributesService, |
||||
|
SubscriptionServiceStatistics stats, TelemetryWebSocketSessionRef sessionRef, int cmdId) { |
||||
|
super(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, sessionRef, cmdId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void fetchData() { |
||||
|
result = (int) entityService.countEntitiesByQuery(getTenantId(), getCustomerId(), query); |
||||
|
wsService.sendWsMsg(sessionRef.getSessionId(), new EntityCountUpdate(cmdId, result)); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void update() { |
||||
|
int newCount = (int) entityService.countEntitiesByQuery(getTenantId(), getCustomerId(), query); |
||||
|
if (newCount != result) { |
||||
|
result = newCount; |
||||
|
wsService.sendWsMsg(sessionRef.getSessionId(), new EntityCountUpdate(cmdId, result)); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,33 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.telemetry.cmd.v2; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; |
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.Data; |
||||
|
|
||||
|
@Data |
||||
|
@AllArgsConstructor |
||||
|
@JsonIgnoreProperties(ignoreUnknown = true) |
||||
|
public abstract class CmdUpdate { |
||||
|
|
||||
|
private final int cmdId; |
||||
|
private final int errorCode; |
||||
|
private final String errorMsg; |
||||
|
|
||||
|
public abstract CmdUpdateType getCmdUpdateType(); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,35 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.telemetry.cmd.v2; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonCreator; |
||||
|
import com.fasterxml.jackson.annotation.JsonProperty; |
||||
|
import lombok.Getter; |
||||
|
import org.thingsboard.server.common.data.query.EntityCountQuery; |
||||
|
import org.thingsboard.server.common.data.query.EntityDataQuery; |
||||
|
|
||||
|
public class EntityCountCmd extends DataCmd { |
||||
|
|
||||
|
@Getter |
||||
|
private final EntityCountQuery query; |
||||
|
|
||||
|
@JsonCreator |
||||
|
public EntityCountCmd(@JsonProperty("cmdId") int cmdId, |
||||
|
@JsonProperty("query") EntityCountQuery query) { |
||||
|
super(cmdId); |
||||
|
this.query = query; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,25 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.telemetry.cmd.v2; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
|
||||
|
@Data |
||||
|
public class EntityCountUnsubscribeCmd implements UnsubscribeCmd { |
||||
|
|
||||
|
private final int cmdId; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,57 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.telemetry.cmd.v2; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonCreator; |
||||
|
import com.fasterxml.jackson.annotation.JsonProperty; |
||||
|
import lombok.Getter; |
||||
|
import lombok.ToString; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.query.EntityData; |
||||
|
import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; |
||||
|
|
||||
|
import java.util.List; |
||||
|
|
||||
|
@ToString |
||||
|
public class EntityCountUpdate extends CmdUpdate { |
||||
|
|
||||
|
@Getter |
||||
|
private int count; |
||||
|
|
||||
|
public EntityCountUpdate(int cmdId, int count) { |
||||
|
super(cmdId, SubscriptionErrorCode.NO_ERROR.getCode(), null); |
||||
|
this.count = count; |
||||
|
} |
||||
|
|
||||
|
public EntityCountUpdate(int cmdId, int errorCode, String errorMsg) { |
||||
|
super(cmdId, errorCode, errorMsg); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public CmdUpdateType getCmdUpdateType() { |
||||
|
return CmdUpdateType.COUNT_DATA; |
||||
|
} |
||||
|
|
||||
|
@JsonCreator |
||||
|
public EntityCountUpdate(@JsonProperty("cmdId") int cmdId, |
||||
|
@JsonProperty("count") int count, |
||||
|
@JsonProperty("errorCode") int errorCode, |
||||
|
@JsonProperty("errorMsg") String errorMsg) { |
||||
|
super(cmdId, errorCode, errorMsg); |
||||
|
this.count = count; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,30 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.common.data.device.profile; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.common.data.query.EntityKeyValueType; |
||||
|
import org.thingsboard.server.common.data.query.KeyFilterPredicate; |
||||
|
|
||||
|
@Data |
||||
|
public class AlarmConditionFilter { |
||||
|
|
||||
|
private AlarmConditionFilterKey key; |
||||
|
private EntityKeyValueType valueType; |
||||
|
private Object value; |
||||
|
private KeyFilterPredicate predicate; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,26 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.common.data.device.profile; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
|
||||
|
@Data |
||||
|
public class AlarmConditionFilterKey { |
||||
|
|
||||
|
private final AlarmConditionKeyType type; |
||||
|
private final String key; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,23 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.common.data.device.profile; |
||||
|
|
||||
|
public enum AlarmConditionKeyType { |
||||
|
ATTRIBUTE, |
||||
|
TIME_SERIES, |
||||
|
ENTITY_FIELD, |
||||
|
CONSTANT |
||||
|
} |
||||
@ -0,0 +1,30 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.common.data.query; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
|
||||
|
@Data |
||||
|
public class EntityTypeFilter implements EntityFilter { |
||||
|
@Override |
||||
|
public EntityFilterType getType() { |
||||
|
return EntityFilterType.ENTITY_TYPE; |
||||
|
} |
||||
|
|
||||
|
private EntityType entityType; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,39 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.attributes; |
||||
|
|
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.EqualsAndHashCode; |
||||
|
import lombok.Getter; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
|
||||
|
import java.io.Serializable; |
||||
|
|
||||
|
@EqualsAndHashCode |
||||
|
@Getter |
||||
|
@AllArgsConstructor |
||||
|
public class AttributeCacheKey implements Serializable { |
||||
|
private static final long serialVersionUID = 2013369077925351881L; |
||||
|
|
||||
|
private final String scope; |
||||
|
private final EntityId entityId; |
||||
|
private final String key; |
||||
|
|
||||
|
@Override |
||||
|
public String toString() { |
||||
|
return entityId + "_" + scope + "_" + key; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,39 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.attributes; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
||||
|
import org.thingsboard.server.dao.exception.IncorrectParameterException; |
||||
|
import org.thingsboard.server.dao.service.Validator; |
||||
|
|
||||
|
public class AttributeUtils { |
||||
|
public static void validate(EntityId id, String scope) { |
||||
|
Validator.validateId(id.getId(), "Incorrect id " + id); |
||||
|
Validator.validateString(scope, "Incorrect scope " + scope); |
||||
|
} |
||||
|
|
||||
|
public static void validate(AttributeKvEntry kvEntry) { |
||||
|
if (kvEntry == null) { |
||||
|
throw new IncorrectParameterException("Key value entry can't be null"); |
||||
|
} else if (kvEntry.getDataType() == null) { |
||||
|
throw new IncorrectParameterException("Incorrect kvEntry. Data type can't be null"); |
||||
|
} else { |
||||
|
Validator.validateString(kvEntry.getKey(), "Incorrect kvEntry. Key can't be empty"); |
||||
|
Validator.validatePositiveNumber(kvEntry.getLastUpdateTs(), "Incorrect last update ts. Ts should be positive"); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,63 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.attributes; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
||||
|
import org.springframework.cache.Cache; |
||||
|
import org.springframework.cache.CacheManager; |
||||
|
import org.springframework.context.annotation.Primary; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.CacheConstants.ATTRIBUTES_CACHE; |
||||
|
|
||||
|
@Service |
||||
|
@ConditionalOnProperty(prefix = "cache.attributes", value = "enabled", havingValue = "true") |
||||
|
@Primary |
||||
|
@Slf4j |
||||
|
public class AttributesCacheWrapper { |
||||
|
private final Cache attributesCache; |
||||
|
|
||||
|
public AttributesCacheWrapper(CacheManager cacheManager) { |
||||
|
this.attributesCache = cacheManager.getCache(ATTRIBUTES_CACHE); |
||||
|
} |
||||
|
|
||||
|
public Cache.ValueWrapper get(AttributeCacheKey attributeCacheKey) { |
||||
|
try { |
||||
|
return attributesCache.get(attributeCacheKey); |
||||
|
} catch (Exception e) { |
||||
|
log.debug("Failed to retrieve element from cache for key {}. Reason - {}.", attributeCacheKey, e.getMessage()); |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void put(AttributeCacheKey attributeCacheKey, AttributeKvEntry attributeKvEntry) { |
||||
|
try { |
||||
|
attributesCache.put(attributeCacheKey, attributeKvEntry); |
||||
|
} catch (Exception e) { |
||||
|
log.debug("Failed to put element from cache for key {}. Reason - {}.", attributeCacheKey, e.getMessage()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void evict(AttributeCacheKey attributeCacheKey) { |
||||
|
try { |
||||
|
attributesCache.evict(attributeCacheKey); |
||||
|
} catch (Exception e) { |
||||
|
log.debug("Failed to evict element from cache for key {}. Reason - {}.", attributeCacheKey, e.getMessage()); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,193 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2021 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.attributes; |
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
||||
|
import org.springframework.cache.Cache; |
||||
|
import org.springframework.cache.CacheManager; |
||||
|
import org.springframework.context.annotation.Primary; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
||||
|
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.stats.DefaultCounter; |
||||
|
import org.thingsboard.server.common.stats.StatsFactory; |
||||
|
import org.thingsboard.server.dao.service.Validator; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.Collection; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.HashSet; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.Objects; |
||||
|
import java.util.Optional; |
||||
|
import java.util.Set; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.CacheConstants.ATTRIBUTES_CACHE; |
||||
|
import static org.thingsboard.server.dao.attributes.AttributeUtils.validate; |
||||
|
|
||||
|
@Service |
||||
|
@ConditionalOnProperty(prefix = "cache.attributes", value = "enabled", havingValue = "true") |
||||
|
@Primary |
||||
|
@Slf4j |
||||
|
public class CachedAttributesService implements AttributesService { |
||||
|
private static final String STATS_NAME = "attributes.cache"; |
||||
|
|
||||
|
private final AttributesDao attributesDao; |
||||
|
private final AttributesCacheWrapper cacheWrapper; |
||||
|
private final DefaultCounter hitCounter; |
||||
|
private final DefaultCounter missCounter; |
||||
|
|
||||
|
public CachedAttributesService(AttributesDao attributesDao, |
||||
|
AttributesCacheWrapper cacheWrapper, |
||||
|
StatsFactory statsFactory) { |
||||
|
this.attributesDao = attributesDao; |
||||
|
this.cacheWrapper = cacheWrapper; |
||||
|
|
||||
|
this.hitCounter = statsFactory.createDefaultCounter(STATS_NAME, "result", "hit"); |
||||
|
this.missCounter = statsFactory.createDefaultCounter(STATS_NAME, "result", "miss"); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, String attributeKey) { |
||||
|
validate(entityId, scope); |
||||
|
Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey); |
||||
|
|
||||
|
AttributeCacheKey attributeCacheKey = new AttributeCacheKey(scope, entityId, attributeKey); |
||||
|
Cache.ValueWrapper cachedAttributeValue = cacheWrapper.get(attributeCacheKey); |
||||
|
if (cachedAttributeValue != null) { |
||||
|
hitCounter.increment(); |
||||
|
AttributeKvEntry cachedAttributeKvEntry = (AttributeKvEntry) cachedAttributeValue.get(); |
||||
|
return Futures.immediateFuture(Optional.ofNullable(cachedAttributeKvEntry)); |
||||
|
} else { |
||||
|
missCounter.increment(); |
||||
|
ListenableFuture<Optional<AttributeKvEntry>> result = attributesDao.find(tenantId, entityId, scope, attributeKey); |
||||
|
return Futures.transform(result, foundAttrKvEntry -> { |
||||
|
// TODO: think if it's a good idea to store 'empty' attributes
|
||||
|
cacheWrapper.put(attributeCacheKey, foundAttrKvEntry.orElse(null)); |
||||
|
return foundAttrKvEntry; |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, Collection<String> attributeKeys) { |
||||
|
validate(entityId, scope); |
||||
|
attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey)); |
||||
|
|
||||
|
Map<String, Cache.ValueWrapper> wrappedCachedAttributes = findCachedAttributes(entityId, scope, attributeKeys); |
||||
|
|
||||
|
List<AttributeKvEntry> cachedAttributes = wrappedCachedAttributes.values().stream() |
||||
|
.map(wrappedCachedAttribute -> (AttributeKvEntry) wrappedCachedAttribute.get()) |
||||
|
.filter(Objects::nonNull) |
||||
|
.collect(Collectors.toList()); |
||||
|
if (wrappedCachedAttributes.size() == attributeKeys.size()) { |
||||
|
return Futures.immediateFuture(cachedAttributes); |
||||
|
} |
||||
|
|
||||
|
Set<String> notFoundAttributeKeys = new HashSet<>(attributeKeys); |
||||
|
notFoundAttributeKeys.removeAll(wrappedCachedAttributes.keySet()); |
||||
|
|
||||
|
ListenableFuture<List<AttributeKvEntry>> result = attributesDao.find(tenantId, entityId, scope, notFoundAttributeKeys); |
||||
|
return Futures.transform(result, foundInDbAttributes -> mergeDbAndCacheAttributes(entityId, scope, cachedAttributes, notFoundAttributeKeys, foundInDbAttributes), MoreExecutors.directExecutor()); |
||||
|
|
||||
|
} |
||||
|
|
||||
|
private Map<String, Cache.ValueWrapper> findCachedAttributes(EntityId entityId, String scope, Collection<String> attributeKeys) { |
||||
|
Map<String, Cache.ValueWrapper> cachedAttributes = new HashMap<>(); |
||||
|
for (String attributeKey : attributeKeys) { |
||||
|
Cache.ValueWrapper cachedAttributeValue = cacheWrapper.get(new AttributeCacheKey(scope, entityId, attributeKey)); |
||||
|
if (cachedAttributeValue != null) { |
||||
|
hitCounter.increment(); |
||||
|
cachedAttributes.put(attributeKey, cachedAttributeValue); |
||||
|
} else { |
||||
|
missCounter.increment(); |
||||
|
} |
||||
|
} |
||||
|
return cachedAttributes; |
||||
|
} |
||||
|
|
||||
|
private List<AttributeKvEntry> mergeDbAndCacheAttributes(EntityId entityId, String scope, List<AttributeKvEntry> cachedAttributes, Set<String> notFoundAttributeKeys, List<AttributeKvEntry> foundInDbAttributes) { |
||||
|
for (AttributeKvEntry foundInDbAttribute : foundInDbAttributes) { |
||||
|
AttributeCacheKey attributeCacheKey = new AttributeCacheKey(scope, entityId, foundInDbAttribute.getKey()); |
||||
|
cacheWrapper.put(attributeCacheKey, foundInDbAttribute); |
||||
|
notFoundAttributeKeys.remove(foundInDbAttribute.getKey()); |
||||
|
} |
||||
|
for (String key : notFoundAttributeKeys){ |
||||
|
cacheWrapper.put(new AttributeCacheKey(scope, entityId, key), null); |
||||
|
} |
||||
|
List<AttributeKvEntry> mergedAttributes = new ArrayList<>(cachedAttributes); |
||||
|
mergedAttributes.addAll(foundInDbAttributes); |
||||
|
return mergedAttributes; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, String scope) { |
||||
|
validate(entityId, scope); |
||||
|
return attributesDao.findAll(tenantId, entityId, scope); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId) { |
||||
|
return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds) { |
||||
|
return attributesDao.findAllKeysByEntityIds(tenantId, entityType, entityIds); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<List<Void>> save(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes) { |
||||
|
validate(entityId, scope); |
||||
|
attributes.forEach(AttributeUtils::validate); |
||||
|
|
||||
|
List<ListenableFuture<Void>> saveFutures = attributes.stream().map(attribute -> attributesDao.save(tenantId, entityId, scope, attribute)).collect(Collectors.toList()); |
||||
|
ListenableFuture<List<Void>> future = Futures.allAsList(saveFutures); |
||||
|
|
||||
|
// TODO: can do if (attributesCache.get() != null) attributesCache.put() instead, but will be more twice more requests to cache
|
||||
|
List<String> attributeKeys = attributes.stream().map(KvEntry::getKey).collect(Collectors.toList()); |
||||
|
future.addListener(() -> evictAttributesFromCache(tenantId, entityId, scope, attributeKeys), MoreExecutors.directExecutor()); |
||||
|
return future; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<List<Void>> removeAll(TenantId tenantId, EntityId entityId, String scope, List<String> attributeKeys) { |
||||
|
validate(entityId, scope); |
||||
|
ListenableFuture<List<Void>> future = attributesDao.removeAll(tenantId, entityId, scope, attributeKeys); |
||||
|
future.addListener(() -> evictAttributesFromCache(tenantId, entityId, scope, attributeKeys), MoreExecutors.directExecutor()); |
||||
|
return future; |
||||
|
} |
||||
|
|
||||
|
private void evictAttributesFromCache(TenantId tenantId, EntityId entityId, String scope, List<String> attributeKeys) { |
||||
|
try { |
||||
|
for (String attributeKey : attributeKeys) { |
||||
|
cacheWrapper.evict(new AttributeCacheKey(scope, entityId, attributeKey)); |
||||
|
} |
||||
|
} catch (Exception e) { |
||||
|
log.error("[{}][{}] Failed to remove values from cache.", tenantId, entityId, e); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue