|
|
|
@ -15,7 +15,6 @@ |
|
|
|
*/ |
|
|
|
package org.thingsboard.server.service.telemetry; |
|
|
|
|
|
|
|
import com.google.common.util.concurrent.AsyncFunction; |
|
|
|
import com.google.common.util.concurrent.FutureCallback; |
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
@ -42,6 +41,7 @@ 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.kv.AttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult; |
|
|
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; |
|
|
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
|
|
|
@ -52,7 +52,6 @@ import org.thingsboard.server.dao.util.KvUtils; |
|
|
|
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
|
|
import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; |
|
|
|
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldAttributeUpdateRequest; |
|
|
|
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTimeSeriesUpdateRequest; |
|
|
|
import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService; |
|
|
|
import org.thingsboard.server.service.subscription.TbSubscriptionUtils; |
|
|
|
|
|
|
|
@ -127,7 +126,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
boolean sysTenant = TenantId.SYS_TENANT_ID.equals(tenantId) || tenantId == null; |
|
|
|
if (sysTenant || request.isOnlyLatest() || apiUsageStateService.getApiUsageState(tenantId).isDbStorageEnabled()) { |
|
|
|
KvUtils.validate(request.getEntries(), valueNoXssValidation); |
|
|
|
ListenableFuture<Integer> future = saveTimeseriesInternal(request); |
|
|
|
ListenableFuture<TimeseriesSaveResult> future = saveTimeseriesInternal(request); |
|
|
|
if (!request.isOnlyLatest()) { |
|
|
|
Futures.addCallback(future, getApiUsageCallback(tenantId, request.getCustomerId(), sysTenant), tsCallBackExecutor); |
|
|
|
} |
|
|
|
@ -137,32 +136,25 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public ListenableFuture<Integer> saveTimeseriesInternal(TimeseriesSaveRequest request) { |
|
|
|
public ListenableFuture<TimeseriesSaveResult> saveTimeseriesInternal(TimeseriesSaveRequest request) { |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
ListenableFuture<Integer> saveFuture; |
|
|
|
ListenableFuture<TimeseriesSaveResult> resultFuture; |
|
|
|
if (request.isOnlyLatest()) { |
|
|
|
saveFuture = Futures.transform(tsService.saveLatest(tenantId, entityId, request.getEntries()), result -> 0, MoreExecutors.directExecutor()); |
|
|
|
resultFuture = tsService.saveLatest(tenantId, entityId, request.getEntries()); |
|
|
|
} else if (request.isSaveLatest()) { |
|
|
|
saveFuture = tsService.save(tenantId, entityId, request.getEntries(), request.getTtl()); |
|
|
|
resultFuture = tsService.save(tenantId, entityId, request.getEntries(), request.getTtl()); |
|
|
|
} else { |
|
|
|
saveFuture = tsService.saveWithoutLatest(tenantId, entityId, request.getEntries(), request.getTtl()); |
|
|
|
resultFuture = tsService.saveWithoutLatest(tenantId, entityId, request.getEntries(), request.getTtl()); |
|
|
|
} |
|
|
|
// We need to guarantee, that the message is successfully pushed to the calculated fields service before we execute any callbacks.
|
|
|
|
// saveFuture = Futures.transformAsync(saveFuture, new AsyncFunction<Integer, Integer>() {
|
|
|
|
// @Override
|
|
|
|
// public ListenableFuture<Integer> apply(Integer input) throws Exception {
|
|
|
|
// calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(request));
|
|
|
|
// return input;
|
|
|
|
// }
|
|
|
|
// });
|
|
|
|
addMainCallback(saveFuture, request.getCallback()); |
|
|
|
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, request.getEntries())); |
|
|
|
DonAsynchron.withCallback(resultFuture, result -> { |
|
|
|
calculatedFieldExecutionService.pushRequestToQueue(request, result); |
|
|
|
}, safeCallback(request.getCallback()), tsCallBackExecutor); |
|
|
|
addWsCallback(resultFuture, success -> onTimeSeriesUpdate(tenantId, entityId, request.getEntries())); |
|
|
|
if (request.isSaveLatest() && !request.isOnlyLatest()) { |
|
|
|
addEntityViewCallback(tenantId, entityId, request.getEntries()); |
|
|
|
} |
|
|
|
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(request)), tsCallBackExecutor); |
|
|
|
return saveFuture; |
|
|
|
return resultFuture; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -176,6 +168,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
log.trace("Executing saveInternal [{}]", request); |
|
|
|
ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries()); |
|
|
|
addMainCallback(saveFuture, request.getCallback()); |
|
|
|
//TODO: IM to push to CF queue
|
|
|
|
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice())); |
|
|
|
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor); |
|
|
|
} |
|
|
|
@ -273,27 +266,21 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
} |
|
|
|
|
|
|
|
private void onAttributesUpdate(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice) { |
|
|
|
forwardToSubscriptionManagerService(tenantId, entityId, subscriptionManagerService -> { |
|
|
|
subscriptionManagerService.onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice, TbCallback.EMPTY); |
|
|
|
}, () -> { |
|
|
|
return TbSubscriptionUtils.toAttributesUpdateProto(tenantId, entityId, scope, attributes); |
|
|
|
}); |
|
|
|
forwardToSubscriptionManagerService(tenantId, entityId, |
|
|
|
subscriptionManagerService -> subscriptionManagerService.onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice, TbCallback.EMPTY), |
|
|
|
() -> TbSubscriptionUtils.toAttributesUpdateProto(tenantId, entityId, scope, attributes)); |
|
|
|
} |
|
|
|
|
|
|
|
private void onAttributesDelete(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice) { |
|
|
|
forwardToSubscriptionManagerService(tenantId, entityId, subscriptionManagerService -> { |
|
|
|
subscriptionManagerService.onAttributesDelete(tenantId, entityId, scope, keys, notifyDevice, TbCallback.EMPTY); |
|
|
|
}, () -> { |
|
|
|
return TbSubscriptionUtils.toAttributesDeleteProto(tenantId, entityId, scope, keys, notifyDevice); |
|
|
|
}); |
|
|
|
forwardToSubscriptionManagerService(tenantId, entityId, |
|
|
|
subscriptionManagerService -> subscriptionManagerService.onAttributesDelete(tenantId, entityId, scope, keys, notifyDevice, TbCallback.EMPTY), |
|
|
|
() -> TbSubscriptionUtils.toAttributesDeleteProto(tenantId, entityId, scope, keys, notifyDevice)); |
|
|
|
} |
|
|
|
|
|
|
|
private void onTimeSeriesUpdate(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts) { |
|
|
|
forwardToSubscriptionManagerService(tenantId, entityId, subscriptionManagerService -> { |
|
|
|
subscriptionManagerService.onTimeSeriesUpdate(tenantId, entityId, ts, TbCallback.EMPTY); |
|
|
|
}, () -> { |
|
|
|
return TbSubscriptionUtils.toTimeseriesUpdateProto(tenantId, entityId, ts); |
|
|
|
}); |
|
|
|
forwardToSubscriptionManagerService(tenantId, entityId, |
|
|
|
subscriptionManagerService -> subscriptionManagerService.onTimeSeriesUpdate(tenantId, entityId, ts, TbCallback.EMPTY), |
|
|
|
() -> TbSubscriptionUtils.toTimeseriesUpdateProto(tenantId, entityId, ts)); |
|
|
|
} |
|
|
|
|
|
|
|
private void onTimeSeriesDelete(TenantId tenantId, EntityId entityId, List<String> keys, List<TsKvLatestRemovingResult> ts) { |
|
|
|
@ -313,9 +300,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
|
|
|
|
subscriptionManagerService.onTimeSeriesUpdate(tenantId, entityId, updated, TbCallback.EMPTY); |
|
|
|
subscriptionManagerService.onTimeSeriesDelete(tenantId, entityId, deleted, TbCallback.EMPTY); |
|
|
|
}, () -> { |
|
|
|
return TbSubscriptionUtils.toTimeseriesDeleteProto(tenantId, entityId, keys); |
|
|
|
}); |
|
|
|
}, () -> TbSubscriptionUtils.toTimeseriesDeleteProto(tenantId, entityId, keys)); |
|
|
|
} |
|
|
|
|
|
|
|
private <S> void addMainCallback(ListenableFuture<S> saveFuture, final FutureCallback<Void> callback) { |
|
|
|
@ -333,18 +318,18 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private FutureCallback<Integer> getApiUsageCallback(TenantId tenantId, CustomerId customerId, boolean sysTenant) { |
|
|
|
private FutureCallback<TimeseriesSaveResult> getApiUsageCallback(TenantId tenantId, CustomerId customerId, boolean sysTenant) { |
|
|
|
return new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(Integer result) { |
|
|
|
if (!sysTenant && result != null && result > 0) { |
|
|
|
apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, result); |
|
|
|
public void onSuccess(TimeseriesSaveResult result) { |
|
|
|
Integer dataPoints = result.getDataPoints(); |
|
|
|
if (!sysTenant && dataPoints != null && dataPoints > 0) { |
|
|
|
apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, dataPoints); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
|
|
} |
|
|
|
}; |
|
|
|
} |
|
|
|
|