|
|
|
@ -116,14 +116,6 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
super.shutdownExecutor(); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public ListenableFuture<Void> saveAndNotify(TimeseriesSaveRequest request) { |
|
|
|
SettableFuture<Void> future = SettableFuture.create(); |
|
|
|
request.setCallback(new VoidFutureCallback(future)); |
|
|
|
save(request); |
|
|
|
return future; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void save(TimeseriesSaveRequest request) { |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
@ -132,51 +124,39 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
boolean sysTenant = TenantId.SYS_TENANT_ID.equals(tenantId) || tenantId == null; |
|
|
|
if (sysTenant || apiUsageStateService.getApiUsageState(tenantId).isDbStorageEnabled()) { |
|
|
|
KvUtils.validate(request.getEntries(), valueNoXssValidation); |
|
|
|
FutureCallback<Integer> callback = getCallback(tenantId, request.getCustomerId(), sysTenant, request.getCallback()); |
|
|
|
if (request.isSaveLatest()) { |
|
|
|
saveAndNotifyInternal(tenantId, entityId, request.getEntries(), request.getTtl(), callback); |
|
|
|
} else { |
|
|
|
saveWithoutLatestAndNotifyInternal(tenantId, entityId, request.getEntries(), request.getTtl(), callback); |
|
|
|
} |
|
|
|
FutureCallback<Integer> callback = getApiUsageCallback(tenantId, request.getCustomerId(), sysTenant, request.getCallback()); |
|
|
|
ListenableFuture<Integer> future = saveInternal(request); |
|
|
|
Futures.addCallback(future, callback, tsCallBackExecutor); |
|
|
|
} else { |
|
|
|
request.getCallback().onFailure(new RuntimeException("DB storage writes are disabled due to API limits!")); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private FutureCallback<Integer> getCallback(TenantId tenantId, CustomerId customerId, boolean sysTenant, FutureCallback<Void> callback) { |
|
|
|
return new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(Integer result) { |
|
|
|
if (!sysTenant && result != null && result > 0) { |
|
|
|
apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, result); |
|
|
|
} |
|
|
|
callback.onSuccess(null); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
callback.onFailure(t); |
|
|
|
} |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Integer> callback) { |
|
|
|
saveAndNotifyInternal(tenantId, entityId, ts, 0L, callback); |
|
|
|
public ListenableFuture<Void> saveAndNotify(TimeseriesSaveRequest request) { |
|
|
|
SettableFuture<Void> future = SettableFuture.create(); |
|
|
|
request.setCallback(new VoidFutureCallback(future)); |
|
|
|
save(request); |
|
|
|
return future; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Integer> callback) { |
|
|
|
ListenableFuture<Integer> saveFuture = tsService.save(tenantId, entityId, ts, ttl); |
|
|
|
addMainCallback(saveFuture, callback); |
|
|
|
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); |
|
|
|
addEntityViewCallback(tenantId, entityId, ts); |
|
|
|
} |
|
|
|
public ListenableFuture<Integer> saveInternal(TimeseriesSaveRequest request) { |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
ListenableFuture<Integer> saveFuture; |
|
|
|
if (request.isSaveLatest()) { |
|
|
|
saveFuture = tsService.save(tenantId, entityId, request.getEntries(), request.getTtl()); |
|
|
|
} else { |
|
|
|
saveFuture = tsService.saveWithoutLatest(tenantId, entityId, request.getEntries(), request.getTtl()); |
|
|
|
} |
|
|
|
|
|
|
|
private void saveWithoutLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Integer> callback) { |
|
|
|
ListenableFuture<Integer> saveFuture = tsService.saveWithoutLatest(tenantId, entityId, ts, ttl); |
|
|
|
addMainCallback(saveFuture, callback); |
|
|
|
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); |
|
|
|
addMainCallback(saveFuture, request.getCallback()); |
|
|
|
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, request.getEntries())); |
|
|
|
if (request.isSaveLatest()) { |
|
|
|
addEntityViewCallback(tenantId, entityId, request.getEntries()); |
|
|
|
} |
|
|
|
return saveFuture; |
|
|
|
} |
|
|
|
|
|
|
|
private void addEntityViewCallback(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts) { |
|
|
|
@ -452,11 +432,11 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
}, tsCallBackExecutor); |
|
|
|
} |
|
|
|
|
|
|
|
private <S> void addMainCallback(ListenableFuture<S> saveFuture, final FutureCallback<S> callback) { |
|
|
|
private <S> void addMainCallback(ListenableFuture<S> saveFuture, final FutureCallback<Void> callback) { |
|
|
|
Futures.addCallback(saveFuture, new FutureCallback<S>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(@Nullable S result) { |
|
|
|
callback.onSuccess(result); |
|
|
|
callback.onSuccess(null); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -472,6 +452,23 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private FutureCallback<Integer> getApiUsageCallback(TenantId tenantId, CustomerId customerId, boolean sysTenant, FutureCallback<Void> callback) { |
|
|
|
return new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(Integer result) { |
|
|
|
if (!sysTenant && result != null && result > 0) { |
|
|
|
apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, result); |
|
|
|
} |
|
|
|
callback.onSuccess(null); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
callback.onFailure(t); |
|
|
|
} |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
private static class VoidFutureCallback implements FutureCallback<Void> { |
|
|
|
private final SettableFuture<Void> future; |
|
|
|
|
|
|
|
|