211 changed files with 12380 additions and 3801 deletions
@ -0,0 +1,49 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.actors.plugin; |
||||
|
|
||||
|
import com.hazelcast.util.function.Consumer; |
||||
|
import org.thingsboard.server.extensions.api.exception.UnauthorizedException; |
||||
|
import org.thingsboard.server.extensions.api.plugins.PluginCallback; |
||||
|
import org.thingsboard.server.extensions.api.plugins.PluginContext; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 21.02.17. |
||||
|
*/ |
||||
|
public class ValidationCallback implements PluginCallback<Boolean> { |
||||
|
|
||||
|
private final PluginCallback<?> callback; |
||||
|
private final Consumer<PluginContext> action; |
||||
|
|
||||
|
public ValidationCallback(PluginCallback<?> callback, Consumer<PluginContext> action) { |
||||
|
this.callback = callback; |
||||
|
this.action = action; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onSuccess(PluginContext ctx, Boolean value) { |
||||
|
if (value) { |
||||
|
action.accept(ctx); |
||||
|
} else { |
||||
|
onFailure(ctx, new UnauthorizedException()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(PluginContext ctx, Exception e) { |
||||
|
callback.onFailure(ctx, e); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,25 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.kv; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 20.02.17. |
||||
|
*/ |
||||
|
public enum Aggregation { |
||||
|
|
||||
|
MIN, MAX, AVG, SUM, COUNT, NONE; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,42 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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; |
||||
|
|
||||
|
import javax.annotation.PostConstruct; |
||||
|
import javax.annotation.PreDestroy; |
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.concurrent.Executors; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 21.02.17. |
||||
|
*/ |
||||
|
public abstract class AbstractAsyncDao extends AbstractDao { |
||||
|
|
||||
|
protected ExecutorService readResultsProcessingExecutor; |
||||
|
|
||||
|
@PostConstruct |
||||
|
public void startExecutor() { |
||||
|
readResultsProcessingExecutor = Executors.newCachedThreadPool(); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void stopExecutor() { |
||||
|
if (readResultsProcessingExecutor != null) { |
||||
|
readResultsProcessingExecutor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,193 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.timeseries; |
||||
|
|
||||
|
import com.datastax.driver.core.ResultSet; |
||||
|
import com.datastax.driver.core.Row; |
||||
|
import org.thingsboard.server.common.data.kv.*; |
||||
|
|
||||
|
import javax.annotation.Nullable; |
||||
|
import java.util.List; |
||||
|
import java.util.Optional; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 20.02.17. |
||||
|
*/ |
||||
|
public class AggregatePartitionsFunction implements com.google.common.base.Function<List<ResultSet>, Optional<TsKvEntry>> { |
||||
|
|
||||
|
private static final int LONG_CNT_POS = 0; |
||||
|
private static final int DOUBLE_CNT_POS = 1; |
||||
|
private static final int BOOL_CNT_POS = 2; |
||||
|
private static final int STR_CNT_POS = 3; |
||||
|
private static final int LONG_POS = 4; |
||||
|
private static final int DOUBLE_POS = 5; |
||||
|
private static final int BOOL_POS = 6; |
||||
|
private static final int STR_POS = 7; |
||||
|
|
||||
|
private final Aggregation aggregation; |
||||
|
private final String key; |
||||
|
private final long ts; |
||||
|
|
||||
|
public AggregatePartitionsFunction(Aggregation aggregation, String key, long ts) { |
||||
|
this.aggregation = aggregation; |
||||
|
this.key = key; |
||||
|
this.ts = ts; |
||||
|
} |
||||
|
|
||||
|
@Nullable |
||||
|
@Override |
||||
|
public Optional<TsKvEntry> apply(@Nullable List<ResultSet> rsList) { |
||||
|
if (rsList == null || rsList.size() == 0) { |
||||
|
return Optional.empty(); |
||||
|
} |
||||
|
long count = 0; |
||||
|
DataType dataType = null; |
||||
|
|
||||
|
Boolean bValue = null; |
||||
|
String sValue = null; |
||||
|
Double dValue = null; |
||||
|
Long lValue = null; |
||||
|
|
||||
|
for (ResultSet rs : rsList) { |
||||
|
for (Row row : rs.all()) { |
||||
|
long curCount; |
||||
|
|
||||
|
Long curLValue = null; |
||||
|
Double curDValue = null; |
||||
|
Boolean curBValue = null; |
||||
|
String curSValue = null; |
||||
|
|
||||
|
long longCount = row.getLong(LONG_CNT_POS); |
||||
|
long doubleCount = row.getLong(DOUBLE_CNT_POS); |
||||
|
long boolCount = row.getLong(BOOL_CNT_POS); |
||||
|
long strCount = row.getLong(STR_CNT_POS); |
||||
|
|
||||
|
if (longCount > 0) { |
||||
|
dataType = DataType.LONG; |
||||
|
curCount = longCount; |
||||
|
curLValue = getLongValue(row); |
||||
|
} else if (doubleCount > 0) { |
||||
|
dataType = DataType.DOUBLE; |
||||
|
curCount = doubleCount; |
||||
|
curDValue = getDoubleValue(row); |
||||
|
} else if (boolCount > 0) { |
||||
|
dataType = DataType.BOOLEAN; |
||||
|
curCount = boolCount; |
||||
|
curBValue = getBooleanValue(row); |
||||
|
} else if (strCount > 0) { |
||||
|
dataType = DataType.STRING; |
||||
|
curCount = strCount; |
||||
|
curSValue = getStringValue(row); |
||||
|
} else { |
||||
|
continue; |
||||
|
} |
||||
|
|
||||
|
if (aggregation == Aggregation.COUNT) { |
||||
|
count += curCount; |
||||
|
} else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { |
||||
|
count += curCount; |
||||
|
if (curDValue != null) { |
||||
|
dValue = dValue == null ? curDValue : dValue + curDValue; |
||||
|
} else if (curLValue != null) { |
||||
|
lValue = lValue == null ? curLValue : lValue + curLValue; |
||||
|
} |
||||
|
} else if (aggregation == Aggregation.MIN) { |
||||
|
if (curDValue != null) { |
||||
|
dValue = dValue == null ? curDValue : Math.min(dValue, curDValue); |
||||
|
} else if (curLValue != null) { |
||||
|
lValue = lValue == null ? curLValue : Math.min(lValue, curLValue); |
||||
|
} else if (curBValue != null) { |
||||
|
bValue = bValue == null ? curBValue : bValue && curBValue; |
||||
|
} else if (curSValue != null) { |
||||
|
if (sValue == null || curSValue.compareTo(sValue) < 0) { |
||||
|
sValue = curSValue; |
||||
|
} |
||||
|
} |
||||
|
} else if (aggregation == Aggregation.MAX) { |
||||
|
if (curDValue != null) { |
||||
|
dValue = dValue == null ? curDValue : Math.max(dValue, curDValue); |
||||
|
} else if (curLValue != null) { |
||||
|
lValue = lValue == null ? curLValue : Math.max(lValue, curLValue); |
||||
|
} else if (curBValue != null) { |
||||
|
bValue = bValue == null ? curBValue : bValue || curBValue; |
||||
|
} else if (curSValue != null) { |
||||
|
if (sValue == null || curSValue.compareTo(sValue) > 0) { |
||||
|
sValue = curSValue; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
if (dataType == null) { |
||||
|
return Optional.empty(); |
||||
|
} else if (aggregation == Aggregation.COUNT) { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, (long) count))); |
||||
|
} else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { |
||||
|
if (count == 0 || (dataType == DataType.DOUBLE && dValue == null) || (dataType == DataType.LONG && lValue == null)) { |
||||
|
return Optional.empty(); |
||||
|
} else if (dataType == DataType.DOUBLE) { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, aggregation == Aggregation.SUM ? dValue : (dValue / count)))); |
||||
|
} else if (dataType == DataType.LONG) { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggregation == Aggregation.SUM ? lValue : (lValue / count)))); |
||||
|
} |
||||
|
} else if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) { |
||||
|
if (dataType == DataType.DOUBLE) { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, dValue))); |
||||
|
} else if (dataType == DataType.LONG) { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, lValue))); |
||||
|
} else if (dataType == DataType.STRING) { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new StringDataEntry(key, sValue))); |
||||
|
} else { |
||||
|
return Optional.of(new BasicTsKvEntry(ts, new BooleanDataEntry(key, bValue))); |
||||
|
} |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
private Boolean getBooleanValue(Row row) { |
||||
|
if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) { |
||||
|
return row.getBool(BOOL_POS); |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private String getStringValue(Row row) { |
||||
|
if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) { |
||||
|
return row.getString(STR_POS); |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private Long getLongValue(Row row) { |
||||
|
if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX |
||||
|
|| aggregation == Aggregation.SUM || aggregation == Aggregation.AVG) { |
||||
|
return row.getLong(LONG_POS); |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private Double getDoubleValue(Row row) { |
||||
|
if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX |
||||
|
|| aggregation == Aggregation.SUM || aggregation == Aggregation.AVG) { |
||||
|
return row.getDouble(DOUBLE_POS); |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,29 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.timeseries; |
||||
|
|
||||
|
import com.google.common.util.concurrent.AbstractFuture; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 21.02.17. |
||||
|
*/ |
||||
|
public class SimpleListenableFuture<V> extends AbstractFuture<V> { |
||||
|
|
||||
|
public boolean set(V value) { |
||||
|
return super.set(value); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,82 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.timeseries; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import lombok.Getter; |
||||
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.TsKvQuery; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 21.02.17. |
||||
|
*/ |
||||
|
public class TsKvQueryCursor { |
||||
|
@Getter |
||||
|
private final String entityType; |
||||
|
@Getter |
||||
|
private final UUID entityId; |
||||
|
@Getter |
||||
|
private final String key; |
||||
|
@Getter |
||||
|
private final long startTs; |
||||
|
@Getter |
||||
|
private final long endTs; |
||||
|
private final List<Long> partitions; |
||||
|
@Getter |
||||
|
private final List<TsKvEntry> data; |
||||
|
|
||||
|
private int partitionIndex; |
||||
|
private int currentLimit; |
||||
|
|
||||
|
public TsKvQueryCursor(String entityType, UUID entityId, TsKvQuery baseQuery, List<Long> partitions) { |
||||
|
this.entityType = entityType; |
||||
|
this.entityId = entityId; |
||||
|
this.key = baseQuery.getKey(); |
||||
|
this.startTs = baseQuery.getStartTs(); |
||||
|
this.endTs = baseQuery.getEndTs(); |
||||
|
this.partitions = partitions; |
||||
|
this.partitionIndex = partitions.size() - 1; |
||||
|
this.data = new ArrayList<>(); |
||||
|
this.currentLimit = baseQuery.getLimit(); |
||||
|
} |
||||
|
|
||||
|
public boolean hasNextPartition() { |
||||
|
return partitionIndex >= 0; |
||||
|
} |
||||
|
|
||||
|
public boolean isFull() { |
||||
|
return currentLimit <= 0; |
||||
|
} |
||||
|
|
||||
|
public long getNextPartition() { |
||||
|
long partition = partitions.get(partitionIndex); |
||||
|
partitionIndex--; |
||||
|
return partition; |
||||
|
} |
||||
|
|
||||
|
public int getCurrentLimit() { |
||||
|
return currentLimit; |
||||
|
} |
||||
|
|
||||
|
public void addData(List<TsKvEntry> newData) { |
||||
|
currentLimit -= newData.size(); |
||||
|
data.addAll(newData); |
||||
|
} |
||||
|
} |
||||
File diff suppressed because one or more lines are too long
@ -0,0 +1,22 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.extensions.api.exception; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 21.02.17. |
||||
|
*/ |
||||
|
public class UnauthorizedException extends Exception { |
||||
|
} |
||||
@ -0,0 +1,74 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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.extensions.core.plugin.telemetry.handlers; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.extensions.api.plugins.PluginCallback; |
||||
|
import org.thingsboard.server.extensions.api.plugins.PluginContext; |
||||
|
|
||||
|
/** |
||||
|
* Created by ashvayka on 21.02.17. |
||||
|
*/ |
||||
|
@Slf4j |
||||
|
public abstract class BiPluginCallBack<V1, V2> { |
||||
|
|
||||
|
private V1 v1; |
||||
|
private V2 v2; |
||||
|
|
||||
|
public PluginCallback<V1> getV1Callback() { |
||||
|
return new PluginCallback<V1>() { |
||||
|
@Override |
||||
|
public void onSuccess(PluginContext ctx, V1 value) { |
||||
|
synchronized (BiPluginCallBack.this) { |
||||
|
BiPluginCallBack.this.v1 = value; |
||||
|
if (v2 != null) { |
||||
|
BiPluginCallBack.this.onSuccess(ctx, v1, v2); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(PluginContext ctx, Exception e) { |
||||
|
BiPluginCallBack.this.onFailure(ctx, e); |
||||
|
} |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
public PluginCallback<V2> getV2Callback() { |
||||
|
return new PluginCallback<V2>() { |
||||
|
@Override |
||||
|
public void onSuccess(PluginContext ctx, V2 value) { |
||||
|
synchronized (BiPluginCallBack.this) { |
||||
|
BiPluginCallBack.this.v2 = value; |
||||
|
if (v1 != null) { |
||||
|
BiPluginCallBack.this.onSuccess(ctx, v1, v2); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(PluginContext ctx, Exception e) { |
||||
|
BiPluginCallBack.this.onFailure(ctx, e); |
||||
|
} |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
abstract public void onSuccess(PluginContext ctx, V1 v1, V2 v2); |
||||
|
|
||||
|
abstract public void onFailure(PluginContext ctx, Exception e); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,274 @@ |
|||||
|
/* |
||||
|
* Copyright © 2016-2017 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. |
||||
|
*/ |
||||
|
|
||||
|
export default class DataAggregator { |
||||
|
|
||||
|
constructor(onDataCb, tsKeyNames, startTs, limit, aggregationType, timeWindow, interval, types, $timeout, $filter) { |
||||
|
this.onDataCb = onDataCb; |
||||
|
this.tsKeyNames = tsKeyNames; |
||||
|
this.startTs = startTs; |
||||
|
this.aggregationType = aggregationType; |
||||
|
this.types = types; |
||||
|
this.$timeout = $timeout; |
||||
|
this.$filter = $filter; |
||||
|
this.dataReceived = false; |
||||
|
this.resetPending = false; |
||||
|
this.noAggregation = aggregationType === types.aggregation.none.value; |
||||
|
this.limit = limit; |
||||
|
this.timeWindow = timeWindow; |
||||
|
this.interval = interval; |
||||
|
this.aggregationTimeout = Math.max(this.interval, 1000); |
||||
|
switch (aggregationType) { |
||||
|
case types.aggregation.min.value: |
||||
|
this.aggFunction = min; |
||||
|
break; |
||||
|
case types.aggregation.max.value: |
||||
|
this.aggFunction = max; |
||||
|
break; |
||||
|
case types.aggregation.avg.value: |
||||
|
this.aggFunction = avg; |
||||
|
break; |
||||
|
case types.aggregation.sum.value: |
||||
|
this.aggFunction = sum; |
||||
|
break; |
||||
|
case types.aggregation.count.value: |
||||
|
this.aggFunction = count; |
||||
|
break; |
||||
|
case types.aggregation.none.value: |
||||
|
this.aggFunction = none; |
||||
|
break; |
||||
|
default: |
||||
|
this.aggFunction = avg; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
reset(startTs, timeWindow, interval) { |
||||
|
if (this.intervalTimeoutHandle) { |
||||
|
this.$timeout.cancel(this.intervalTimeoutHandle); |
||||
|
this.intervalTimeoutHandle = null; |
||||
|
} |
||||
|
this.intervalScheduledTime = currentTime(); |
||||
|
this.startTs = startTs; |
||||
|
this.timeWindow = timeWindow; |
||||
|
this.interval = interval; |
||||
|
this.endTs = this.startTs + this.timeWindow; |
||||
|
this.elapsed = 0; |
||||
|
this.aggregationTimeout = Math.max(this.interval, 1000); |
||||
|
this.resetPending = true; |
||||
|
var self = this; |
||||
|
this.intervalTimeoutHandle = this.$timeout(function() { |
||||
|
self.onInterval(); |
||||
|
}, this.aggregationTimeout, false); |
||||
|
} |
||||
|
|
||||
|
onData(data, update, history, apply) { |
||||
|
if (!this.dataReceived || this.resetPending) { |
||||
|
var updateIntervalScheduledTime = true; |
||||
|
if (!this.dataReceived) { |
||||
|
this.elapsed = 0; |
||||
|
this.dataReceived = true; |
||||
|
this.endTs = this.startTs + this.timeWindow; |
||||
|
} |
||||
|
if (this.resetPending) { |
||||
|
this.resetPending = false; |
||||
|
updateIntervalScheduledTime = false; |
||||
|
} |
||||
|
if (update) { |
||||
|
this.aggregationMap = {}; |
||||
|
updateAggregatedData(this.aggregationMap, this.aggregationType === this.types.aggregation.count.value, |
||||
|
this.noAggregation, this.aggFunction, data.data, this.interval, this.startTs); |
||||
|
} else { |
||||
|
this.aggregationMap = processAggregatedData(data.data, this.aggregationType === this.types.aggregation.count.value, this.noAggregation); |
||||
|
} |
||||
|
if (updateIntervalScheduledTime) { |
||||
|
this.intervalScheduledTime = currentTime(); |
||||
|
} |
||||
|
this.onInterval(history, apply); |
||||
|
} else { |
||||
|
updateAggregatedData(this.aggregationMap, this.aggregationType === this.types.aggregation.count.value, |
||||
|
this.noAggregation, this.aggFunction, data.data, this.interval, this.startTs); |
||||
|
if (history) { |
||||
|
this.intervalScheduledTime = currentTime(); |
||||
|
this.onInterval(history, apply); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
onInterval(history, apply) { |
||||
|
var now = currentTime(); |
||||
|
this.elapsed += now - this.intervalScheduledTime; |
||||
|
this.intervalScheduledTime = now; |
||||
|
if (this.intervalTimeoutHandle) { |
||||
|
this.$timeout.cancel(this.intervalTimeoutHandle); |
||||
|
this.intervalTimeoutHandle = null; |
||||
|
} |
||||
|
if (!history) { |
||||
|
var delta = Math.floor(this.elapsed / this.interval); |
||||
|
if (delta || !this.data) { |
||||
|
this.startTs += delta * this.interval; |
||||
|
this.endTs += delta * this.interval; |
||||
|
this.data = toData(this.tsKeyNames, this.aggregationMap, this.startTs, this.endTs, this.$filter, this.limit); |
||||
|
this.elapsed = this.elapsed - delta * this.interval; |
||||
|
} |
||||
|
} else { |
||||
|
this.data = toData(this.tsKeyNames, this.aggregationMap, this.startTs, this.endTs, this.$filter, this.limit); |
||||
|
} |
||||
|
if (this.onDataCb) { |
||||
|
this.onDataCb(this.data, this.startTs, this.endTs, apply); |
||||
|
} |
||||
|
|
||||
|
var self = this; |
||||
|
if (!history) { |
||||
|
this.intervalTimeoutHandle = this.$timeout(function() { |
||||
|
self.onInterval(); |
||||
|
}, this.aggregationTimeout, false); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
destroy() { |
||||
|
if (this.intervalTimeoutHandle) { |
||||
|
this.$timeout.cancel(this.intervalTimeoutHandle); |
||||
|
this.intervalTimeoutHandle = null; |
||||
|
} |
||||
|
this.aggregationMap = null; |
||||
|
} |
||||
|
|
||||
|
} |
||||
|
|
||||
|
/* eslint-disable */ |
||||
|
function currentTime() { |
||||
|
return window.performance && window.performance.now ? |
||||
|
window.performance.now() : Date.now(); |
||||
|
} |
||||
|
/* eslint-enable */ |
||||
|
|
||||
|
function processAggregatedData(data, isCount, noAggregation) { |
||||
|
var aggregationMap = {}; |
||||
|
for (var key in data) { |
||||
|
var aggKeyData = aggregationMap[key]; |
||||
|
if (!aggKeyData) { |
||||
|
aggKeyData = {}; |
||||
|
aggregationMap[key] = aggKeyData; |
||||
|
} |
||||
|
var keyData = data[key]; |
||||
|
for (var i in keyData) { |
||||
|
var kvPair = keyData[i]; |
||||
|
var timestamp = kvPair[0]; |
||||
|
var value = convertValue(kvPair[1], noAggregation); |
||||
|
var aggKey = timestamp; |
||||
|
var aggData = { |
||||
|
count: isCount ? value : 1, |
||||
|
sum: value, |
||||
|
aggValue: value |
||||
|
} |
||||
|
aggKeyData[aggKey] = aggData; |
||||
|
} |
||||
|
} |
||||
|
return aggregationMap; |
||||
|
} |
||||
|
|
||||
|
function updateAggregatedData(aggregationMap, isCount, noAggregation, aggFunction, data, interval, startTs) { |
||||
|
for (var key in data) { |
||||
|
var aggKeyData = aggregationMap[key]; |
||||
|
if (!aggKeyData) { |
||||
|
aggKeyData = {}; |
||||
|
aggregationMap[key] = aggKeyData; |
||||
|
} |
||||
|
var keyData = data[key]; |
||||
|
for (var i in keyData) { |
||||
|
var kvPair = keyData[i]; |
||||
|
var timestamp = kvPair[0]; |
||||
|
var value = convertValue(kvPair[1], noAggregation); |
||||
|
var aggTimestamp = noAggregation ? timestamp : (startTs + Math.floor((timestamp - startTs) / interval) * interval + interval/2); |
||||
|
var aggData = aggKeyData[aggTimestamp]; |
||||
|
if (!aggData) { |
||||
|
aggData = { |
||||
|
count: 1, |
||||
|
sum: value, |
||||
|
aggValue: isCount ? 1 : value |
||||
|
} |
||||
|
aggKeyData[aggTimestamp] = aggData; |
||||
|
} else { |
||||
|
aggFunction(aggData, value); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
function toData(tsKeyNames, aggregationMap, startTs, endTs, $filter, limit) { |
||||
|
var data = {}; |
||||
|
for (var k in tsKeyNames) { |
||||
|
data[tsKeyNames[k]] = []; |
||||
|
} |
||||
|
for (var key in aggregationMap) { |
||||
|
var aggKeyData = aggregationMap[key]; |
||||
|
var keyData = data[key]; |
||||
|
for (var aggTimestamp in aggKeyData) { |
||||
|
if (aggTimestamp <= startTs) { |
||||
|
delete aggKeyData[aggTimestamp]; |
||||
|
} else if (aggTimestamp <= endTs) { |
||||
|
var aggData = aggKeyData[aggTimestamp]; |
||||
|
var kvPair = [Number(aggTimestamp), aggData.aggValue]; |
||||
|
keyData.push(kvPair); |
||||
|
} |
||||
|
} |
||||
|
keyData = $filter('orderBy')(keyData, '+this[0]'); |
||||
|
if (keyData.length > limit) { |
||||
|
keyData = keyData.slice(keyData.length - limit); |
||||
|
} |
||||
|
data[key] = keyData; |
||||
|
} |
||||
|
return data; |
||||
|
} |
||||
|
|
||||
|
function convertValue(value, noAggregation) { |
||||
|
if (!noAggregation || value && isNumeric(value)) { |
||||
|
return Number(value); |
||||
|
} else { |
||||
|
return value; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
function isNumeric(value) { |
||||
|
return (value - parseFloat( value ) + 1) >= 0; |
||||
|
} |
||||
|
|
||||
|
function avg(aggData, value) { |
||||
|
aggData.count++; |
||||
|
aggData.sum += value; |
||||
|
aggData.aggValue = aggData.sum / aggData.count; |
||||
|
} |
||||
|
|
||||
|
function min(aggData, value) { |
||||
|
aggData.aggValue = Math.min(aggData.aggValue, value); |
||||
|
} |
||||
|
|
||||
|
function max(aggData, value) { |
||||
|
aggData.aggValue = Math.max(aggData.aggValue, value); |
||||
|
} |
||||
|
|
||||
|
function sum(aggData, value) { |
||||
|
aggData.aggValue = aggData.aggValue + value; |
||||
|
} |
||||
|
|
||||
|
function count(aggData) { |
||||
|
aggData.count++; |
||||
|
aggData.aggValue = aggData.count; |
||||
|
} |
||||
|
|
||||
|
function none(aggData, value) { |
||||
|
aggData.aggValue = value; |
||||
|
} |
||||
@ -0,0 +1,330 @@ |
|||||
|
/* |
||||
|
* Copyright © 2016-2017 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. |
||||
|
*/ |
||||
|
export default angular.module('thingsboard.api.time', []) |
||||
|
.factory('timeService', TimeService) |
||||
|
.name; |
||||
|
|
||||
|
const SECOND = 1000; |
||||
|
const MINUTE = 60 * SECOND; |
||||
|
const HOUR = 60 * MINUTE; |
||||
|
const DAY = 24 * HOUR; |
||||
|
|
||||
|
const MIN_INTERVAL = SECOND; |
||||
|
const MAX_INTERVAL = 365 * 20 * DAY; |
||||
|
|
||||
|
const MIN_LIMIT = 10; |
||||
|
const AVG_LIMIT = 200; |
||||
|
const MAX_LIMIT = 500; |
||||
|
|
||||
|
/*@ngInject*/ |
||||
|
function TimeService($translate, types) { |
||||
|
|
||||
|
var predefIntervals = [ |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.seconds-interval', {seconds: 1}, 'messageformat'), |
||||
|
value: 1 * SECOND |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.seconds-interval', {seconds: 5}, 'messageformat'), |
||||
|
value: 5 * SECOND |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.seconds-interval', {seconds: 10}, 'messageformat'), |
||||
|
value: 10 * SECOND |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.seconds-interval', {seconds: 15}, 'messageformat'), |
||||
|
value: 15 * SECOND |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.seconds-interval', {seconds: 30}, 'messageformat'), |
||||
|
value: 30 * SECOND |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.minutes-interval', {minutes: 1}, 'messageformat'), |
||||
|
value: 1 * MINUTE |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.minutes-interval', {minutes: 2}, 'messageformat'), |
||||
|
value: 2 * MINUTE |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.minutes-interval', {minutes: 5}, 'messageformat'), |
||||
|
value: 5 * MINUTE |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.minutes-interval', {minutes: 10}, 'messageformat'), |
||||
|
value: 10 * MINUTE |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.minutes-interval', {minutes: 15}, 'messageformat'), |
||||
|
value: 15 * MINUTE |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.minutes-interval', {minutes: 30}, 'messageformat'), |
||||
|
value: 30 * MINUTE |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.hours-interval', {hours: 1}, 'messageformat'), |
||||
|
value: 1 * HOUR |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.hours-interval', {hours: 2}, 'messageformat'), |
||||
|
value: 2 * HOUR |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.hours-interval', {hours: 5}, 'messageformat'), |
||||
|
value: 5 * HOUR |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.hours-interval', {hours: 10}, 'messageformat'), |
||||
|
value: 10 * HOUR |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.hours-interval', {hours: 12}, 'messageformat'), |
||||
|
value: 12 * HOUR |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.days-interval', {days: 1}, 'messageformat'), |
||||
|
value: 1 * DAY |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.days-interval', {days: 7}, 'messageformat'), |
||||
|
value: 7 * DAY |
||||
|
}, |
||||
|
{ |
||||
|
name: $translate.instant('timeinterval.days-interval', {days: 30}, 'messageformat'), |
||||
|
value: 30 * DAY |
||||
|
} |
||||
|
]; |
||||
|
|
||||
|
var service = { |
||||
|
minIntervalLimit: minIntervalLimit, |
||||
|
maxIntervalLimit: maxIntervalLimit, |
||||
|
boundMinInterval: boundMinInterval, |
||||
|
boundMaxInterval: boundMaxInterval, |
||||
|
getIntervals: getIntervals, |
||||
|
matchesExistingInterval: matchesExistingInterval, |
||||
|
boundToPredefinedInterval: boundToPredefinedInterval, |
||||
|
defaultTimewindow: defaultTimewindow, |
||||
|
toHistoryTimewindow: toHistoryTimewindow, |
||||
|
createSubscriptionTimewindow: createSubscriptionTimewindow, |
||||
|
avgAggregationLimit: function () { |
||||
|
return AVG_LIMIT; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
return service; |
||||
|
|
||||
|
function minIntervalLimit(timewindow) { |
||||
|
var min = timewindow / MAX_LIMIT; |
||||
|
return boundMinInterval(min); |
||||
|
} |
||||
|
|
||||
|
function avgInterval(timewindow) { |
||||
|
var avg = timewindow / AVG_LIMIT; |
||||
|
return boundMinInterval(avg); |
||||
|
} |
||||
|
|
||||
|
function maxIntervalLimit(timewindow) { |
||||
|
var max = timewindow / MIN_LIMIT; |
||||
|
return boundMaxInterval(max); |
||||
|
} |
||||
|
|
||||
|
function boundMinInterval(min) { |
||||
|
return toBound(min, MIN_INTERVAL, MAX_INTERVAL, MIN_INTERVAL); |
||||
|
} |
||||
|
|
||||
|
function boundMaxInterval(max) { |
||||
|
return toBound(max, MIN_INTERVAL, MAX_INTERVAL, MAX_INTERVAL); |
||||
|
} |
||||
|
|
||||
|
function toBound(value, min, max, defValue) { |
||||
|
if (angular.isDefined(value)) { |
||||
|
value = Math.max(value, min); |
||||
|
value = Math.min(value, max); |
||||
|
return value; |
||||
|
} else { |
||||
|
return defValue; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
function getIntervals(min, max) { |
||||
|
min = boundMinInterval(min); |
||||
|
max = boundMaxInterval(max); |
||||
|
var intervals = []; |
||||
|
for (var i in predefIntervals) { |
||||
|
var interval = predefIntervals[i]; |
||||
|
if (interval.value >= min && interval.value <= max) { |
||||
|
intervals.push(interval); |
||||
|
} |
||||
|
} |
||||
|
return intervals; |
||||
|
} |
||||
|
|
||||
|
function matchesExistingInterval(min, max, intervalMs) { |
||||
|
var intervals = getIntervals(min, max); |
||||
|
for (var i in intervals) { |
||||
|
var interval = intervals[i]; |
||||
|
if (intervalMs === interval.value) { |
||||
|
return true; |
||||
|
} |
||||
|
} |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
function boundToPredefinedInterval(min, max, intervalMs) { |
||||
|
var intervals = getIntervals(min, max); |
||||
|
var minDelta = MAX_INTERVAL; |
||||
|
var boundedInterval = intervalMs || min; |
||||
|
var matchedInterval; |
||||
|
for (var i in intervals) { |
||||
|
var interval = intervals[i]; |
||||
|
var delta = Math.abs(interval.value - boundedInterval); |
||||
|
if (delta < minDelta) { |
||||
|
matchedInterval = interval; |
||||
|
minDelta = delta; |
||||
|
} |
||||
|
} |
||||
|
boundedInterval = matchedInterval.value; |
||||
|
return boundedInterval; |
||||
|
} |
||||
|
|
||||
|
function defaultTimewindow() { |
||||
|
var currentTime = (new Date).getTime(); |
||||
|
var timewindow = { |
||||
|
displayValue: "", |
||||
|
selectedTab: 0, |
||||
|
realtime: { |
||||
|
interval: SECOND, |
||||
|
timewindowMs: MINUTE // 1 min by default
|
||||
|
}, |
||||
|
history: { |
||||
|
historyType: 0, |
||||
|
interval: SECOND, |
||||
|
timewindowMs: MINUTE, // 1 min by default
|
||||
|
fixedTimewindow: { |
||||
|
startTimeMs: currentTime - DAY, // 1 day by default
|
||||
|
endTimeMs: currentTime |
||||
|
} |
||||
|
}, |
||||
|
aggregation: { |
||||
|
type: types.aggregation.avg.value, |
||||
|
limit: AVG_LIMIT |
||||
|
} |
||||
|
} |
||||
|
return timewindow; |
||||
|
} |
||||
|
|
||||
|
function toHistoryTimewindow(timewindow, startTimeMs, endTimeMs) { |
||||
|
|
||||
|
var interval = 0; |
||||
|
if (timewindow.history) { |
||||
|
interval = timewindow.history.interval; |
||||
|
} else if (timewindow.realtime) { |
||||
|
interval = timewindow.realtime.interval; |
||||
|
} |
||||
|
|
||||
|
var historyTimewindow = { |
||||
|
history: { |
||||
|
fixedTimewindow: { |
||||
|
startTimeMs: startTimeMs, |
||||
|
endTimeMs: endTimeMs |
||||
|
}, |
||||
|
interval: boundIntervalToTimewindow(endTimeMs - startTimeMs, interval) |
||||
|
}, |
||||
|
aggregation: { |
||||
|
|
||||
|
} |
||||
|
} |
||||
|
if (timewindow.aggregation) { |
||||
|
historyTimewindow.aggregation.type = timewindow.aggregation.type || types.aggregation.avg.value; |
||||
|
} else { |
||||
|
historyTimewindow.aggregation.type = types.aggregation.avg.value; |
||||
|
} |
||||
|
|
||||
|
return historyTimewindow; |
||||
|
} |
||||
|
|
||||
|
function createSubscriptionTimewindow(timewindow, stDiff) { |
||||
|
|
||||
|
var subscriptionTimewindow = { |
||||
|
fixedWindow: null, |
||||
|
realtimeWindowMs: null, |
||||
|
aggregation: { |
||||
|
interval: SECOND, |
||||
|
limit: AVG_LIMIT, |
||||
|
type: types.aggregation.avg.value |
||||
|
} |
||||
|
}; |
||||
|
var aggTimewindow = 0; |
||||
|
|
||||
|
if (angular.isDefined(timewindow.aggregation)) { |
||||
|
subscriptionTimewindow.aggregation = { |
||||
|
type: timewindow.aggregation.type || types.aggregation.avg.value, |
||||
|
limit: timewindow.aggregation.limit || AVG_LIMIT |
||||
|
}; |
||||
|
} |
||||
|
if (angular.isDefined(timewindow.realtime)) { |
||||
|
subscriptionTimewindow.realtimeWindowMs = timewindow.realtime.timewindowMs; |
||||
|
subscriptionTimewindow.aggregation.interval = |
||||
|
boundIntervalToTimewindow(subscriptionTimewindow.realtimeWindowMs, timewindow.realtime.interval); |
||||
|
subscriptionTimewindow.startTs = (new Date).getTime() + stDiff - subscriptionTimewindow.realtimeWindowMs; |
||||
|
var startDiff = subscriptionTimewindow.startTs % subscriptionTimewindow.aggregation.interval; |
||||
|
aggTimewindow = subscriptionTimewindow.realtimeWindowMs; |
||||
|
if (startDiff) { |
||||
|
subscriptionTimewindow.startTs -= startDiff; |
||||
|
aggTimewindow += subscriptionTimewindow.aggregation.interval; |
||||
|
} |
||||
|
} else if (angular.isDefined(timewindow.history)) { |
||||
|
if (angular.isDefined(timewindow.history.timewindowMs)) { |
||||
|
var currentTime = (new Date).getTime(); |
||||
|
subscriptionTimewindow.fixedWindow = { |
||||
|
startTimeMs: currentTime - timewindow.history.timewindowMs, |
||||
|
endTimeMs: currentTime |
||||
|
} |
||||
|
aggTimewindow = timewindow.history.timewindowMs; |
||||
|
|
||||
|
} else { |
||||
|
subscriptionTimewindow.fixedWindow = { |
||||
|
startTimeMs: timewindow.history.fixedTimewindow.startTimeMs, |
||||
|
endTimeMs: timewindow.history.fixedTimewindow.endTimeMs |
||||
|
} |
||||
|
aggTimewindow = subscriptionTimewindow.fixedWindow.endTimeMs - subscriptionTimewindow.fixedWindow.startTimeMs; |
||||
|
} |
||||
|
subscriptionTimewindow.startTs = subscriptionTimewindow.fixedWindow.startTimeMs; |
||||
|
subscriptionTimewindow.aggregation.interval = boundIntervalToTimewindow(aggTimewindow, timewindow.history.interval); |
||||
|
} |
||||
|
var aggregation = subscriptionTimewindow.aggregation; |
||||
|
aggregation.timeWindow = aggTimewindow; |
||||
|
if (aggregation.type !== types.aggregation.none.value) { |
||||
|
aggregation.limit = Math.ceil(aggTimewindow / subscriptionTimewindow.aggregation.interval); |
||||
|
} |
||||
|
return subscriptionTimewindow; |
||||
|
} |
||||
|
|
||||
|
function boundIntervalToTimewindow(timewindow, intervalMs) { |
||||
|
var min = minIntervalLimit(timewindow); |
||||
|
var max = maxIntervalLimit(timewindow); |
||||
|
if (intervalMs) { |
||||
|
return toBound(intervalMs, min, max, intervalMs); |
||||
|
} else { |
||||
|
return boundToPredefinedInterval(min, max, avgInterval(timewindow)); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
|
||||
|
} |
||||
@ -0,0 +1,50 @@ |
|||||
|
/* |
||||
|
* Copyright © 2016-2017 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. |
||||
|
*/ |
||||
|
|
||||
|
export default angular.module('thingsboard.raf', []) |
||||
|
.provider('tbRaf', TbRAFProvider) |
||||
|
.name; |
||||
|
|
||||
|
function TbRAFProvider() { |
||||
|
/*@ngInject*/ |
||||
|
this.$get = function($window, $timeout) { |
||||
|
var requestAnimationFrame = $window.requestAnimationFrame || |
||||
|
$window.webkitRequestAnimationFrame; |
||||
|
|
||||
|
var cancelAnimationFrame = $window.cancelAnimationFrame || |
||||
|
$window.webkitCancelAnimationFrame || |
||||
|
$window.webkitCancelRequestAnimationFrame; |
||||
|
|
||||
|
var rafSupported = !!requestAnimationFrame; |
||||
|
var raf = rafSupported |
||||
|
? function(fn) { |
||||
|
var id = requestAnimationFrame(fn); |
||||
|
return function() { |
||||
|
cancelAnimationFrame(id); |
||||
|
}; |
||||
|
} |
||||
|
: function(fn) { |
||||
|
var timer = $timeout(fn, 16.66, false); |
||||
|
return function() { |
||||
|
$timeout.cancel(timer); |
||||
|
}; |
||||
|
}; |
||||
|
|
||||
|
raf.supported = rafSupported; |
||||
|
|
||||
|
return raf; |
||||
|
}; |
||||
|
} |
||||
@ -0,0 +1,217 @@ |
|||||
|
/* |
||||
|
* Copyright © 2016-2017 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. |
||||
|
*/ |
||||
|
|
||||
|
/* eslint-disable import/no-unresolved, import/default */ |
||||
|
|
||||
|
import deviceFilterTemplate from './device-filter.tpl.html'; |
||||
|
|
||||
|
/* eslint-enable import/no-unresolved, import/default */ |
||||
|
|
||||
|
import './device-filter.scss'; |
||||
|
|
||||
|
export default angular.module('thingsboard.directives.deviceFilter', []) |
||||
|
.directive('tbDeviceFilter', DeviceFilter) |
||||
|
.name; |
||||
|
|
||||
|
/*@ngInject*/ |
||||
|
function DeviceFilter($compile, $templateCache, $q, deviceService) { |
||||
|
|
||||
|
var linker = function (scope, element, attrs, ngModelCtrl) { |
||||
|
|
||||
|
var template = $templateCache.get(deviceFilterTemplate); |
||||
|
element.html(template); |
||||
|
|
||||
|
scope.ngModelCtrl = ngModelCtrl; |
||||
|
|
||||
|
scope.fetchDevices = function(searchText, limit) { |
||||
|
var pageLink = {limit: limit, textSearch: searchText}; |
||||
|
|
||||
|
var deferred = $q.defer(); |
||||
|
|
||||
|
deviceService.getTenantDevices(pageLink).then(function success(result) { |
||||
|
deferred.resolve(result.data); |
||||
|
}, function fail() { |
||||
|
deferred.reject(); |
||||
|
}); |
||||
|
|
||||
|
return deferred.promise; |
||||
|
} |
||||
|
|
||||
|
scope.updateValidity = function() { |
||||
|
if (ngModelCtrl.$viewValue) { |
||||
|
var value = ngModelCtrl.$viewValue; |
||||
|
var valid; |
||||
|
if (value.useFilter) { |
||||
|
ngModelCtrl.$setValidity('deviceList', true); |
||||
|
if (angular.isDefined(value.deviceNameFilter) && value.deviceNameFilter.length > 0) { |
||||
|
ngModelCtrl.$setValidity('deviceNameFilter', true); |
||||
|
valid = angular.isDefined(scope.model.matchingFilterDevice) && scope.model.matchingFilterDevice != null; |
||||
|
ngModelCtrl.$setValidity('deviceNameFilterDeviceMatch', valid); |
||||
|
} else { |
||||
|
ngModelCtrl.$setValidity('deviceNameFilter', false); |
||||
|
} |
||||
|
} else { |
||||
|
ngModelCtrl.$setValidity('deviceNameFilter', true); |
||||
|
ngModelCtrl.$setValidity('deviceNameFilterDeviceMatch', true); |
||||
|
valid = angular.isDefined(value.deviceList) && value.deviceList.length > 0; |
||||
|
ngModelCtrl.$setValidity('deviceList', valid); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
ngModelCtrl.$render = function () { |
||||
|
destroyWatchers(); |
||||
|
scope.model = { |
||||
|
useFilter: false, |
||||
|
deviceList: [], |
||||
|
deviceNameFilter: '' |
||||
|
} |
||||
|
if (ngModelCtrl.$viewValue) { |
||||
|
var value = ngModelCtrl.$viewValue; |
||||
|
var model = scope.model; |
||||
|
model.useFilter = value.useFilter === true ? true: false; |
||||
|
model.deviceList = []; |
||||
|
model.deviceNameFilter = value.deviceNameFilter || ''; |
||||
|
processDeviceNameFilter(model.deviceNameFilter).then( |
||||
|
function(device) { |
||||
|
scope.model.matchingFilterDevice = device; |
||||
|
if (value.deviceList && value.deviceList.length > 0) { |
||||
|
deviceService.getDevices(value.deviceList).then(function (devices) { |
||||
|
model.deviceList = devices; |
||||
|
updateMatchingDevice(); |
||||
|
initWatchers(); |
||||
|
}); |
||||
|
} else { |
||||
|
updateMatchingDevice(); |
||||
|
initWatchers(); |
||||
|
} |
||||
|
} |
||||
|
) |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
function updateMatchingDevice() { |
||||
|
if (scope.model.useFilter) { |
||||
|
scope.model.matchingDevice = scope.model.matchingFilterDevice; |
||||
|
} else { |
||||
|
if (scope.model.deviceList && scope.model.deviceList.length > 0) { |
||||
|
scope.model.matchingDevice = scope.model.deviceList[0]; |
||||
|
} else { |
||||
|
scope.model.matchingDevice = null; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
function processDeviceNameFilter(deviceNameFilter) { |
||||
|
var deferred = $q.defer(); |
||||
|
if (angular.isDefined(deviceNameFilter) && deviceNameFilter.length > 0) { |
||||
|
scope.fetchDevices(deviceNameFilter, 1).then(function (devices) { |
||||
|
if (devices && devices.length > 0) { |
||||
|
deferred.resolve(devices[0]); |
||||
|
} else { |
||||
|
deferred.resolve(null); |
||||
|
} |
||||
|
}); |
||||
|
} else { |
||||
|
deferred.resolve(null); |
||||
|
} |
||||
|
return deferred.promise; |
||||
|
} |
||||
|
|
||||
|
function destroyWatchers() { |
||||
|
if (scope.deviceListDeregistration) { |
||||
|
scope.deviceListDeregistration(); |
||||
|
scope.deviceListDeregistration = null; |
||||
|
} |
||||
|
if (scope.useFilterDeregistration) { |
||||
|
scope.useFilterDeregistration(); |
||||
|
scope.useFilterDeregistration = null; |
||||
|
} |
||||
|
if (scope.deviceNameFilterDeregistration) { |
||||
|
scope.deviceNameFilterDeregistration(); |
||||
|
scope.deviceNameFilterDeregistration = null; |
||||
|
} |
||||
|
if (scope.matchingDeviceDeregistration) { |
||||
|
scope.matchingDeviceDeregistration(); |
||||
|
scope.matchingDeviceDeregistration = null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
function initWatchers() { |
||||
|
scope.deviceListDeregistration = scope.$watch('model.deviceList', function () { |
||||
|
if (ngModelCtrl.$viewValue) { |
||||
|
var value = ngModelCtrl.$viewValue; |
||||
|
value.deviceList = []; |
||||
|
if (scope.model.deviceList && scope.model.deviceList.length > 0) { |
||||
|
for (var i in scope.model.deviceList) { |
||||
|
value.deviceList.push(scope.model.deviceList[i].id.id); |
||||
|
} |
||||
|
} |
||||
|
updateMatchingDevice(); |
||||
|
ngModelCtrl.$setViewValue(value); |
||||
|
scope.updateValidity(); |
||||
|
} |
||||
|
}, true); |
||||
|
scope.useFilterDeregistration = scope.$watch('model.useFilter', function () { |
||||
|
if (ngModelCtrl.$viewValue) { |
||||
|
var value = ngModelCtrl.$viewValue; |
||||
|
value.useFilter = scope.model.useFilter; |
||||
|
updateMatchingDevice(); |
||||
|
ngModelCtrl.$setViewValue(value); |
||||
|
scope.updateValidity(); |
||||
|
} |
||||
|
}); |
||||
|
scope.deviceNameFilterDeregistration = scope.$watch('model.deviceNameFilter', function (newNameFilter, prevNameFilter) { |
||||
|
if (ngModelCtrl.$viewValue) { |
||||
|
if (!angular.equals(newNameFilter, prevNameFilter)) { |
||||
|
var value = ngModelCtrl.$viewValue; |
||||
|
value.deviceNameFilter = scope.model.deviceNameFilter; |
||||
|
processDeviceNameFilter(value.deviceNameFilter).then( |
||||
|
function(device) { |
||||
|
scope.model.matchingFilterDevice = device; |
||||
|
updateMatchingDevice(); |
||||
|
ngModelCtrl.$setViewValue(value); |
||||
|
scope.updateValidity(); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
} |
||||
|
}); |
||||
|
|
||||
|
scope.matchingDeviceDeregistration = scope.$watch('model.matchingDevice', function (newMatchingDevice, prevMatchingDevice) { |
||||
|
if (!angular.equals(newMatchingDevice, prevMatchingDevice)) { |
||||
|
if (scope.onMatchingDeviceChange) { |
||||
|
scope.onMatchingDeviceChange({device: newMatchingDevice}); |
||||
|
} |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
$compile(element.contents())(scope); |
||||
|
|
||||
|
} |
||||
|
|
||||
|
return { |
||||
|
restrict: "E", |
||||
|
require: "^ngModel", |
||||
|
link: linker, |
||||
|
scope: { |
||||
|
isEdit: '=', |
||||
|
onMatchingDeviceChange: '&' |
||||
|
} |
||||
|
}; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,45 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2017 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. |
||||
|
*/ |
||||
|
.tb-device-filter { |
||||
|
#device_list_chips { |
||||
|
.md-chips { |
||||
|
padding-bottom: 1px; |
||||
|
} |
||||
|
} |
||||
|
.device-name-filter-input { |
||||
|
margin-top: 10px; |
||||
|
margin-bottom: 0px; |
||||
|
.md-errors-spacer { |
||||
|
min-height: 0px; |
||||
|
} |
||||
|
} |
||||
|
.tb-filter-switch { |
||||
|
padding-left: 10px; |
||||
|
.filter-switch { |
||||
|
margin: 0; |
||||
|
} |
||||
|
.filter-label { |
||||
|
margin: 5px 0; |
||||
|
} |
||||
|
} |
||||
|
.tb-error-messages { |
||||
|
margin-top: -11px; |
||||
|
height: 35px; |
||||
|
.tb-error-message { |
||||
|
padding-left: 1px; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,67 @@ |
|||||
|
<!-- |
||||
|
|
||||
|
Copyright © 2016-2017 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. |
||||
|
|
||||
|
--> |
||||
|
<section layout='column' class="tb-device-filter"> |
||||
|
<section layout='row'> |
||||
|
<section layout="column" flex ng-show="!model.useFilter"> |
||||
|
<md-chips flex |
||||
|
id="device_list_chips" |
||||
|
ng-required="!useFilter" |
||||
|
ng-model="model.deviceList" md-autocomplete-snap |
||||
|
md-require-match="true"> |
||||
|
<md-autocomplete |
||||
|
md-no-cache="true" |
||||
|
id="device" |
||||
|
md-selected-item="selectedDevice" |
||||
|
md-search-text="deviceSearchText" |
||||
|
md-items="item in fetchDevices(deviceSearchText, 10)" |
||||
|
md-item-text="item.name" |
||||
|
md-min-length="0" |
||||
|
placeholder="{{ 'device.device-list' | translate }}"> |
||||
|
<md-item-template> |
||||
|
<span md-highlight-text="deviceSearchText" md-highlight-flags="^i">{{item.name}}</span> |
||||
|
</md-item-template> |
||||
|
<md-not-found> |
||||
|
<span translate translate-values='{ device: deviceSearchText }'>device.no-devices-matching</span> |
||||
|
</md-not-found> |
||||
|
</md-autocomplete> |
||||
|
<md-chip-template> |
||||
|
<span> |
||||
|
<strong>{{$chip.name}}</strong> |
||||
|
</span> |
||||
|
</md-chip-template> |
||||
|
</md-chips> |
||||
|
</section> |
||||
|
<section layout="row" flex ng-show="model.useFilter"> |
||||
|
<md-input-container flex class="device-name-filter-input"> |
||||
|
<label translate>device.name-starts-with</label> |
||||
|
<input ng-model="model.deviceNameFilter" aria-label="{{ 'device.name-starts-with' | translate }}"> |
||||
|
</md-input-container> |
||||
|
</section> |
||||
|
<section class="tb-filter-switch" layout="column" layout-align="center center"> |
||||
|
<label class="tb-small filter-label" translate>device.use-device-name-filter</label> |
||||
|
<md-switch class="filter-switch" ng-model="model.useFilter" aria-label="use-filter-switcher"> |
||||
|
</md-switch> |
||||
|
</section> |
||||
|
</section> |
||||
|
<div class="tb-error-messages" ng-messages="ngModelCtrl.$error" role="alert"> |
||||
|
<div translate ng-message="deviceList" class="tb-error-message">device.device-list-empty</div> |
||||
|
<div translate ng-message="deviceNameFilter" class="tb-error-message">device.device-name-filter-required</div> |
||||
|
<div translate translate-values='{ device: model.deviceNameFilter }' ng-message="deviceNameFilterDeviceMatch" |
||||
|
class="tb-error-message">device.device-name-filter-no-device-matched</div> |
||||
|
</div> |
||||
|
</section> |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue