130 changed files with 4959 additions and 544 deletions
@ -0,0 +1,36 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.ctx.state.aggregation.single; |
|||
|
|||
import lombok.AllArgsConstructor; |
|||
import lombok.Data; |
|||
|
|||
@Data |
|||
@AllArgsConstructor |
|||
public class AggIntervalEntry { |
|||
|
|||
private Long startTs; |
|||
private Long endTs; |
|||
|
|||
public boolean belongsToInterval(long ts) { |
|||
return ts >= startTs && ts < endTs; |
|||
} |
|||
|
|||
public long getIntervalDuration() { |
|||
return endTs - startTs; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,45 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.ctx.state.aggregation.single; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonIgnore; |
|||
import lombok.AllArgsConstructor; |
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
@AllArgsConstructor |
|||
public class AggIntervalEntryStatus { |
|||
|
|||
private long lastArgsRefreshTs = -1; |
|||
|
|||
private long lastMetricsEvalTs = -1; |
|||
|
|||
public AggIntervalEntryStatus(long lastArgsRefreshTs) { |
|||
this.lastArgsRefreshTs = lastArgsRefreshTs; |
|||
} |
|||
|
|||
public boolean intervalPassed(long checkInterval) { |
|||
return lastMetricsEvalTs <= System.currentTimeMillis() - checkInterval; |
|||
} |
|||
|
|||
@JsonIgnore |
|||
public boolean argsUpdated() { |
|||
return lastArgsRefreshTs > -1; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,87 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.ctx.state.aggregation.single; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.Data; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.script.api.tbel.TbelCfArg; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; |
|||
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|||
|
|||
import java.util.Map; |
|||
|
|||
@Data |
|||
public class EntityAggregationArgumentEntry implements ArgumentEntry { |
|||
|
|||
private Map<AggIntervalEntry, AggIntervalEntryStatus> aggIntervals; |
|||
|
|||
private boolean forceResetPrevious; |
|||
|
|||
public EntityAggregationArgumentEntry(Map<AggIntervalEntry, AggIntervalEntryStatus> aggIntervals) { |
|||
this.aggIntervals = aggIntervals; |
|||
} |
|||
|
|||
@Override |
|||
public ArgumentEntryType getType() { |
|||
return ArgumentEntryType.ENTITY_AGGREGATION; |
|||
} |
|||
|
|||
@Override |
|||
public Object getValue() { |
|||
return aggIntervals; |
|||
} |
|||
|
|||
@Override |
|||
public boolean updateEntry(ArgumentEntry entry) { |
|||
boolean updated = false; |
|||
if (entry instanceof EntityAggregationArgumentEntry entityAggEntry) { |
|||
aggIntervals.putAll(entityAggEntry.getAggIntervals()); |
|||
} else if (entry instanceof SingleValueArgumentEntry singleValueArgEntry) { |
|||
long entryTs = singleValueArgEntry.getTs(); |
|||
long argUpdateTs = System.currentTimeMillis(); |
|||
for (Map.Entry<AggIntervalEntry, AggIntervalEntryStatus> aggIntervalEntry : aggIntervals.entrySet()) { |
|||
if (singleValueArgEntry.isForceResetPrevious()) { |
|||
aggIntervalEntry.getValue().setLastArgsRefreshTs(argUpdateTs); |
|||
updated = true; |
|||
continue; |
|||
} |
|||
if (aggIntervalEntry.getKey().belongsToInterval(entryTs)) { |
|||
aggIntervalEntry.getValue().setLastArgsRefreshTs(argUpdateTs); |
|||
return true; |
|||
} |
|||
} |
|||
} |
|||
return updated; |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return aggIntervals.isEmpty(); |
|||
} |
|||
|
|||
@Override |
|||
public JsonNode jsonValue() { |
|||
return JacksonUtil.valueToTree(aggIntervals); |
|||
} |
|||
|
|||
@Override |
|||
public TbelCfArg toTbelCfArg() { |
|||
return null; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,271 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.ctx.state.aggregation.single; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ArrayNode; |
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.script.api.tbel.TbUtils; |
|||
import org.thingsboard.server.actors.TbActorRef; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Output; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.AggInterval; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.Watermark; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; |
|||
import org.thingsboard.server.service.cf.CalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
|
|||
import java.time.Instant; |
|||
import java.time.ZoneId; |
|||
import java.time.ZonedDateTime; |
|||
import java.util.ArrayList; |
|||
import java.util.Comparator; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultMetricArgumentEntry; |
|||
|
|||
public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState { |
|||
|
|||
private AggInterval interval; |
|||
private long watermarkDuration; |
|||
private long checkInterval; |
|||
private Map<String, AggMetric> metrics; |
|||
|
|||
private CalculatedFieldProcessingService cfProcessingService; |
|||
|
|||
public EntityAggregationCalculatedFieldState(EntityId entityId) { |
|||
super(entityId); |
|||
} |
|||
|
|||
@Override |
|||
public void setCtx(CalculatedFieldCtx ctx, TbActorRef actorCtx) { |
|||
super.setCtx(ctx, actorCtx); |
|||
this.cfProcessingService = ctx.getCfProcessingService(); |
|||
var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
|||
Watermark watermark = configuration.getWatermark(); |
|||
watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); |
|||
checkInterval = TimeUnit.SECONDS.toMillis(ctx.getSystemContext().getCfCheckInterval()); |
|||
interval = configuration.getInterval(); |
|||
metrics = configuration.getMetrics(); |
|||
} |
|||
|
|||
@Override |
|||
public void init(boolean restored) { |
|||
super.init(restored); |
|||
if (restored) { |
|||
fillMissingIntervals(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public CalculatedFieldType getType() { |
|||
return CalculatedFieldType.ENTITY_AGGREGATION; |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception { |
|||
createIntervalIfNotExist(); |
|||
long now = System.currentTimeMillis(); |
|||
|
|||
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results = new HashMap<>(); |
|||
List<AggIntervalEntry> expiredIntervals = new ArrayList<>(); |
|||
getIntervals().forEach((intervalEntry, argIntervalStatuses) -> { |
|||
processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results); |
|||
}); |
|||
removeExpiredIntervals(expiredIntervals); |
|||
|
|||
Output output = ctx.getOutput(); |
|||
ArrayNode result = toResult(results, output.getDecimalsByDefault()); |
|||
if (result.isEmpty()) { |
|||
return Futures.immediateFuture(TelemetryCalculatedFieldResult.EMPTY); |
|||
} |
|||
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() |
|||
.type(output.getType()) |
|||
.scope(output.getScope()) |
|||
.result(result) |
|||
.build()); |
|||
} |
|||
|
|||
private void removeExpiredIntervals(List<AggIntervalEntry> expiredIntervals) { |
|||
expiredIntervals.forEach(expiredInterval -> { |
|||
arguments.values().stream() |
|||
.map(EntityAggregationArgumentEntry.class::cast) |
|||
.forEach(arg -> arg.getAggIntervals().remove(expiredInterval)); |
|||
}); |
|||
} |
|||
|
|||
private void createIntervalIfNotExist() { |
|||
AggIntervalEntry currentInterval = new AggIntervalEntry(interval.getCurrentIntervalStartTs(), interval.getCurrentIntervalEndTs()); |
|||
arguments.forEach((argName, argumentEntry) -> { |
|||
var entityAggEntry = (EntityAggregationArgumentEntry) argumentEntry; |
|||
entityAggEntry.getAggIntervals().computeIfAbsent(currentInterval, current -> new AggIntervalEntryStatus()); |
|||
}); |
|||
} |
|||
|
|||
private void fillMissingIntervals() { |
|||
ZoneId zoneId = interval.getZoneId(); |
|||
long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); |
|||
|
|||
Map<AggIntervalEntry, Map<String, AggIntervalEntryStatus>> intervals = getIntervals(); |
|||
AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null); |
|||
if (lastIntervalEntry == null) { |
|||
return; |
|||
} |
|||
|
|||
ZonedDateTime nextStart = Instant.ofEpochMilli(lastIntervalEntry.getEndTs()).atZone(zoneId); |
|||
ZonedDateTime nextEnd = interval.getNextIntervalStart(nextStart); |
|||
|
|||
while (nextEnd.toInstant().toEpochMilli() <= currentIntervalEndTs) { |
|||
long nextStartTs = nextStart.toInstant().toEpochMilli(); |
|||
long nextEndTs = nextEnd.toInstant().toEpochMilli(); |
|||
AggIntervalEntry missing = new AggIntervalEntry(nextStartTs, nextEndTs); |
|||
|
|||
arguments.forEach((argName, argumentEntry) -> { |
|||
var entityAggEntry = (EntityAggregationArgumentEntry) argumentEntry; |
|||
AggIntervalEntryStatus intervalEntryStatus = new AggIntervalEntryStatus(System.currentTimeMillis()); |
|||
entityAggEntry.getAggIntervals().computeIfAbsent(missing, missingInterval -> intervalEntryStatus); |
|||
}); |
|||
|
|||
nextStart = nextEnd; |
|||
nextEnd = interval.getNextIntervalStart(nextStart); |
|||
} |
|||
} |
|||
|
|||
private Map<AggIntervalEntry, Map<String, AggIntervalEntryStatus>> getIntervals() { |
|||
Map<AggIntervalEntry, Map<String, AggIntervalEntryStatus>> intervals = new HashMap<>(); |
|||
arguments.forEach((argName, entry) -> { |
|||
var argEntry = (EntityAggregationArgumentEntry) entry; |
|||
argEntry.getAggIntervals().forEach((intervalEntry, status) -> |
|||
intervals.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(argName, status) |
|||
); |
|||
}); |
|||
return intervals; |
|||
} |
|||
|
|||
private void processInterval(long now, |
|||
AggIntervalEntry intervalEntry, |
|||
Map<String, AggIntervalEntryStatus> args, |
|||
List<AggIntervalEntry> expiredIntervals, |
|||
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) { |
|||
long startTs = intervalEntry.getStartTs(); |
|||
long endTs = intervalEntry.getEndTs(); |
|||
|
|||
if (now - endTs > watermarkDuration) { |
|||
handleExpiredInterval(intervalEntry, args, results); |
|||
expiredIntervals.add(intervalEntry); |
|||
} else if (now - startTs >= intervalEntry.getIntervalDuration()) { |
|||
handleActiveInterval(intervalEntry, args, results); |
|||
} |
|||
} |
|||
|
|||
private void handleExpiredInterval(AggIntervalEntry intervalEntry, |
|||
Map<String, AggIntervalEntryStatus> args, |
|||
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) { |
|||
args.forEach((argName, argEntryIntervalStatus) -> { |
|||
if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) { |
|||
argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); |
|||
processMetric(intervalEntry, argName, false, results); |
|||
} else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { |
|||
argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); |
|||
processMetric(intervalEntry, argName, true, results); |
|||
} |
|||
}); |
|||
} |
|||
|
|||
private void handleActiveInterval(AggIntervalEntry intervalEntry, |
|||
Map<String, AggIntervalEntryStatus> args, |
|||
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) { |
|||
args.forEach((argName, argEntryIntervalStatus) -> { |
|||
if (argEntryIntervalStatus.intervalPassed(checkInterval)) { |
|||
if (argEntryIntervalStatus.argsUpdated()) { |
|||
argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); |
|||
argEntryIntervalStatus.setLastArgsRefreshTs(-1); |
|||
processMetric(intervalEntry, argName, false, results); |
|||
} else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { |
|||
argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); |
|||
processMetric(intervalEntry, argName, true, results); |
|||
} |
|||
} |
|||
}); |
|||
} |
|||
|
|||
private void processMetric(AggIntervalEntry intervalEntry, |
|||
String argName, |
|||
boolean useDefault, |
|||
Map<AggIntervalEntry, Map<String, ArgumentEntry>> results) { |
|||
String metricName = findMetricName(argName); |
|||
if (metricName != null) { |
|||
AggMetric metric = metrics.get(metricName); |
|||
String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey(); |
|||
ArgumentEntry metricEntry = useDefault |
|||
? createDefaultMetricArgumentEntry(argKey, metric) |
|||
: cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry); |
|||
if (!metricEntry.isEmpty()) { |
|||
results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metricEntry); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private String findMetricName(String argName) { |
|||
return metrics.entrySet().stream() |
|||
.filter(e -> ((AggKeyInput) e.getValue().getInput()).getKey().equals(argName)) |
|||
.map(Map.Entry::getKey) |
|||
.findFirst() |
|||
.orElse(null); |
|||
} |
|||
|
|||
protected ArrayNode toResult(Map<AggIntervalEntry, Map<String, ArgumentEntry>> results, Integer precision) { |
|||
ArrayNode result = JacksonUtil.newArrayNode(); |
|||
results.forEach((interval, args) -> { |
|||
ObjectNode metricsNode = JacksonUtil.newObjectNode(); |
|||
for (Map.Entry<String, ArgumentEntry> entry : args.entrySet()) { |
|||
String metricName = entry.getKey(); |
|||
ArgumentEntry argumentEntry = entry.getValue(); |
|||
if (!argumentEntry.isEmpty()) { |
|||
Object resultValue = argumentEntry.getValue() instanceof Number number |
|||
? TbUtils.roundResult(number.doubleValue(), precision) |
|||
: argumentEntry.getValue(); |
|||
metricsNode.put(metricName, JacksonUtil.toString(resultValue)); |
|||
} |
|||
} |
|||
if (!metricsNode.isEmpty()) { |
|||
ObjectNode resultNode = JacksonUtil.newObjectNode(); |
|||
resultNode.put("ts", interval.getEndTs() - 1); |
|||
resultNode.set("values", metricsNode); |
|||
result.add(resultNode); |
|||
} |
|||
}); |
|||
return result; |
|||
} |
|||
|
|||
@Override |
|||
public boolean isReady() { |
|||
return true; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,272 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.system; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import com.google.common.hash.Hashing; |
|||
import jakarta.annotation.PostConstruct; |
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.jdbc.core.JdbcTemplate; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.widget.WidgetTypeDetails; |
|||
import org.thingsboard.server.dao.widget.WidgetTypeService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.install.DatabaseSchemaSettingsService; |
|||
import org.thingsboard.server.service.install.InstallScripts; |
|||
import org.thingsboard.server.service.install.update.DefaultDataUpdateService; |
|||
|
|||
import java.io.IOException; |
|||
import java.io.UncheckedIOException; |
|||
import java.nio.file.Files; |
|||
import java.nio.file.NoSuchFileException; |
|||
import java.nio.file.Path; |
|||
import java.util.Objects; |
|||
import java.util.concurrent.ExecutorService; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.atomic.AtomicInteger; |
|||
import java.util.stream.Stream; |
|||
|
|||
/** |
|||
* Runs at application startup and applies no-downtime data updates |
|||
* when the package PATCH version increases (e.g., 4.2.1.0 -> 4.2.1.1). |
|||
*/ |
|||
@Slf4j |
|||
@Component |
|||
@TbCoreComponent |
|||
@RequiredArgsConstructor |
|||
public class SystemPatchApplier { |
|||
|
|||
private static final long ADVISORY_LOCK_ID = 7536891047216478431L; |
|||
|
|||
private final JdbcTemplate jdbcTemplate; |
|||
private final InstallScripts installScripts; |
|||
private final DatabaseSchemaSettingsService schemaSettingsService; |
|||
private final WidgetTypeService widgetTypeService; |
|||
|
|||
@PostConstruct |
|||
private void init() { |
|||
ExecutorService executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("system-patch-applier")); |
|||
executor.submit(() -> { |
|||
try { |
|||
applyPatchIfNeeded(); |
|||
} catch (Exception e) { |
|||
log.error("Failed to apply system data patch updates", e); |
|||
} finally { |
|||
executor.shutdown(); |
|||
} |
|||
}); |
|||
} |
|||
|
|||
private void applyPatchIfNeeded() { |
|||
boolean skipVersionCheck = DefaultDataUpdateService.getEnv("SKIP_PATCH_VERSION_CHECK", false); |
|||
if (!skipVersionCheck && !isVersionChanged()) { |
|||
return; |
|||
} |
|||
|
|||
if (!acquireAdvisoryLock()) { |
|||
log.trace("Could not acquire advisory lock. Another node is processing patch updates."); |
|||
return; |
|||
} |
|||
|
|||
try { |
|||
int updated = updateWidgetTypes(); |
|||
log.info("Updated {} widget types", updated); |
|||
|
|||
schemaSettingsService.updateSchemaVersion(); |
|||
log.info("System data patch update completed successfully"); |
|||
|
|||
} finally { |
|||
releaseAdvisoryLock(); |
|||
} |
|||
} |
|||
|
|||
private boolean isVersionChanged() { |
|||
String packageVersion = schemaSettingsService.getPackageSchemaVersion(); |
|||
String dbVersion = schemaSettingsService.getDbSchemaVersion(); |
|||
|
|||
log.trace("Package version: {}, DB schema version: {}", packageVersion, dbVersion); |
|||
|
|||
VersionInfo packageVersionInfo = parseVersion(packageVersion); |
|||
VersionInfo dbVersionInfo = parseVersion(dbVersion); |
|||
|
|||
if (packageVersionInfo == null || dbVersionInfo == null) { |
|||
log.warn("Unable to parse versions. Package: {}, DB: {}", packageVersion, dbVersion); |
|||
return false; |
|||
} |
|||
|
|||
if (!isPatchVersionChanged(packageVersionInfo, dbVersionInfo)) { |
|||
return false; |
|||
} |
|||
|
|||
log.info("Patch version increased from {} to {}. Starting system data update.", dbVersion, packageVersion); |
|||
return true; |
|||
} |
|||
|
|||
private boolean isPatchVersionChanged(VersionInfo packageVersion, VersionInfo dbVersion) { |
|||
return packageVersion.major == dbVersion.major && packageVersion.minor == dbVersion.minor |
|||
&& packageVersion.maintenance == dbVersion.maintenance && packageVersion.patch > dbVersion.patch; |
|||
} |
|||
|
|||
private int updateWidgetTypes() { |
|||
AtomicInteger updated = new AtomicInteger(); |
|||
Path widgetTypesDir = installScripts.getWidgetTypesDir(); |
|||
|
|||
if (!Files.exists(widgetTypesDir)) { |
|||
log.trace("Widget types directory does not exist: {}", widgetTypesDir); |
|||
return 0; |
|||
} |
|||
|
|||
try (Stream<Path> dirStream = listDir(widgetTypesDir).filter(path -> path.toString().endsWith(InstallScripts.JSON_EXT))) { |
|||
dirStream.forEach( |
|||
path -> { |
|||
try { |
|||
if (updateWidgetTypeFromFile(path)) { |
|||
updated.incrementAndGet(); |
|||
} |
|||
} catch (Exception e) { |
|||
log.error("Unable to update widget type from json: [{}]", path.toString()); |
|||
throw new RuntimeException("Unable to update widget type from json", e); |
|||
} |
|||
} |
|||
); |
|||
} |
|||
|
|||
return updated.get(); |
|||
} |
|||
|
|||
private boolean updateWidgetTypeFromFile(Path filePath) { |
|||
JsonNode json = JacksonUtil.toJsonNode(filePath.toFile()); |
|||
WidgetTypeDetails fileWidgetType = JacksonUtil.treeToValue(json, WidgetTypeDetails.class); |
|||
String fqn = fileWidgetType.getFqn(); |
|||
|
|||
WidgetTypeDetails existingWidgetType = widgetTypeService.findWidgetTypeDetailsByTenantIdAndFqn(TenantId.SYS_TENANT_ID, fqn); |
|||
if (existingWidgetType == null) { |
|||
// We expect only update here, so it's probably never happening, but for test purpose leave it like this:
|
|||
throw new RuntimeException("Widget type not found: " + fqn); |
|||
} |
|||
if (isWidgetTypeChanged(existingWidgetType, fileWidgetType)) { |
|||
existingWidgetType.setDescription(fileWidgetType.getDescription()); |
|||
existingWidgetType.setName(fileWidgetType.getName()); |
|||
existingWidgetType.setDescriptor(fileWidgetType.getDescriptor()); |
|||
widgetTypeService.saveWidgetType(existingWidgetType); |
|||
log.trace("Updated widget type: {}", fqn); |
|||
return true; |
|||
} |
|||
|
|||
log.trace("Widget type unchanged: {}", fqn); |
|||
return false; |
|||
} |
|||
|
|||
private boolean isWidgetTypeChanged(WidgetTypeDetails existing, WidgetTypeDetails file) { |
|||
if (!isDescriptorEqual(existing.getDescriptor(), file.getDescriptor())) { |
|||
return true; |
|||
} |
|||
|
|||
if (!Objects.equals(existing.getName(), file.getName())) { |
|||
return true; |
|||
} |
|||
|
|||
return !Objects.equals(existing.getDescription(), file.getDescription()); |
|||
} |
|||
|
|||
private boolean isDescriptorEqual(JsonNode desc1, JsonNode desc2) { |
|||
if (desc1 == null && desc2 == null) { |
|||
return true; |
|||
} |
|||
if (desc1 == null || desc2 == null) { |
|||
return false; |
|||
} |
|||
|
|||
try { |
|||
String hash1 = computeChecksum(desc1); |
|||
String hash2 = computeChecksum(desc2); |
|||
return Objects.equals(hash1, hash2); |
|||
} catch (Exception e) { |
|||
log.warn("Failed to compare descriptors using checksum, falling back to equals", e); |
|||
return desc1.equals(desc2); |
|||
} |
|||
} |
|||
|
|||
private String computeChecksum(JsonNode node) { |
|||
String canonicalString = JacksonUtil.toCanonicalString(node); |
|||
if (canonicalString == null) { |
|||
return null; |
|||
} |
|||
return Hashing.sha256().hashBytes(canonicalString.getBytes()).toString(); |
|||
} |
|||
|
|||
private boolean acquireAdvisoryLock() { |
|||
try { |
|||
Boolean acquired = jdbcTemplate.queryForObject( |
|||
"SELECT pg_try_advisory_lock(?)", |
|||
Boolean.class, |
|||
ADVISORY_LOCK_ID |
|||
); |
|||
if (Boolean.TRUE.equals(acquired)) { |
|||
log.trace("Acquired advisory lock"); |
|||
return true; |
|||
} |
|||
return false; |
|||
} catch (Exception e) { |
|||
log.error("Failed to acquire advisory lock", e); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
private void releaseAdvisoryLock() { |
|||
try { |
|||
jdbcTemplate.queryForObject( |
|||
"SELECT pg_advisory_unlock(?)", |
|||
Boolean.class, |
|||
ADVISORY_LOCK_ID |
|||
); |
|||
log.debug("Released advisory lock"); |
|||
} catch (Exception e) { |
|||
log.error("Failed to release advisory lock", e); |
|||
} |
|||
} |
|||
|
|||
private VersionInfo parseVersion(String version) { |
|||
try { |
|||
String[] parts = version.split("\\."); |
|||
int major = Integer.parseInt(parts[0]); |
|||
int minor = parts.length > 1 ? Integer.parseInt(parts[1]) : 0; |
|||
int maintenance = parts.length > 2 ? Integer.parseInt(parts[2]) : 0; |
|||
int patch = parts.length > 3 ? Integer.parseInt(parts[3]) : 0; |
|||
return new VersionInfo(major, minor, maintenance, patch); |
|||
} catch (Exception e) { |
|||
log.error("Failed to parse version: {}", version, e); |
|||
return null; |
|||
} |
|||
} |
|||
|
|||
private Stream<Path> listDir(Path dir) { |
|||
try { |
|||
return Files.list(dir); |
|||
} catch (NoSuchFileException e) { |
|||
return Stream.empty(); |
|||
} catch (IOException e) { |
|||
throw new UncheckedIOException(e); |
|||
} |
|||
} |
|||
|
|||
public record VersionInfo(int major, int minor, int maintenance, int patch) {} |
|||
|
|||
} |
|||
@ -0,0 +1,255 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import org.junit.After; |
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.springframework.test.annotation.DirtiesContext; |
|||
import org.springframework.test.context.TestPropertySource; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.cf.CalculatedField; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Argument; |
|||
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Output; |
|||
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
|||
import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.AggInterval; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.CustomInterval; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.Watermark; |
|||
import org.thingsboard.server.common.data.debug.DebugSettings; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.security.Authority; |
|||
import org.thingsboard.server.controller.AbstractControllerTest; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
import java.util.HashMap; |
|||
import java.util.Map; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.awaitility.Awaitility.await; |
|||
import static org.thingsboard.server.cf.CalculatedFieldIntegrationTest.POLL_INTERVAL; |
|||
|
|||
@DaoSqlTest |
|||
@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_EACH_TEST_METHOD) |
|||
@TestPropertySource(properties = { |
|||
"actors.calculated_fields.check_interval=1" |
|||
}) |
|||
public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest { |
|||
|
|||
private Tenant savedTenant; |
|||
|
|||
@Before |
|||
public void beforeEach() throws Exception { |
|||
loginSysAdmin(); |
|||
|
|||
updateDefaultTenantProfileConfig(tenantProfileConfig -> { |
|||
tenantProfileConfig.setMinAllowedDeduplicationIntervalInSecForCF(1); |
|||
tenantProfileConfig.setMinAllowedAggregationIntervalInSecForCF(1); |
|||
}); |
|||
|
|||
Tenant tenant = new Tenant(); |
|||
tenant.setTitle("My tenant"); |
|||
savedTenant = saveTenant(tenant); |
|||
assertThat(savedTenant).isNotNull(); |
|||
|
|||
User tenantAdmin = new User(); |
|||
tenantAdmin.setAuthority(Authority.TENANT_ADMIN); |
|||
tenantAdmin.setTenantId(savedTenant.getId()); |
|||
tenantAdmin.setEmail("tenant@thingsboard.org"); |
|||
tenantAdmin.setFirstName("John"); |
|||
tenantAdmin.setLastName("Doe"); |
|||
|
|||
createUserAndLogin(tenantAdmin, "testPassword"); |
|||
} |
|||
|
|||
@After |
|||
public void afterTest() throws Exception { |
|||
loginSysAdmin(); |
|||
|
|||
deleteTenant(savedTenant.getId()); |
|||
} |
|||
|
|||
@Test |
|||
public void testCreateCfAndNoTelemetryDuringInterval_checkAggregation() throws Exception { |
|||
Device device = createDevice("Device", "1234567890111"); |
|||
|
|||
CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); |
|||
long intervalEndTs = customInterval.getCurrentIntervalEndTs(); |
|||
|
|||
CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, null); |
|||
long interval = customInterval.getCurrentIntervalDurationMillis(); |
|||
|
|||
await().alias("create CF and no telemetry during interval -> save metric with default value") |
|||
.atMost(2 * interval, TimeUnit.MILLISECONDS) |
|||
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) |
|||
.untilAsserted(() -> { |
|||
ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); |
|||
assertThat(result).isNotNull(); |
|||
assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("9999"); |
|||
}); |
|||
} |
|||
|
|||
@Test |
|||
public void testCreateCfWithoutWatermark_checkAggregation() throws Exception { |
|||
Device device = createDevice("Device", "1234567890111"); |
|||
|
|||
CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); |
|||
long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); |
|||
long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); |
|||
|
|||
long tsBeforeInterval = currentIntervalStartTs - 1000; |
|||
long tsInInterval_1 = currentIntervalStartTs + 1000; |
|||
long tsInInterval_2 = currentIntervalStartTs + 500; |
|||
long tsInInterval_3 = currentIntervalStartTs + 200; |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsBeforeInterval)); |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":100}}", tsInInterval_1)); |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); |
|||
|
|||
long interval = customInterval.getCurrentIntervalDurationMillis(); |
|||
CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, null); |
|||
|
|||
await().alias("create CF -> perform aggregation after interval end") |
|||
.atMost(2 * interval, TimeUnit.MILLISECONDS) |
|||
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) |
|||
.untilAsserted(() -> { |
|||
ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); |
|||
assertThat(result).isNotNull(); |
|||
assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); |
|||
}); |
|||
|
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":500}}", tsInInterval_1)); |
|||
|
|||
await().alias("update telemetry that belongs to previous interval -> no aggregation since watermark is not set ") |
|||
.atMost(2 * interval, TimeUnit.MILLISECONDS) |
|||
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) |
|||
.untilAsserted(() -> { |
|||
ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); |
|||
assertThat(result).isNotNull(); |
|||
assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); |
|||
}); |
|||
} |
|||
|
|||
@Test |
|||
public void testCreateCfWithWatermark_checkAggregationDuringWatermark() throws Exception { |
|||
Device device = createDevice("Device", "1234567890111"); |
|||
|
|||
CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); |
|||
long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); |
|||
long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); |
|||
|
|||
long tsBeforeInterval = currentIntervalStartTs - 1000L; |
|||
long tsInInterval_1 = currentIntervalStartTs + 1000L; |
|||
long tsInInterval_2 = currentIntervalStartTs + 500L; |
|||
long tsInInterval_3 = currentIntervalStartTs + 200L; |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsBeforeInterval)); |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":100}}", tsInInterval_1)); |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); |
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); |
|||
|
|||
long interval = customInterval.getCurrentIntervalDurationMillis(); |
|||
Watermark watermark = new Watermark(10); |
|||
CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); |
|||
|
|||
await().alias("create CF -> perform aggregation after interval end") |
|||
.atMost(2 * interval, TimeUnit.MILLISECONDS) |
|||
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) |
|||
.untilAsserted(() -> { |
|||
ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); |
|||
assertThat(result).isNotNull(); |
|||
assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); |
|||
}); |
|||
|
|||
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":300}}", tsInInterval_1)); |
|||
|
|||
await().alias("update telemetry during watermark -> perform aggregation") |
|||
.atMost(2 * 10, TimeUnit.SECONDS) |
|||
.pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) |
|||
.untilAsserted(() -> { |
|||
ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); |
|||
assertThat(result).isNotNull(); |
|||
assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("600"); |
|||
}); |
|||
} |
|||
|
|||
private CalculatedField createTotalConsumptionCF(EntityId entityId, AggInterval aggInterval, Watermark watermark) { |
|||
Map<String, Argument> arguments = new HashMap<>(); |
|||
Argument argument = new Argument(); |
|||
argument.setRefEntityKey(new ReferencedEntityKey("energy", ArgumentType.TS_LATEST, null)); |
|||
arguments.put("en", argument); |
|||
|
|||
Map<String, AggMetric> aggMetrics = new HashMap<>(); |
|||
|
|||
AggMetric consumption = new AggMetric(); |
|||
consumption.setFunction(AggFunction.SUM); |
|||
consumption.setInput(new AggKeyInput("en")); |
|||
consumption.setDefaultValue(9999L); |
|||
aggMetrics.put("consumption", consumption); |
|||
|
|||
Output output = new Output(); |
|||
output.setType(OutputType.TIME_SERIES); |
|||
output.setDecimalsByDefault(0); |
|||
|
|||
return createAggCf("Consumption per minute", entityId, |
|||
aggInterval, |
|||
watermark, |
|||
arguments, |
|||
aggMetrics, |
|||
output); |
|||
} |
|||
|
|||
private CalculatedField createAggCf(String name, |
|||
EntityId entityId, |
|||
AggInterval aggInterval, |
|||
Watermark watermark, |
|||
Map<String, Argument> inputs, |
|||
Map<String, AggMetric> metrics, |
|||
Output output) { |
|||
CalculatedField calculatedField = new CalculatedField(); |
|||
calculatedField.setName(name); |
|||
calculatedField.setEntityId(entityId); |
|||
calculatedField.setType(CalculatedFieldType.ENTITY_AGGREGATION); |
|||
|
|||
EntityAggregationCalculatedFieldConfiguration configuration = new EntityAggregationCalculatedFieldConfiguration(); |
|||
|
|||
configuration.setArguments(inputs); |
|||
configuration.setMetrics(metrics); |
|||
configuration.setInterval(aggInterval); |
|||
if (watermark != null) { |
|||
configuration.setWatermark(watermark); |
|||
} |
|||
configuration.setOutput(output); |
|||
|
|||
calculatedField.setConfiguration(configuration); |
|||
calculatedField.setDebugSettings(DebugSettings.all()); |
|||
return saveCalculatedField(calculatedField); |
|||
} |
|||
|
|||
private ObjectNode getLatestTelemetry(EntityId entityId, String... keys) throws Exception { |
|||
return doGetAsync("/api/plugins/telemetry/" + entityId.getEntityType() + "/" + entityId.getId() + "/values/timeseries?keys=" + String.join(",", keys), ObjectNode.class); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,410 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.system; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.junit.jupiter.api.io.TempDir; |
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.Arguments; |
|||
import org.junit.jupiter.params.provider.CsvSource; |
|||
import org.junit.jupiter.params.provider.MethodSource; |
|||
import org.mockito.InjectMocks; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.springframework.jdbc.core.JdbcTemplate; |
|||
import org.springframework.test.util.ReflectionTestUtils; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.id.WidgetTypeId; |
|||
import org.thingsboard.server.common.data.widget.WidgetTypeDetails; |
|||
import org.thingsboard.server.dao.widget.WidgetTypeService; |
|||
import org.thingsboard.server.service.install.InstallScripts; |
|||
import org.thingsboard.server.service.system.SystemPatchApplier; |
|||
|
|||
import java.nio.file.Files; |
|||
import java.nio.file.Path; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.CountDownLatch; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.atomic.AtomicBoolean; |
|||
import java.util.stream.Stream; |
|||
|
|||
import static org.junit.jupiter.api.Assertions.assertEquals; |
|||
import static org.junit.jupiter.api.Assertions.assertFalse; |
|||
import static org.junit.jupiter.api.Assertions.assertNotEquals; |
|||
import static org.junit.jupiter.api.Assertions.assertNotNull; |
|||
import static org.junit.jupiter.api.Assertions.assertNull; |
|||
import static org.junit.jupiter.api.Assertions.assertThrows; |
|||
import static org.junit.jupiter.api.Assertions.assertTrue; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.anyLong; |
|||
import static org.mockito.ArgumentMatchers.anyString; |
|||
import static org.mockito.ArgumentMatchers.argThat; |
|||
import static org.mockito.ArgumentMatchers.contains; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class SystemPatchApplierTest { |
|||
|
|||
@Mock |
|||
private JdbcTemplate jdbcTemplate; |
|||
|
|||
@Mock |
|||
private InstallScripts installScripts; |
|||
|
|||
@Mock |
|||
private WidgetTypeService widgetTypeService; |
|||
|
|||
@InjectMocks |
|||
private SystemPatchApplier reconciler; |
|||
|
|||
@TempDir |
|||
Path tempDir; |
|||
|
|||
@ParameterizedTest(name = "Parse version {0} should return major={1}, minor={2}, patch={3}") |
|||
@CsvSource({ |
|||
"4.2.1, 4, 2, 1, 0", |
|||
"4.2.0, 4, 2, 0, 0", |
|||
"4.2, 4, 2, 0, 0", |
|||
"4.0.1.2, 4, 0, 1, 2", |
|||
"4, 4, 0, 0, 0", |
|||
"1.0.5.7, 1, 0, 5, 7", |
|||
"10.20.30.40, 10, 20, 30, 40", |
|||
"0.0.1, 0, 0, 1, 0" |
|||
}) |
|||
void testParseVersion(String versionString, int expectedMajor, int expectedMinor, int expectedMaintenance, int expectedPatch) { |
|||
SystemPatchApplier.VersionInfo version = ReflectionTestUtils.invokeMethod(reconciler, "parseVersion", versionString); |
|||
|
|||
assertNotNull(version, "Version should not be null for: " + versionString); |
|||
assertEquals(expectedMajor, version.major(), "Major version mismatch"); |
|||
assertEquals(expectedMinor, version.minor(), "Minor version mismatch"); |
|||
assertEquals(expectedMaintenance, version.maintenance(), "Maintenance version mismatch"); |
|||
assertEquals(expectedPatch, version.patch(), "Patch version mismatch"); |
|||
} |
|||
|
|||
@ParameterizedTest(name = "Parse invalid version: {0}") |
|||
@CsvSource({ |
|||
"invalid", |
|||
"a.b.c", |
|||
"1.2.y.x", |
|||
"''", |
|||
"1.x.3" |
|||
}) |
|||
void testParseInvalidVersion(String invalidVersion) { |
|||
SystemPatchApplier.VersionInfo version = ReflectionTestUtils.invokeMethod(reconciler, "parseVersion", invalidVersion); |
|||
assertNull(version, "Version should be null for invalid input: " + invalidVersion); |
|||
} |
|||
|
|||
@Test |
|||
void whenLockIsNotAcquired_thenAcquiredIsSuccess() { |
|||
when(jdbcTemplate.queryForObject(anyString(), eq(Boolean.class), anyLong())).thenReturn(true); |
|||
|
|||
Boolean acquired = ReflectionTestUtils.invokeMethod(reconciler, "acquireAdvisoryLock"); |
|||
|
|||
assertEquals(Boolean.TRUE, acquired); |
|||
verify(jdbcTemplate).queryForObject(contains("pg_try_advisory_lock"), eq(Boolean.class), anyLong()); |
|||
} |
|||
|
|||
@Test |
|||
void whenLockIsAlreadyAcquired_thenAcquiredIsFailed() { |
|||
when(jdbcTemplate.queryForObject(anyString(), eq(Boolean.class), anyLong())).thenReturn(false); |
|||
|
|||
Boolean acquired = ReflectionTestUtils.invokeMethod(reconciler, "acquireAdvisoryLock"); |
|||
|
|||
assertNotEquals(Boolean.TRUE, acquired); |
|||
} |
|||
|
|||
@Test |
|||
void testReleaseAdvisoryLock() { |
|||
when(jdbcTemplate.queryForObject(anyString(), eq(Boolean.class), anyLong())) |
|||
.thenReturn(true); |
|||
|
|||
ReflectionTestUtils.invokeMethod(reconciler, "releaseAdvisoryLock"); |
|||
|
|||
verify(jdbcTemplate).queryForObject( |
|||
contains("pg_advisory_unlock"), eq(Boolean.class), anyLong()); |
|||
} |
|||
|
|||
@Test |
|||
void whenWidgetNotFound_thenThrowException() throws Exception { |
|||
Path widgetTypesDir = tempDir.resolve("widget_types"); |
|||
Files.createDirectories(widgetTypesDir); |
|||
when(installScripts.getWidgetTypesDir()).thenReturn(widgetTypesDir); |
|||
|
|||
WidgetTypeDetails testWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
String json = JacksonUtil.toString(testWidget); |
|||
assertNotNull(json); |
|||
Files.writeString(widgetTypesDir.resolve("test_widget.json"), json); |
|||
|
|||
when(widgetTypeService.findWidgetTypeDetailsByTenantIdAndFqn(TenantId.SYS_TENANT_ID, "test_widget")).thenReturn(null); |
|||
|
|||
assertThrows(RuntimeException.class, () -> ReflectionTestUtils.invokeMethod(reconciler, "updateWidgetTypes")); |
|||
} |
|||
|
|||
@Test |
|||
void whenDescriptorChanged_thenUpdateTheExistingWidget() throws Exception { |
|||
Path widgetTypesDir = tempDir.resolve("widget_types"); |
|||
Files.createDirectories(widgetTypesDir); |
|||
when(installScripts.getWidgetTypesDir()).thenReturn(widgetTypesDir); |
|||
|
|||
WidgetTypeDetails fileWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
fileWidget.setDescriptor(JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":2}")); |
|||
String json = JacksonUtil.toString(fileWidget); |
|||
assertNotNull(json); |
|||
Files.writeString(widgetTypesDir.resolve("test_widget.json"), json); |
|||
|
|||
WidgetTypeDetails existingWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
existingWidget.setId(new WidgetTypeId(UUID.randomUUID())); |
|||
existingWidget.setDescriptor(JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":1}")); |
|||
|
|||
when(widgetTypeService.findWidgetTypeDetailsByTenantIdAndFqn(TenantId.SYS_TENANT_ID, "test_widget")) |
|||
.thenReturn(existingWidget); |
|||
|
|||
Integer updated = ReflectionTestUtils.invokeMethod(reconciler, "updateWidgetTypes"); |
|||
|
|||
assertEquals(1, updated); |
|||
verify(widgetTypeService).saveWidgetType(argThat(w -> |
|||
w.getDescriptor().get("version").asInt() == 2 |
|||
)); |
|||
} |
|||
|
|||
@Test |
|||
void whenNameChanged_thenUpdateTheExistingWidget() throws Exception { |
|||
Path widgetTypesDir = tempDir.resolve("widget_types"); |
|||
Files.createDirectories(widgetTypesDir); |
|||
when(installScripts.getWidgetTypesDir()).thenReturn(widgetTypesDir); |
|||
|
|||
WidgetTypeDetails fileWidget = createTestWidgetType("test_widget", "New Name"); |
|||
String json = JacksonUtil.toString(fileWidget); |
|||
assertNotNull(json); |
|||
Files.writeString(widgetTypesDir.resolve("test_widget.json"), json); |
|||
|
|||
WidgetTypeDetails existingWidget = createTestWidgetType("test_widget", "Old Name"); |
|||
existingWidget.setId(new WidgetTypeId(UUID.randomUUID())); |
|||
|
|||
when(widgetTypeService.findWidgetTypeDetailsByTenantIdAndFqn(TenantId.SYS_TENANT_ID, "test_widget")) |
|||
.thenReturn(existingWidget); |
|||
|
|||
Integer updated = ReflectionTestUtils.invokeMethod(reconciler, "updateWidgetTypes"); |
|||
|
|||
assertEquals(1, updated); |
|||
verify(widgetTypeService).saveWidgetType(argThat(w -> "New Name".equals(w.getName()))); |
|||
} |
|||
|
|||
@Test |
|||
void whenNothingChanged_thenSkipTheUpdateOfTheExistingWidget() throws Exception { |
|||
Path widgetTypesDir = tempDir.resolve("widget_types"); |
|||
Files.createDirectories(widgetTypesDir); |
|||
when(installScripts.getWidgetTypesDir()).thenReturn(widgetTypesDir); |
|||
|
|||
WidgetTypeDetails fileWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
String json = JacksonUtil.toString(fileWidget); |
|||
assertNotNull(json); |
|||
Files.writeString(widgetTypesDir.resolve("test_widget.json"), json); |
|||
|
|||
WidgetTypeDetails existingWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
existingWidget.setId(new WidgetTypeId(UUID.randomUUID())); |
|||
|
|||
when(widgetTypeService.findWidgetTypeDetailsByTenantIdAndFqn(TenantId.SYS_TENANT_ID, "test_widget")) |
|||
.thenReturn(existingWidget); |
|||
|
|||
Integer updated = ReflectionTestUtils.invokeMethod(reconciler, "updateWidgetTypes"); |
|||
|
|||
assertEquals(0, updated); |
|||
verify(widgetTypeService, never()).saveWidgetType(any()); |
|||
} |
|||
|
|||
@ParameterizedTest(name = "{0}") |
|||
@MethodSource("provideDescriptorComparisonTestCases") |
|||
void testIfDescriptorsAreEqual(String testName, JsonNode desc1, JsonNode desc2, boolean expectedEqual) { |
|||
Boolean result = ReflectionTestUtils.invokeMethod(reconciler, "isDescriptorEqual", desc1, desc2); |
|||
assertEquals(expectedEqual, result, testName); |
|||
} |
|||
|
|||
@Test |
|||
void whenDescriptorChanged_thenReturnWidgetTypeChanged() { |
|||
WidgetTypeDetails existing = createTestWidgetType("test", "Test"); |
|||
existing.setDescriptor(JacksonUtil.toJsonNode("{\"version\":1}")); |
|||
|
|||
WidgetTypeDetails file = createTestWidgetType("test", "Test"); |
|||
file.setDescriptor(JacksonUtil.toJsonNode("{\"version\":2}")); |
|||
|
|||
boolean result = Boolean.TRUE.equals(ReflectionTestUtils.invokeMethod(reconciler, "isWidgetTypeChanged", existing, file)); |
|||
assertTrue(result); |
|||
} |
|||
|
|||
@Test |
|||
void whenNameChanged_thenReturnWidgetTypeChanged() { |
|||
WidgetTypeDetails existing = createTestWidgetType("test", "Old Name"); |
|||
WidgetTypeDetails file = createTestWidgetType("test", "New Name"); |
|||
|
|||
boolean result = Boolean.TRUE.equals(ReflectionTestUtils.invokeMethod(reconciler, "isWidgetTypeChanged", existing, file)); |
|||
assertTrue(result); |
|||
} |
|||
|
|||
@Test |
|||
void whenDescriptionChanged_thenReturnWidgetTypeChanged() { |
|||
WidgetTypeDetails existing = createTestWidgetType("test", "Test"); |
|||
existing.setDescription("Old description"); |
|||
|
|||
WidgetTypeDetails file = createTestWidgetType("test", "Test"); |
|||
file.setDescription("New description"); |
|||
|
|||
boolean result = Boolean.TRUE.equals(ReflectionTestUtils.invokeMethod(reconciler, "isWidgetTypeChanged", existing, file)); |
|||
assertTrue(result); |
|||
} |
|||
|
|||
@Test |
|||
void whenWidgetTypeAreIdentical_thenNoUpdateIsPerformed() { |
|||
WidgetTypeDetails existing = createTestWidgetType("test", "Test"); |
|||
WidgetTypeDetails file = createTestWidgetType("test", "Test"); |
|||
|
|||
boolean result = Boolean.TRUE.equals(ReflectionTestUtils.invokeMethod(reconciler, "isWidgetTypeChanged", existing, file)); |
|||
assertFalse(result); |
|||
} |
|||
|
|||
@Test |
|||
void whenLockIsHeldByOneThread_thenSecondThreadCannotAcquireLock() throws Exception { |
|||
CountDownLatch lockAcquired = new CountDownLatch(1); |
|||
CountDownLatch startSecondThread = new CountDownLatch(1); |
|||
CountDownLatch testComplete = new CountDownLatch(1); |
|||
|
|||
AtomicBoolean firstThreadAcquiredLock = new AtomicBoolean(false); |
|||
AtomicBoolean secondThreadAcquiredLock = new AtomicBoolean(false); |
|||
AtomicBoolean firstThreadSavedWidget = new AtomicBoolean(false); |
|||
AtomicBoolean secondThreadSavedWidget = new AtomicBoolean(false); |
|||
|
|||
Path widgetTypesDir = tempDir.resolve("widget_types"); |
|||
Files.createDirectories(widgetTypesDir); |
|||
when(installScripts.getWidgetTypesDir()).thenReturn(widgetTypesDir); |
|||
|
|||
WidgetTypeDetails fileWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
fileWidget.setDescriptor(JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":2}")); |
|||
String toString = JacksonUtil.toCanonicalString(fileWidget); |
|||
assertNotNull(toString); |
|||
Files.writeString(widgetTypesDir.resolve("test_widget.json"), toString); |
|||
|
|||
WidgetTypeDetails existingWidget = createTestWidgetType("test_widget", "Test Widget"); |
|||
existingWidget.setId(new WidgetTypeId(UUID.randomUUID())); |
|||
existingWidget.setDescriptor(JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":1}")); |
|||
|
|||
when(widgetTypeService.findWidgetTypeDetailsByTenantIdAndFqn(TenantId.SYS_TENANT_ID, "test_widget")).thenReturn(existingWidget); |
|||
|
|||
when(jdbcTemplate.queryForObject(contains("pg_try_advisory_lock"), eq(Boolean.class), anyLong())) |
|||
.thenReturn(true) |
|||
.thenReturn(false); |
|||
|
|||
when(jdbcTemplate.queryForObject(contains("pg_advisory_unlock"), eq(Boolean.class), anyLong())) |
|||
.thenReturn(true); |
|||
|
|||
// The first thread-acquires lock and performs update
|
|||
Thread firstThread = new Thread(() -> { |
|||
try { |
|||
Boolean acquired = ReflectionTestUtils.invokeMethod(reconciler, "acquireAdvisoryLock"); |
|||
firstThreadAcquiredLock.set(Boolean.TRUE.equals(acquired)); |
|||
|
|||
if (firstThreadAcquiredLock.get()) { |
|||
lockAcquired.countDown(); |
|||
startSecondThread.await(5, TimeUnit.SECONDS); |
|||
|
|||
// Simulate work while holding lock
|
|||
Thread.sleep(100); |
|||
|
|||
Integer updated = ReflectionTestUtils.invokeMethod(reconciler, "updateWidgetTypes"); |
|||
firstThreadSavedWidget.set(updated != null && updated > 0); |
|||
|
|||
ReflectionTestUtils.invokeMethod(reconciler, "releaseAdvisoryLock"); |
|||
} |
|||
} catch (Exception ignored) { |
|||
} finally { |
|||
testComplete.countDown(); |
|||
} |
|||
}); |
|||
|
|||
// Second thread - attempts to acquire lock but fails
|
|||
Thread secondThread = new Thread(() -> { |
|||
try { |
|||
lockAcquired.await(5, TimeUnit.SECONDS); |
|||
startSecondThread.countDown(); |
|||
|
|||
Boolean acquired = ReflectionTestUtils.invokeMethod(reconciler, "acquireAdvisoryLock"); |
|||
secondThreadAcquiredLock.set(Boolean.TRUE.equals(acquired)); |
|||
|
|||
if (secondThreadAcquiredLock.get()) { |
|||
Integer updated = ReflectionTestUtils.invokeMethod(reconciler, "updateWidgetTypes"); |
|||
secondThreadSavedWidget.set(updated != null && updated > 0); |
|||
|
|||
ReflectionTestUtils.invokeMethod(reconciler, "releaseAdvisoryLock"); |
|||
} |
|||
} catch (Exception ignored) {} |
|||
}); |
|||
|
|||
firstThread.start(); |
|||
secondThread.start(); |
|||
|
|||
assertTrue(testComplete.await(10, TimeUnit.SECONDS), "Test should complete within timeout"); |
|||
firstThread.join(1000); |
|||
secondThread.join(1000); |
|||
|
|||
assertTrue(firstThreadAcquiredLock.get(), "First thread should acquire lock"); |
|||
assertFalse(secondThreadAcquiredLock.get(), "Second thread should NOT acquire lock"); |
|||
assertTrue(firstThreadSavedWidget.get(), "First thread should save widget"); |
|||
assertFalse(secondThreadSavedWidget.get(), "Second thread should NOT save widget"); |
|||
|
|||
verify(widgetTypeService, times(1)).saveWidgetType(any()); |
|||
} |
|||
|
|||
private static Stream<Arguments> provideDescriptorComparisonTestCases() { |
|||
return Stream.of( |
|||
Arguments.of("Both null", null, null, true), |
|||
Arguments.of("First null", null, JacksonUtil.newObjectNode(), false), |
|||
Arguments.of("Second null", JacksonUtil.newObjectNode(), null, false), |
|||
Arguments.of("Same content", |
|||
JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":1}"), |
|||
JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":1}"), |
|||
true), |
|||
Arguments.of("Different content", |
|||
JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":1}"), |
|||
JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":2}"), |
|||
false), |
|||
Arguments.of("Different key order but same content", |
|||
JacksonUtil.toJsonNode("{\"version\":1,\"type\":\"latest\"}"), |
|||
JacksonUtil.toJsonNode("{\"type\":\"latest\",\"version\":1}"), |
|||
true), |
|||
Arguments.of("Empty objects", |
|||
JacksonUtil.toJsonNode("{}"), |
|||
JacksonUtil.toJsonNode("{}"), |
|||
true) |
|||
); |
|||
} |
|||
|
|||
private WidgetTypeDetails createTestWidgetType(String fqn, String name) { |
|||
WidgetTypeDetails widget = new WidgetTypeDetails(); |
|||
widget.setFqn(fqn); |
|||
widget.setName(name); |
|||
widget.setDescription("Test description"); |
|||
widget.setTenantId(TenantId.SYS_TENANT_ID); |
|||
widget.setDescriptor(JacksonUtil.toJsonNode("{\"type\":\"latest\"}")); |
|||
return widget; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,40 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.transport.lwm2m.security.cid; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.test.context.TestPropertySource; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
|
|||
@TestPropertySource(properties = { |
|||
"transport.lwm2m.dtls.connection_id_length=16" |
|||
}) |
|||
|
|||
@DaoSqlTest |
|||
@Slf4j |
|||
public abstract class AbstractSecurityLwM2MIntegrationDtlsCidLength16Test extends AbstractSecurityLwM2MIntegrationDtlsCidLengthTest { |
|||
|
|||
private static final Integer serverDtlsCidLength = 16; |
|||
|
|||
protected void testNoSecDtlsCidLength(Integer clientDtlsCidLength) throws Exception { |
|||
testNoSecDtlsCidLength(clientDtlsCidLength, serverDtlsCidLength); |
|||
} |
|||
|
|||
protected void testPskDtlsCidLength(Integer clientDtlsCidLength) throws Exception { |
|||
testPskDtlsCidLength(clientDtlsCidLength, serverDtlsCidLength); |
|||
} |
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.transport.lwm2m.security.cid; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.test.context.TestPropertySource; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
|
|||
@TestPropertySource(properties = { |
|||
"transport.lwm2m.dtls.connection_id_length=2" |
|||
}) |
|||
|
|||
@DaoSqlTest |
|||
@Slf4j |
|||
public abstract class AbstractSecurityLwM2MIntegrationDtlsCidLength2Test extends AbstractSecurityLwM2MIntegrationDtlsCidLengthTest { |
|||
|
|||
private static final Integer serverDtlsCidLength = 2; |
|||
|
|||
protected void testNoSecDtlsCidLength(Integer dtlsCidLength) throws Exception { |
|||
testNoSecDtlsCidLength(dtlsCidLength, serverDtlsCidLength); |
|||
} |
|||
protected void testPskDtlsCidLength(Integer dtlsCidLength) throws Exception { |
|||
testPskDtlsCidLength(dtlsCidLength, serverDtlsCidLength); |
|||
} |
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.transport.lwm2m.security.cid; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.test.context.TestPropertySource; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
|
|||
@TestPropertySource(properties = { |
|||
"transport.lwm2m.dtls.connection_id_length=4" |
|||
}) |
|||
|
|||
@DaoSqlTest |
|||
@Slf4j |
|||
public abstract class AbstractSecurityLwM2MIntegrationDtlsCidLength4Test extends AbstractSecurityLwM2MIntegrationDtlsCidLengthTest { |
|||
|
|||
private static final Integer serverDtlsCidLength = 4; |
|||
|
|||
protected void testNoSecDtlsCidLength(Integer dtlsCidLength) throws Exception { |
|||
testNoSecDtlsCidLength(dtlsCidLength, serverDtlsCidLength); |
|||
} |
|||
protected void testPskDtlsCidLength(Integer dtlsCidLength) throws Exception { |
|||
testPskDtlsCidLength(dtlsCidLength, serverDtlsCidLength); |
|||
} |
|||
} |
|||
@ -0,0 +1,64 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.transport.lwm2m.security.cid.serverDtlsCidLength_1; |
|||
|
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.thingsboard.server.transport.lwm2m.security.cid.AbstractSecurityLwM2MIntegrationDtlsCidLength0Test; |
|||
import org.thingsboard.server.transport.lwm2m.security.cid.AbstractSecurityLwM2MIntegrationDtlsCidLength1Test; |
|||
|
|||
import static org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode.PSK; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE; |
|||
|
|||
public class PskLwm2mIntegrationDtlsCidLengthTest extends AbstractSecurityLwM2MIntegrationDtlsCidLength1Test { |
|||
|
|||
@Before |
|||
public void createProfileRpc() { |
|||
transportConfiguration = getTransportConfiguration(OBSERVE_ATTRIBUTES_WITHOUT_PARAMS, getBootstrapServerCredentialsSecure(PSK, NONE)); |
|||
awaitAlias = "await on client state (Psk_Lwm2m) serverDtlsCidLength = 1"; |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_Null() throws Exception { |
|||
testPskDtlsCidLength(null); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_0() throws Exception { |
|||
testPskDtlsCidLength(0); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_1() throws Exception { |
|||
testPskDtlsCidLength(1); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_2() throws Exception { |
|||
testPskDtlsCidLength(2); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_4() throws Exception { |
|||
testPskDtlsCidLength(4); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_16() throws Exception { |
|||
testPskDtlsCidLength(16); |
|||
} |
|||
} |
|||
|
|||
@ -0,0 +1,64 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.transport.lwm2m.security.cid.serverDtlsCidLength_16; |
|||
|
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.thingsboard.server.transport.lwm2m.security.cid.AbstractSecurityLwM2MIntegrationDtlsCidLength16Test; |
|||
import org.thingsboard.server.transport.lwm2m.security.cid.AbstractSecurityLwM2MIntegrationDtlsCidLength4Test; |
|||
|
|||
import static org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode.PSK; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE; |
|||
|
|||
public class PskLwm2mIntegrationDtlsCidLengthTest extends AbstractSecurityLwM2MIntegrationDtlsCidLength16Test { |
|||
|
|||
@Before |
|||
public void createProfileRpc() { |
|||
transportConfiguration = getTransportConfiguration(OBSERVE_ATTRIBUTES_WITHOUT_PARAMS, getBootstrapServerCredentialsSecure(PSK, NONE)); |
|||
awaitAlias = "await on client state (Psk_Lwm2m) serverDtlsCidLength = 16"; |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_Null() throws Exception { |
|||
testPskDtlsCidLength(null); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_0() throws Exception { |
|||
testPskDtlsCidLength(0); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_1() throws Exception { |
|||
testPskDtlsCidLength(1); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_2() throws Exception { |
|||
testPskDtlsCidLength(2); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_4() throws Exception { |
|||
testPskDtlsCidLength(4); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_16() throws Exception { |
|||
testPskDtlsCidLength(16); |
|||
} |
|||
} |
|||
|
|||
@ -0,0 +1,63 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.transport.lwm2m.security.cid.serverDtlsCidLength_4; |
|||
|
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.thingsboard.server.transport.lwm2m.security.cid.AbstractSecurityLwM2MIntegrationDtlsCidLength4Test; |
|||
|
|||
import static org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode.PSK; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE; |
|||
|
|||
public class PskLwm2mIntegrationDtlsCidLengthTest extends AbstractSecurityLwM2MIntegrationDtlsCidLength4Test { |
|||
|
|||
@Before |
|||
public void createProfileRpc() { |
|||
transportConfiguration = getTransportConfiguration(OBSERVE_ATTRIBUTES_WITHOUT_PARAMS, getBootstrapServerCredentialsSecure(PSK, NONE)); |
|||
awaitAlias = "await on client state (Psk_Lwm2m) serverDtlsCidLength = 4"; |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_Null() throws Exception { |
|||
testPskDtlsCidLength(null); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_0() throws Exception { |
|||
testPskDtlsCidLength(0); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_1() throws Exception { |
|||
testPskDtlsCidLength(1); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_2() throws Exception { |
|||
testPskDtlsCidLength(2); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_4() throws Exception { |
|||
testPskDtlsCidLength(4); |
|||
} |
|||
|
|||
@Test |
|||
public void testWithPskConnectLwm2mSuccessClientDtlsCidLength_16() throws Exception { |
|||
testPskDtlsCidLength(16); |
|||
} |
|||
} |
|||
|
|||
@ -0,0 +1,96 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single; |
|||
|
|||
import jakarta.validation.Valid; |
|||
import jakarta.validation.constraints.NotEmpty; |
|||
import jakarta.validation.constraints.NotNull; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Argument; |
|||
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|||
import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.cf.configuration.Output; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.AggInterval; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.Watermark; |
|||
|
|||
import java.util.Map; |
|||
|
|||
@Data |
|||
public class EntityAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration { |
|||
|
|||
private Map<String, Argument> arguments; |
|||
@Valid |
|||
@NotEmpty |
|||
private Map<String, AggMetric> metrics; |
|||
@Valid |
|||
@NotNull |
|||
private AggInterval interval; |
|||
@Valid |
|||
private Watermark watermark; |
|||
@Valid |
|||
@NotNull |
|||
private Output output; |
|||
|
|||
@Override |
|||
public CalculatedFieldType getType() { |
|||
return CalculatedFieldType.ENTITY_AGGREGATION; |
|||
} |
|||
|
|||
@Override |
|||
public void validate() { |
|||
validateArguments(); |
|||
validateMetrics(); |
|||
validateInterval(); |
|||
} |
|||
|
|||
private void validateArguments() { |
|||
if (arguments.containsKey("ctx")) { |
|||
throw new IllegalArgumentException("Argument name 'ctx' is reserved and cannot be used."); |
|||
} |
|||
if (arguments.values().stream().anyMatch(argument -> !ArgumentType.TS_LATEST.equals(argument.getRefEntityKey().getType()))) { |
|||
throw new IllegalArgumentException("Calculated field with type: '" + getType() + "' support only TS_LATEST arguments."); |
|||
} |
|||
} |
|||
|
|||
private void validateMetrics() { |
|||
if (metrics == null || metrics.isEmpty()) { |
|||
throw new IllegalArgumentException("Metrics map cannot be empty."); |
|||
} |
|||
|
|||
for (AggMetric metric : metrics.values()) { |
|||
if (metric.getInput() instanceof AggKeyInput aggKeyInput) { |
|||
if (!arguments.containsKey(aggKeyInput.getKey())) { |
|||
throw new IllegalArgumentException( |
|||
"Metric references unknown argument: '" + aggKeyInput.getKey() + "'." |
|||
); |
|||
} |
|||
} else { |
|||
throw new IllegalArgumentException("Metric key can only refer to argument."); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private void validateInterval() { |
|||
if (interval == null) { |
|||
throw new IllegalArgumentException("Interval must be defined."); |
|||
} |
|||
interval.validate(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,67 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonIgnore; |
|||
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; |
|||
import com.fasterxml.jackson.annotation.JsonSubTypes; |
|||
import com.fasterxml.jackson.annotation.JsonTypeInfo; |
|||
|
|||
import java.time.ZoneId; |
|||
import java.time.ZonedDateTime; |
|||
|
|||
@JsonTypeInfo( |
|||
use = JsonTypeInfo.Id.NAME, |
|||
include = JsonTypeInfo.As.PROPERTY, |
|||
property = "type" |
|||
) |
|||
@JsonSubTypes({ |
|||
@JsonSubTypes.Type(value = HourInterval.class, name = "HOUR"), |
|||
@JsonSubTypes.Type(value = DayInterval.class, name = "DAY"), |
|||
@JsonSubTypes.Type(value = WeekInterval.class, name = "WEEK"), |
|||
@JsonSubTypes.Type(value = WeekSunSatInterval.class, name = "WEEK_SUN_SAT"), |
|||
@JsonSubTypes.Type(value = MonthInterval.class, name = "MONTH"), |
|||
@JsonSubTypes.Type(value = QuarterInterval.class, name = "QUARTER"), |
|||
@JsonSubTypes.Type(value = YearInterval.class, name = "YEAR"), |
|||
@JsonSubTypes.Type(value = CustomInterval.class, name = "CUSTOM") |
|||
}) |
|||
@JsonIgnoreProperties(ignoreUnknown = true) |
|||
public interface AggInterval { |
|||
|
|||
@JsonIgnore |
|||
AggIntervalType getType(); |
|||
|
|||
@JsonIgnore |
|||
ZoneId getZoneId(); |
|||
|
|||
@JsonIgnore |
|||
long getCurrentIntervalDurationMillis(); |
|||
|
|||
@JsonIgnore |
|||
long getCurrentIntervalStartTs(); |
|||
|
|||
long getDateTimeIntervalStartTs(ZonedDateTime dateTime); |
|||
|
|||
@JsonIgnore |
|||
long getCurrentIntervalEndTs(); |
|||
|
|||
long getDateTimeIntervalEndTs(ZonedDateTime dateTime); |
|||
|
|||
ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart); |
|||
|
|||
void validate(); |
|||
|
|||
} |
|||
@ -0,0 +1,29 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
public enum AggIntervalType { |
|||
|
|||
HOUR, |
|||
DAY, |
|||
WEEK, |
|||
WEEK_SUN_SAT, |
|||
MONTH, |
|||
QUARTER, |
|||
YEAR, |
|||
CUSTOM |
|||
|
|||
} |
|||
@ -0,0 +1,108 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonInclude; |
|||
import jakarta.validation.constraints.NotBlank; |
|||
import lombok.AllArgsConstructor; |
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.ZoneId; |
|||
import java.time.ZonedDateTime; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@Data |
|||
@JsonInclude(JsonInclude.Include.NON_NULL) |
|||
@AllArgsConstructor |
|||
@NoArgsConstructor |
|||
public abstract class BaseAggInterval implements AggInterval { |
|||
|
|||
@NotBlank |
|||
protected String tz; |
|||
protected Long offsetSec; // delay seconds since start of interval
|
|||
|
|||
@Override |
|||
public ZoneId getZoneId() { |
|||
return ZoneId.of(tz); |
|||
} |
|||
|
|||
protected long getOffsetSafe() { |
|||
return offsetSec != null ? offsetSec : 0L; |
|||
} |
|||
|
|||
@Override |
|||
public long getCurrentIntervalDurationMillis() { |
|||
return getCurrentIntervalEndTs() - getCurrentIntervalStartTs(); |
|||
} |
|||
|
|||
@Override |
|||
public long getCurrentIntervalStartTs() { |
|||
ZoneId zoneId = getZoneId(); |
|||
ZonedDateTime now = ZonedDateTime.now(zoneId); |
|||
return getDateTimeIntervalStartTs(now); |
|||
} |
|||
|
|||
@Override |
|||
public long getDateTimeIntervalStartTs(ZonedDateTime dateTime) { |
|||
long offset = getOffsetSafe(); |
|||
ZonedDateTime shiftedNow = dateTime.minusSeconds(offset); |
|||
ZonedDateTime alignedStart = getAlignedBoundary(shiftedNow, false); |
|||
ZonedDateTime actualStart = alignedStart.plusSeconds(offset); |
|||
return actualStart.toInstant().toEpochMilli(); |
|||
} |
|||
|
|||
@Override |
|||
public long getCurrentIntervalEndTs() { |
|||
ZoneId zoneId = getZoneId(); |
|||
ZonedDateTime now = ZonedDateTime.now(zoneId); |
|||
return getDateTimeIntervalEndTs(now); |
|||
} |
|||
|
|||
@Override |
|||
public long getDateTimeIntervalEndTs(ZonedDateTime dateTime) { |
|||
long offset = getOffsetSafe(); |
|||
ZonedDateTime shiftedNow = dateTime.minusSeconds(offset); |
|||
ZonedDateTime alignedEnd = getAlignedBoundary(shiftedNow, true); |
|||
ZonedDateTime actualEnd = alignedEnd.plusSeconds(offset); |
|||
return actualEnd.toInstant().toEpochMilli(); |
|||
} |
|||
|
|||
protected abstract ZonedDateTime alignToIntervalStart(ZonedDateTime reference); |
|||
|
|||
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) { |
|||
ZonedDateTime base = alignToIntervalStart(reference); |
|||
return next ? getNextIntervalStart(base) : base; |
|||
} |
|||
|
|||
@Override |
|||
public void validate() { |
|||
try { |
|||
getZoneId(); |
|||
} catch (Exception ex) { |
|||
throw new IllegalArgumentException("Invalid timezone in interval: " + ex.getMessage()); |
|||
} |
|||
if (offsetSec != null) { |
|||
if (offsetSec < 0) { |
|||
throw new IllegalArgumentException("Offset cannot be negative."); |
|||
} |
|||
if (TimeUnit.SECONDS.toMillis(offsetSec) >= getCurrentIntervalDurationMillis()) { |
|||
throw new IllegalArgumentException("Offset must be greater than interval duration."); |
|||
} |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,68 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import jakarta.validation.constraints.Min; |
|||
import jakarta.validation.constraints.NotNull; |
|||
import lombok.Data; |
|||
import lombok.EqualsAndHashCode; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.Duration; |
|||
import java.time.ZonedDateTime; |
|||
|
|||
@EqualsAndHashCode(callSuper = true) |
|||
@Data |
|||
@NoArgsConstructor |
|||
public class CustomInterval extends BaseAggInterval { |
|||
|
|||
@NotNull |
|||
@Min(1) |
|||
private Long durationSec; |
|||
|
|||
public CustomInterval(String tz, Long offsetSec, Long durationSec) { |
|||
super(tz, offsetSec); |
|||
this.durationSec = durationSec; |
|||
} |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.CUSTOM; |
|||
} |
|||
|
|||
@Override |
|||
public long getCurrentIntervalDurationMillis() { |
|||
return getDurationMillis(); |
|||
} |
|||
|
|||
private long getDurationMillis() { |
|||
return Duration.ofSeconds(durationSec).toMillis(); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
ZonedDateTime localMidnight = reference.toLocalDate().atStartOfDay(reference.getZone()); |
|||
long secondsFromMidnight = Duration.between(localMidnight, reference).getSeconds(); |
|||
long alignedSecondsFromMidnight = (secondsFromMidnight / durationSec) * durationSec; |
|||
return localMidnight.plusSeconds(alignedSecondsFromMidnight); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusSeconds(durationSec); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,47 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.ZonedDateTime; |
|||
import java.time.temporal.ChronoUnit; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class DayInterval extends BaseAggInterval { |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.DAY; |
|||
} |
|||
|
|||
public DayInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
return reference.truncatedTo(ChronoUnit.DAYS); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusDays(1); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,49 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.EqualsAndHashCode; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.ZonedDateTime; |
|||
import java.time.temporal.ChronoUnit; |
|||
|
|||
@EqualsAndHashCode(callSuper = true) |
|||
@Data |
|||
@NoArgsConstructor |
|||
public class HourInterval extends BaseAggInterval { |
|||
|
|||
public HourInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.HOUR; |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
return reference.truncatedTo(ChronoUnit.HOURS); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusHours(1); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,47 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.ZonedDateTime; |
|||
import java.time.temporal.ChronoUnit; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class MonthInterval extends BaseAggInterval { |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.MONTH; |
|||
} |
|||
|
|||
public MonthInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
return reference.withDayOfMonth(1).truncatedTo(ChronoUnit.DAYS); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusMonths(1); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,53 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.LocalDate; |
|||
import java.time.LocalTime; |
|||
import java.time.ZonedDateTime; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class QuarterInterval extends BaseAggInterval { |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.QUARTER; |
|||
} |
|||
|
|||
public QuarterInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
int month = reference.getMonthValue(); |
|||
int quarterStartMonth = ((month - 1) / 3) * 3 + 1; // 1, 4, 7, 10
|
|||
return ZonedDateTime.of( |
|||
LocalDate.of(reference.getYear(), quarterStartMonth, 1), |
|||
LocalTime.MIDNIGHT, |
|||
reference.getZone()); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusMonths(3); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,31 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import jakarta.validation.constraints.Min; |
|||
import lombok.AllArgsConstructor; |
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
@Data |
|||
@AllArgsConstructor |
|||
@NoArgsConstructor |
|||
public class Watermark { |
|||
|
|||
@Min(0) |
|||
private long duration; |
|||
|
|||
} |
|||
@ -0,0 +1,49 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.DayOfWeek; |
|||
import java.time.ZonedDateTime; |
|||
import java.time.temporal.ChronoUnit; |
|||
import java.time.temporal.TemporalAdjusters; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class WeekInterval extends BaseAggInterval { |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.WEEK; |
|||
} |
|||
|
|||
public WeekInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
return reference.with(TemporalAdjusters.previousOrSame(DayOfWeek.MONDAY)).truncatedTo(ChronoUnit.DAYS); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusWeeks(1); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,49 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.DayOfWeek; |
|||
import java.time.ZonedDateTime; |
|||
import java.time.temporal.ChronoUnit; |
|||
import java.time.temporal.TemporalAdjusters; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class WeekSunSatInterval extends BaseAggInterval { |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.WEEK_SUN_SAT; |
|||
} |
|||
|
|||
public WeekSunSatInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
return reference.with(TemporalAdjusters.previousOrSame(DayOfWeek.SUNDAY)).truncatedTo(ChronoUnit.DAYS); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusWeeks(1); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,51 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import lombok.Data; |
|||
import lombok.NoArgsConstructor; |
|||
|
|||
import java.time.LocalDate; |
|||
import java.time.LocalTime; |
|||
import java.time.ZonedDateTime; |
|||
|
|||
@Data |
|||
@NoArgsConstructor |
|||
public class YearInterval extends BaseAggInterval { |
|||
|
|||
@Override |
|||
public AggIntervalType getType() { |
|||
return AggIntervalType.YEAR; |
|||
} |
|||
|
|||
public YearInterval(String tz, Long offsetSec) { |
|||
super(tz, offsetSec); |
|||
} |
|||
|
|||
@Override |
|||
protected ZonedDateTime alignToIntervalStart(ZonedDateTime reference) { |
|||
return ZonedDateTime.of( |
|||
LocalDate.of(reference.getYear(), 1, 1), |
|||
LocalTime.MIDNIGHT, |
|||
reference.getZone()); |
|||
} |
|||
|
|||
@Override |
|||
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) { |
|||
return currentStart.plusYears(1); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,128 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single; |
|||
|
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.ValueSource; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Argument; |
|||
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Output; |
|||
import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunctionInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.HourInterval; |
|||
|
|||
import java.util.Map; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
|||
|
|||
public class EntityAggregationCalculatedFieldConfigurationTest { |
|||
|
|||
@Test |
|||
void typeShouldBeEntityAggregation() { |
|||
var cfg = new EntityAggregationCalculatedFieldConfiguration(); |
|||
assertThat(cfg.getType()).isEqualTo(CalculatedFieldType.ENTITY_AGGREGATION); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@ValueSource(strings = {"ATTRIBUTE", "TS_ROLLING"}) |
|||
void validateShouldThrowWhenNotTsLatestArgumentUsed(String argumentType) { |
|||
var cfg = new EntityAggregationCalculatedFieldConfiguration(); |
|||
cfg.setArguments(Map.of("k", validArgument(ArgumentType.valueOf(argumentType)))); |
|||
assertThatThrownBy(cfg::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Calculated field with type: '" + cfg.getType() + "' support only TS_LATEST arguments."); |
|||
} |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenMetricMapIsEmpty() { |
|||
var cfg = new EntityAggregationCalculatedFieldConfiguration(); |
|||
|
|||
cfg.setArguments(Map.of("k", validArgument(ArgumentType.TS_LATEST))); |
|||
cfg.setMetrics(Map.of()); |
|||
|
|||
assertThatThrownBy(cfg::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Metrics map cannot be empty."); |
|||
} |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenMetricInputIsNotAggKeyInput() { |
|||
var cfg = new EntityAggregationCalculatedFieldConfiguration(); |
|||
|
|||
cfg.setArguments(Map.of("k", validArgument(ArgumentType.TS_LATEST))); |
|||
|
|||
AggMetric metric = new AggMetric(); |
|||
metric.setInput(new AggFunctionInput()); // cannot be function
|
|||
cfg.setMetrics(Map.of("m", metric)); |
|||
|
|||
cfg.setInterval(new HourInterval("Europe/Kiev", null)); |
|||
cfg.setOutput(new Output()); |
|||
|
|||
assertThatThrownBy(cfg::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Metric key can only refer to argument."); |
|||
} |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenMetricReferencesUnknownArgument() { |
|||
var cfg = new EntityAggregationCalculatedFieldConfiguration(); |
|||
|
|||
cfg.setArguments(Map.of("k", validArgument(ArgumentType.TS_LATEST))); |
|||
|
|||
AggMetric metric = new AggMetric(); |
|||
metric.setInput(new AggKeyInput("unknown")); |
|||
cfg.setMetrics(Map.of("m", metric)); |
|||
|
|||
cfg.setInterval(new HourInterval("Europe/Kiev", null)); |
|||
cfg.setOutput(new Output()); |
|||
|
|||
assertThatThrownBy(cfg::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Metric references unknown argument: 'unknown'."); |
|||
} |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenIntervalIsNull() { |
|||
var cfg = new EntityAggregationCalculatedFieldConfiguration(); |
|||
|
|||
cfg.setArguments(Map.of("k", validArgument(ArgumentType.TS_LATEST))); |
|||
cfg.setMetrics(Map.of("m", validMetric())); |
|||
cfg.setInterval(null); |
|||
cfg.setOutput(new Output()); |
|||
|
|||
assertThatThrownBy(cfg::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Interval must be defined."); |
|||
} |
|||
|
|||
private Argument validArgument(ArgumentType type) { |
|||
Argument a = new Argument(); |
|||
a.setRefEntityKey(new ReferencedEntityKey("key", type, null)); |
|||
return a; |
|||
} |
|||
|
|||
private AggMetric validMetric() { |
|||
AggMetric metric = new AggMetric(); |
|||
metric.setInput(new AggKeyInput("k")); |
|||
return metric; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,168 @@ |
|||
/** |
|||
* Copyright © 2016-2025 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.cf.configuration.aggregation.single.interval; |
|||
|
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.Arguments; |
|||
import org.junit.jupiter.params.provider.MethodSource; |
|||
|
|||
import java.time.Duration; |
|||
import java.time.Instant; |
|||
import java.time.ZoneId; |
|||
import java.time.ZonedDateTime; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.function.Function; |
|||
import java.util.function.LongFunction; |
|||
import java.util.stream.Stream; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
|||
|
|||
public class AggIntervalTest { |
|||
|
|||
private static final String TZ = "Europe/Kiev"; |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenInvalidTimZone() { |
|||
AggInterval interval = new HourInterval("TimeZone", null); |
|||
|
|||
assertThatThrownBy(interval::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessageContaining("Invalid timezone in interval: "); |
|||
} |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenOffsetIsNegative() { |
|||
AggInterval interval = new CustomInterval(TZ, -100L, TimeUnit.HOURS.toSeconds(2)); |
|||
|
|||
assertThatThrownBy(interval::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Offset cannot be negative."); |
|||
} |
|||
|
|||
@Test |
|||
void validateShouldThrowWhenOffsetGreaterThanIntervalDuration() { |
|||
AggInterval interval = new CustomInterval(TZ, TimeUnit.HOURS.toSeconds(2), TimeUnit.HOURS.toSeconds(2)); |
|||
|
|||
assertThatThrownBy(interval::validate) |
|||
.isInstanceOf(IllegalArgumentException.class) |
|||
.hasMessage("Offset must be greater than interval duration."); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@MethodSource("intervals") |
|||
void testGetStartAndEndWithoutOffset(LongFunction<AggInterval> intervalCreator) { |
|||
AggInterval interval = intervalCreator.apply(0L); |
|||
|
|||
ZonedDateTime dateTime = ZonedDateTime.of( |
|||
// 2025.11.11 00:00:00
|
|||
2025, 11, 11, 0, 0, 0, 0, ZoneId.of(TZ) |
|||
); |
|||
long startTs = interval.getDateTimeIntervalStartTs(dateTime); |
|||
long endTs = interval.getDateTimeIntervalEndTs(dateTime); |
|||
|
|||
assertThat(endTs).isGreaterThan(startTs); |
|||
assertThat(endTs - startTs).isEqualTo(interval.getCurrentIntervalDurationMillis()); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@MethodSource("intervals") |
|||
void testApplyOffset(LongFunction<AggInterval> intervalCreator) { |
|||
long offsetSec = TimeUnit.MINUTES.toSeconds(15); |
|||
AggInterval intervalWithOffset = intervalCreator.apply(offsetSec); |
|||
AggInterval intervalNoOffset = intervalCreator.apply(0L); |
|||
|
|||
ZonedDateTime dateTime = ZonedDateTime.of( |
|||
// 2025.11.11 11:20:00 - chosen so 15m offset shifts into a new interval
|
|||
2025, 11, 11, 11, 20, 0, 0, ZoneId.of(TZ) |
|||
); |
|||
|
|||
long startWithOffsetTs = intervalWithOffset.getDateTimeIntervalStartTs(dateTime); |
|||
long startNoOffsetTs = intervalNoOffset.getDateTimeIntervalStartTs(dateTime); |
|||
|
|||
ZonedDateTime startWithOffset = Instant.ofEpochMilli(startWithOffsetTs).atZone(intervalWithOffset.getZoneId()); |
|||
ZonedDateTime startNoOffset = Instant.ofEpochMilli(startNoOffsetTs).atZone(intervalNoOffset.getZoneId()); |
|||
|
|||
long actualOffset = Duration.between(startNoOffset, startWithOffset).toSeconds(); |
|||
assertThat(actualOffset).isEqualTo(offsetSec); |
|||
} |
|||
|
|||
private static Stream<Arguments> intervals() { |
|||
return Stream.of( |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new HourInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new DayInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new WeekInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new WeekSunSatInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new MonthInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new QuarterInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new YearInterval(TZ, offset)), |
|||
Arguments.of((LongFunction<AggInterval>) offset -> new CustomInterval(TZ, offset, TimeUnit.HOURS.toSeconds(4))) |
|||
); |
|||
} |
|||
|
|||
@ParameterizedTest |
|||
@MethodSource("nextIntervalFromExactDate") |
|||
void testNextIntervalFromExactDate(LongFunction<AggInterval> intervalCreator, Function<ZonedDateTime, ZonedDateTime> expectedDateTimeFunction) { |
|||
AggInterval interval = intervalCreator.apply(0L); |
|||
|
|||
ZonedDateTime currentStart = ZonedDateTime.of( |
|||
2025, 11, 11, 0, 0, 0, 0, ZoneId.of(TZ) |
|||
); |
|||
|
|||
ZonedDateTime nextStart = interval.getNextIntervalStart(currentStart); |
|||
|
|||
assertThat(nextStart).isEqualTo(expectedDateTimeFunction.apply(currentStart)); |
|||
} |
|||
|
|||
private static Stream<Arguments> nextIntervalFromExactDate() { |
|||
return Stream.of( |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new HourInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusHours(1) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new DayInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusDays(1) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new WeekInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusWeeks(1) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new WeekSunSatInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusWeeks(1) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new MonthInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusMonths(1) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new QuarterInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusMonths(3) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new YearInterval(TZ, offset), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusYears(1) |
|||
), |
|||
Arguments.of( |
|||
(LongFunction<AggInterval>) offset -> new CustomInterval(TZ, offset, TimeUnit.HOURS.toSeconds(4)), |
|||
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusHours(4) |
|||
) |
|||
); |
|||
} |
|||
|
|||
} |
|||
File diff suppressed because it is too large
@ -0,0 +1,71 @@ |
|||
///
|
|||
/// Copyright © 2016-2025 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.
|
|||
///
|
|||
|
|||
import { ChangeDetectorRef, Component, DestroyRef, forwardRef, Renderer2, ViewContainerRef, } from '@angular/core'; |
|||
import { FormBuilder, NG_VALIDATORS, NG_VALUE_ACCESSOR, } from '@angular/forms'; |
|||
import { TbPopoverService } from '@shared/components/popover.service'; |
|||
import { EntityService } from '@core/http/entity.service'; |
|||
import { Store } from '@ngrx/store'; |
|||
import { AppState } from '@core/core.state'; |
|||
import { |
|||
CalculatedFieldArgumentsTableComponent |
|||
} from '@home/components/calculated-fields/components/calculated-field-arguments/calculated-field-arguments-table.component'; |
|||
import { ArgumentEntityType } from '@shared/models/calculated-field.models'; |
|||
|
|||
@Component({ |
|||
selector: 'tb-entity-aggregation-arguments-table', |
|||
templateUrl: './calculated-field-arguments-table.component.html', |
|||
styleUrls: [`calculated-field-arguments-table.component.scss`], |
|||
providers: [ |
|||
{ |
|||
provide: NG_VALUE_ACCESSOR, |
|||
useExisting: forwardRef(() => EntityAggregationArgumentsTableComponent), |
|||
multi: true |
|||
}, |
|||
{ |
|||
provide: NG_VALIDATORS, |
|||
useExisting: forwardRef(() => EntityAggregationArgumentsTableComponent), |
|||
multi: true |
|||
} |
|||
], |
|||
}) |
|||
export class EntityAggregationArgumentsTableComponent extends CalculatedFieldArgumentsTableComponent { |
|||
|
|||
constructor( |
|||
protected fb: FormBuilder, |
|||
protected popoverService: TbPopoverService, |
|||
protected viewContainerRef: ViewContainerRef, |
|||
protected cd: ChangeDetectorRef, |
|||
protected renderer: Renderer2, |
|||
protected entityService: EntityService, |
|||
protected destroyRef: DestroyRef, |
|||
protected store: Store<AppState> |
|||
) { |
|||
super(fb, popoverService, viewContainerRef, cd, renderer, entityService, destroyRef, store); |
|||
|
|||
this.argumentNameColumn = 'calculated-fields.argument-name'; |
|||
this.displayColumns = ['name', 'type', 'key', 'actions']; |
|||
this.panelAdditionalCtx = { |
|||
hiddenEntityTypes: true, |
|||
argumentEntityTypes: [ArgumentEntityType.Current], |
|||
hint: 'calculated-fields.entity-aggregation.argument-setting-hint', |
|||
hiddenDefaultValue: true, |
|||
hiddenEntityKeyTypes: true, |
|||
}; |
|||
|
|||
this.isScript = false; |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue