8 changed files with 778 additions and 1 deletions
@ -0,0 +1,42 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.deduplicate; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.List; |
|||
|
|||
@Data |
|||
public class DeDuplicateData { |
|||
|
|||
private EntityId entityId; |
|||
private List<DeDuplicatePack> deDuplicateStates; |
|||
private int statesCount; |
|||
|
|||
public DeDuplicateData(EntityId entityId) { |
|||
this.entityId = entityId; |
|||
this.deDuplicateStates = new ArrayList<>(); |
|||
this.statesCount = 0; |
|||
} |
|||
|
|||
public void addDeDuplicatePack(List<TbMsgDeDuplicateState> states, int indexOfFirstStateInPack, int indexOfLastStateInPack) { |
|||
deDuplicateStates.add(new DeDuplicatePack(indexOfFirstStateInPack, indexOfLastStateInPack, states)); |
|||
statesCount += states.size(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.deduplicate; |
|||
|
|||
import lombok.AllArgsConstructor; |
|||
import lombok.Data; |
|||
|
|||
import java.util.List; |
|||
|
|||
@Data |
|||
@AllArgsConstructor |
|||
public class DeDuplicatePack { |
|||
|
|||
private int indexOfFirstStateInPack; |
|||
private int indexOfLastStateInPack; |
|||
private List<TbMsgDeDuplicateState> states; |
|||
|
|||
public TbMsgDeDuplicateState getFirst() { |
|||
return states.get(indexOfFirstStateInPack); |
|||
} |
|||
|
|||
public TbMsgDeDuplicateState getLast() { |
|||
return states.get(indexOfLastStateInPack); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,22 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.deduplicate; |
|||
|
|||
public enum DeDuplicateStrategy { |
|||
|
|||
FIRST, LAST, ALL |
|||
|
|||
} |
|||
@ -0,0 +1,251 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.deduplicate; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ArrayNode; |
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.rule.engine.api.RuleNode; |
|||
import org.thingsboard.rule.engine.api.TbContext; |
|||
import org.thingsboard.rule.engine.api.TbNode; |
|||
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|||
import org.thingsboard.rule.engine.api.TbNodeException; |
|||
import org.thingsboard.rule.engine.api.TbRelationTypes; |
|||
import org.thingsboard.rule.engine.api.util.TbNodeUtils; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.plugin.ComponentType; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Collections; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.stream.Collectors; |
|||
|
|||
@RuleNode( |
|||
type = ComponentType.ACTION, |
|||
name = "de-duplicate", |
|||
configClazz = TbMsgDeDuplicateNodeConfiguration.class, |
|||
nodeDescription = "De-duplicate messages for a configurable period based on a specified de-duplication strategy.", |
|||
nodeDetails = "Rule node allows you to select one of the following strategy to de-duplicate messages: <br></br>" + |
|||
"<b>FIRST</b> - return first message that arrived during de-duplication period.<br></br>" + |
|||
"<b>LAST</b> - return last message that arrived during de-duplication period.<br></br>" + |
|||
"<b>ALL</b> - return all messages as a single JSON array message. Where each element represents object with <b>msg</b> and <b>metadata</b> inner properties.<br></br>" + |
|||
"By default rule node <b>De-duplicate messages by message originator</b>, however, there is an option to de-duplicate messages independently from the incoming message's originator. " + |
|||
"In case of the <b>De-duplicate strategy</b> set to <b>ALL</b> and <b>De-duplicate messages by message originator</b> checkbox set to <b>false</b>, you must configure the <b>Queue Name</b> and <b>Out Message Type</b> for the output array message. " + |
|||
"Also in this case the output message originator will be set to the current tenant id.<br></br>", |
|||
icon = "pause", |
|||
uiResources = {"static/rulenode/rulenode-core-config.js"}, |
|||
configDirective = "tbActionNodeMsgDeDuplicateConfig" |
|||
) |
|||
@Slf4j |
|||
public class TbMsgDeDuplicateNode implements TbNode { |
|||
|
|||
private static final String TB_MSG_DEDUPLICATION_TIMEOUT_MSG = "TbMsgDeDuplicateNodeMsg"; |
|||
|
|||
private TbMsgDeDuplicateNodeConfiguration config; |
|||
|
|||
private Map<EntityId, List<TbMsgDeDuplicateState>> deDuplicateStateMap; |
|||
private long deDuplicationInterval; |
|||
private long lastScheduledTs; |
|||
private boolean deDuplicateByOriginator; |
|||
private UUID nextTickId; |
|||
|
|||
|
|||
@Override |
|||
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { |
|||
this.config = TbNodeUtils.convert(configuration, TbMsgDeDuplicateNodeConfiguration.class); |
|||
this.deDuplicationInterval = TimeUnit.SECONDS.toMillis(config.getInterval()); |
|||
this.deDuplicateByOriginator = config.isDeDuplicateByOriginator(); |
|||
this.deDuplicateStateMap = new HashMap<>(); |
|||
scheduleTickMsg(ctx); |
|||
} |
|||
|
|||
@Override |
|||
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException { |
|||
if (TB_MSG_DEDUPLICATION_TIMEOUT_MSG.equals(msg.getType())) { |
|||
if (msg.getId().equals(nextTickId)) { |
|||
if (deDuplicateByOriginator) { |
|||
List<DeDuplicateData> deDuplicateDataList = getDeDuplicateDataList(); |
|||
if (!deDuplicateDataList.isEmpty()) { |
|||
deDuplicateDataList.forEach(deDuplicateData -> { |
|||
EntityId entityId = deDuplicateData.getEntityId(); |
|||
List<DeDuplicatePack> deDuplicateStates = deDuplicateData.getDeDuplicateStates(); |
|||
deDuplicateStates.forEach(pack -> |
|||
ctx.enqueueForTellNext(getOutputMsg(pack), TbRelationTypes.SUCCESS, |
|||
() -> { |
|||
deDuplicateStateMap.computeIfPresent(entityId, (key, states) -> { |
|||
pack.getStates().forEach(states::remove); |
|||
return states; |
|||
}); |
|||
}, |
|||
throwable -> { |
|||
log.trace("Failed to enqueue de-duplication output message due to: ", throwable); |
|||
deDuplicateStateMap.computeIfPresent(entityId, (key, states) -> { |
|||
List<TbMsgDeDuplicateState> statesToRemove = new ArrayList<>(); |
|||
states.forEach(state -> { |
|||
if (pack.getStates().contains(state)) { |
|||
int retries = state.getRetries(); |
|||
if (retries < config.getMaxRetries()) { |
|||
state.incrementRetries(); |
|||
} else { |
|||
statesToRemove.add(state); |
|||
} |
|||
} |
|||
}); |
|||
states.removeAll(statesToRemove); |
|||
return states; |
|||
}); |
|||
})); |
|||
}); |
|||
} |
|||
} else { |
|||
// todo implement
|
|||
} |
|||
scheduleTickMsg(ctx); |
|||
} else { |
|||
log.trace("[{}][{}] Received outdated tick msg: {}", ctx.getTenantId(), ctx.getSelfId(), msg.getId()); |
|||
} |
|||
} else { |
|||
List<TbMsgDeDuplicateState> deDuplicateStateList = deDuplicateStateMap.computeIfAbsent(msg.getOriginator(), k -> new ArrayList<>()); |
|||
if (deDuplicateStateList.size() < config.getMaxPendingMsgs()) { |
|||
log.trace("[{}] Adding msg: [{}][{}] to the pending msgs map ...", msg.getOriginator(), msg.getId(), msg.getMetaDataTs()); |
|||
deDuplicateStateList.add(new TbMsgDeDuplicateState(msg)); |
|||
ctx.ack(msg); |
|||
} else { |
|||
log.trace("[{}] Max limit of pending messages reached for entity!", msg.getOriginator()); |
|||
ctx.tellFailure(msg, new RuntimeException("Max limit of pending messages reached for entity: " + msg.getOriginator())); |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void destroy() { |
|||
deDuplicateStateMap.clear(); |
|||
} |
|||
|
|||
private List<DeDuplicateData> getDeDuplicateDataList() { |
|||
if (deDuplicateStateMap.isEmpty()) { |
|||
return Collections.emptyList(); |
|||
} else { |
|||
List<DeDuplicateData> result = new ArrayList<>(); |
|||
long deDuplicationTimeoutMs = System.currentTimeMillis(); |
|||
deDuplicateStateMap.forEach((entityId, tbMsgDeDuplicateStates) -> { |
|||
if (!tbMsgDeDuplicateStates.isEmpty()) { |
|||
DeDuplicateData deDuplicateData = new DeDuplicateData(entityId); |
|||
TbMsgDeDuplicateState firstDeDuplicateState = getFirstDeDuplicateState(tbMsgDeDuplicateStates, null); |
|||
TbMsgDeDuplicateState lastStateInPack = firstDeDuplicateState; |
|||
long deDuplicationStartTs = firstDeDuplicateState.getTs(); |
|||
long deDuplicationEndTs = deDuplicationStartTs + deDuplicationInterval; |
|||
boolean hasNextDeduplicationPack = deDuplicationEndTs < deDuplicationTimeoutMs; |
|||
while (hasNextDeduplicationPack) { |
|||
long finalDeDuplicationStartTs = deDuplicationStartTs; |
|||
List<TbMsgDeDuplicateState> statesPack = new ArrayList<>(); |
|||
for (TbMsgDeDuplicateState state : tbMsgDeDuplicateStates) { |
|||
if (state.getTs() >= finalDeDuplicationStartTs && state.getTs() < deDuplicationEndTs) { |
|||
statesPack.add(state); |
|||
if (state.getTs() > lastStateInPack.getTs()) { |
|||
lastStateInPack = state; |
|||
} |
|||
} |
|||
} |
|||
deDuplicateData.addDeDuplicatePack(statesPack, statesPack.indexOf(firstDeDuplicateState), statesPack.indexOf(lastStateInPack)); |
|||
if (deDuplicateData.getStatesCount() == tbMsgDeDuplicateStates.size()) { |
|||
hasNextDeduplicationPack = false; |
|||
} else { |
|||
firstDeDuplicateState = getFirstDeDuplicateState(tbMsgDeDuplicateStates, deDuplicationEndTs); |
|||
if (firstDeDuplicateState == null) { |
|||
hasNextDeduplicationPack = false; |
|||
} else { |
|||
deDuplicationStartTs = firstDeDuplicateState.getTs(); |
|||
deDuplicationEndTs = deDuplicationStartTs + deDuplicationInterval; |
|||
hasNextDeduplicationPack = deDuplicationEndTs < deDuplicationTimeoutMs; |
|||
} |
|||
} |
|||
} |
|||
result.add(deDuplicateData); |
|||
} |
|||
}); |
|||
return result; |
|||
} |
|||
} |
|||
|
|||
private TbMsgDeDuplicateState getFirstDeDuplicateState(List<TbMsgDeDuplicateState> tbMsgDeDuplicateStates, Long previousPackEndTs) { |
|||
if (previousPackEndTs != null) { |
|||
tbMsgDeDuplicateStates = tbMsgDeDuplicateStates.stream().filter(state -> state.getTs() >= previousPackEndTs).collect(Collectors.toList()); |
|||
} |
|||
TbMsgDeDuplicateState first = null; |
|||
for (TbMsgDeDuplicateState state : tbMsgDeDuplicateStates) { |
|||
if (first == null || state.getTs() < first.getTs()) { |
|||
first = state; |
|||
} |
|||
} |
|||
return first; |
|||
} |
|||
|
|||
private TbMsg getOutputMsg(DeDuplicatePack pack) { |
|||
switch (config.getStrategy()) { |
|||
case FIRST: |
|||
return pack.getFirst().getTbMsg(); |
|||
case LAST: |
|||
return pack.getLast().getTbMsg(); |
|||
default: |
|||
EntityId originator = pack.getStates().get(0).getTbMsg().getOriginator(); |
|||
String queueName = pack.getStates().get(0).getTbMsg().getQueueName(); |
|||
String outMsgType = pack.getStates().get(0).getTbMsg().getType(); |
|||
String data = getMergedData(pack.getStates().stream() |
|||
.map(TbMsgDeDuplicateState::getTbMsg) |
|||
.collect(Collectors.toList())); |
|||
return TbMsg.newMsg(queueName, outMsgType, originator, getMetadata(), data); |
|||
} |
|||
} |
|||
|
|||
private void scheduleTickMsg(TbContext ctx) { |
|||
long curTs = System.currentTimeMillis(); |
|||
if (lastScheduledTs == 0L) { |
|||
lastScheduledTs = curTs; |
|||
} |
|||
lastScheduledTs = lastScheduledTs + deDuplicationInterval; |
|||
long curDelay = Math.max(0L, (lastScheduledTs - curTs)); |
|||
TbMsg tickMsg = ctx.newMsg(null, TB_MSG_DEDUPLICATION_TIMEOUT_MSG, ctx.getSelfId(), new TbMsgMetaData(), ""); |
|||
nextTickId = tickMsg.getId(); |
|||
ctx.tellSelf(tickMsg, curDelay); |
|||
} |
|||
|
|||
private String getMergedData(List<TbMsg> msgs) { |
|||
ArrayNode mergedData = JacksonUtil.OBJECT_MAPPER.createArrayNode(); |
|||
msgs.forEach(msg -> { |
|||
ObjectNode msgNode = JacksonUtil.newObjectNode(); |
|||
msgNode.set("msg", JacksonUtil.toJsonNode(msg.getData())); |
|||
msgNode.set("metadata", JacksonUtil.valueToTree(msg.getMetaData().getData())); |
|||
mergedData.add(msgNode); |
|||
}); |
|||
return JacksonUtil.toString(mergedData); |
|||
} |
|||
|
|||
private TbMsgMetaData getMetadata() { |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("ts", String.valueOf(System.currentTimeMillis())); |
|||
return metaData; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,45 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.deduplicate; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.rule.engine.api.NodeConfiguration; |
|||
|
|||
@Data |
|||
public class TbMsgDeDuplicateNodeConfiguration implements NodeConfiguration<TbMsgDeDuplicateNodeConfiguration> { |
|||
|
|||
private int interval; |
|||
private int maxPendingMsgs; |
|||
private int maxRetries; |
|||
private boolean deDuplicateByOriginator; |
|||
private DeDuplicateStrategy strategy; |
|||
|
|||
// only for DeDuplicateStrategy.ALL when deDuplicateByOriginator set to false,
|
|||
// otherwise these parameters will be used from first msg of deDuplication msgs list.
|
|||
private String outMsgType; |
|||
private String queueName; |
|||
|
|||
@Override |
|||
public TbMsgDeDuplicateNodeConfiguration defaultConfiguration() { |
|||
TbMsgDeDuplicateNodeConfiguration configuration = new TbMsgDeDuplicateNodeConfiguration(); |
|||
configuration.setInterval(60); |
|||
configuration.setMaxPendingMsgs(100); |
|||
configuration.setMaxRetries(3); |
|||
configuration.setDeDuplicateByOriginator(true); |
|||
configuration.setStrategy(DeDuplicateStrategy.FIRST); |
|||
return configuration; |
|||
} |
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.deduplicate; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
|
|||
@Data |
|||
public class TbMsgDeDuplicateState { |
|||
|
|||
private TbMsg tbMsg; |
|||
private int retries; |
|||
|
|||
public TbMsgDeDuplicateState(TbMsg tbMsg) { |
|||
this.tbMsg = tbMsg; |
|||
this.retries = 0; |
|||
} |
|||
|
|||
public long getTs() { |
|||
return tbMsg.getMetaDataTs(); |
|||
} |
|||
|
|||
public void incrementRetries() { |
|||
this.retries++; |
|||
} |
|||
} |
|||
@ -0,0 +1,340 @@ |
|||
/** |
|||
* Copyright © 2016-2022 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.rule.engine.action; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ArrayNode; |
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.junit.jupiter.api.AfterEach; |
|||
import org.junit.jupiter.api.Assertions; |
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.mockito.ArgumentMatchers; |
|||
import org.mockito.stubbing.Answer; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|||
import org.thingsboard.rule.engine.api.TbContext; |
|||
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|||
import org.thingsboard.rule.engine.api.TbNodeException; |
|||
import org.thingsboard.rule.engine.api.TbRelationTypes; |
|||
import org.thingsboard.rule.engine.deduplicate.DeDuplicateStrategy; |
|||
import org.thingsboard.rule.engine.deduplicate.TbMsgDeDuplicateNode; |
|||
import org.thingsboard.rule.engine.deduplicate.TbMsgDeDuplicateNodeConfiguration; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.RuleNodeId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.common.msg.session.SessionMsgType; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.List; |
|||
import java.util.Random; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.CountDownLatch; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.ScheduledExecutorService; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.atomic.AtomicInteger; |
|||
import java.util.concurrent.atomic.AtomicLong; |
|||
import java.util.function.Consumer; |
|||
|
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.ArgumentMatchers.isNull; |
|||
import static org.mockito.ArgumentMatchers.nullable; |
|||
import static org.mockito.Mockito.doAnswer; |
|||
import static org.mockito.Mockito.mock; |
|||
import static org.mockito.Mockito.spy; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@Slf4j |
|||
public class TbMsgDeDuplicateNodeTest { |
|||
|
|||
private static final String MAIN_QUEUE_NAME = "Main"; |
|||
private static final String HIGH_PRIORITY_QUEUE_NAME = "HighPriority"; |
|||
private static final String TB_MSG_DEDUPLICATION_TIMEOUT_MSG = "TbMsgDeDuplicateNodeMsg"; |
|||
|
|||
private TbContext ctx; |
|||
|
|||
private final ThingsBoardThreadFactory factory = ThingsBoardThreadFactory.forName("de-duplication-node-test"); |
|||
private final ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor(factory); |
|||
private final int deDuplicationInterval = 1; |
|||
|
|||
private TenantId tenantId; |
|||
|
|||
private TbMsgDeDuplicateNode node; |
|||
private TbMsgDeDuplicateNodeConfiguration config; |
|||
private TbNodeConfiguration nodeConfiguration; |
|||
|
|||
private CountDownLatch awaitTellSelfLatch; |
|||
|
|||
@BeforeEach |
|||
public void init() throws TbNodeException { |
|||
ctx = mock(TbContext.class); |
|||
|
|||
tenantId = TenantId.fromUUID(UUID.randomUUID()); |
|||
RuleNodeId ruleNodeId = new RuleNodeId(UUID.randomUUID()); |
|||
|
|||
when(ctx.getSelfId()).thenReturn(ruleNodeId); |
|||
when(ctx.getTenantId()).thenReturn(tenantId); |
|||
|
|||
doAnswer((Answer<TbMsg>) invocationOnMock -> { |
|||
String type = (String) (invocationOnMock.getArguments())[1]; |
|||
EntityId originator = (EntityId) (invocationOnMock.getArguments())[2]; |
|||
TbMsgMetaData metaData = (TbMsgMetaData) (invocationOnMock.getArguments())[3]; |
|||
String data = (String) (invocationOnMock.getArguments())[4]; |
|||
return TbMsg.newMsg(type, originator, metaData.copy(), data); |
|||
}).when(ctx).newMsg(isNull(), eq(TB_MSG_DEDUPLICATION_TIMEOUT_MSG), nullable(EntityId.class), any(TbMsgMetaData.class), any(String.class)); |
|||
node = spy(new TbMsgDeDuplicateNode()); |
|||
config = new TbMsgDeDuplicateNodeConfiguration().defaultConfiguration(); |
|||
} |
|||
|
|||
private void invokeTellSelf(int maxNumberOfInvocation) { |
|||
invokeTellSelf(maxNumberOfInvocation, false, 0); |
|||
} |
|||
|
|||
private void invokeTellSelf(int maxNumberOfInvocation, boolean delayScheduleTimeout, int delayMultiplier) { |
|||
AtomicLong scheduleTimeout = new AtomicLong(deDuplicationInterval); |
|||
AtomicInteger scheduleCount = new AtomicInteger(0); |
|||
doAnswer((Answer<Void>) invocationOnMock -> { |
|||
scheduleCount.getAndIncrement(); |
|||
if (scheduleCount.get() <= maxNumberOfInvocation) { |
|||
TbMsg msg = (TbMsg) (invocationOnMock.getArguments())[0]; |
|||
executorService.schedule(() -> { |
|||
try { |
|||
node.onMsg(ctx, msg); |
|||
awaitTellSelfLatch.countDown(); |
|||
} catch (ExecutionException | InterruptedException | TbNodeException e) { |
|||
log.error("Failed to execute tellSelf method call due to: ", e); |
|||
} |
|||
}, scheduleTimeout.get(), TimeUnit.SECONDS); |
|||
if (delayScheduleTimeout) { |
|||
scheduleTimeout.set(scheduleTimeout.get() * delayMultiplier); |
|||
} |
|||
} |
|||
|
|||
return null; |
|||
}).when(ctx).tellSelf(ArgumentMatchers.any(TbMsg.class), ArgumentMatchers.anyLong()); |
|||
} |
|||
|
|||
@AfterEach |
|||
public void destroy() { |
|||
executorService.shutdown(); |
|||
node.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void given_100_messages_then_verifyOutputFirst() throws TbNodeException, ExecutionException, InterruptedException { |
|||
int wantedNumberOfTellSelfInvocation = 2; |
|||
int msgCount = 100; |
|||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|||
invokeTellSelf(wantedNumberOfTellSelfInvocation); |
|||
|
|||
config.setInterval(deDuplicationInterval); |
|||
config.setMaxPendingMsgs(msgCount); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctx, nodeConfiguration); |
|||
|
|||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|||
long currentTimeMillis = System.currentTimeMillis(); |
|||
|
|||
List<TbMsg> inputMsgs = getTbMsgs(deviceId, msgCount, currentTimeMillis, 500); |
|||
for (TbMsg msg : inputMsgs) { |
|||
node.onMsg(ctx, msg); |
|||
} |
|||
|
|||
TbMsg msgToReject = createMsg(deviceId, inputMsgs.get(inputMsgs.size() - 1).getMetaDataTs() + 2); |
|||
node.onMsg(ctx, msgToReject); |
|||
|
|||
awaitTellSelfLatch.await(); |
|||
|
|||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|||
|
|||
verify(ctx, times(msgCount)).ack(any()); |
|||
verify(ctx, times(1)).tellFailure(eq(msgToReject), any()); |
|||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation + 1)).onMsg(eq(ctx), any()); |
|||
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|||
Assertions.assertEquals(inputMsgs.get(0), newMsgCaptor.getValue()); |
|||
} |
|||
|
|||
@Test |
|||
public void given_100_messages_then_verifyOutputLast() throws TbNodeException, ExecutionException, InterruptedException { |
|||
int wantedNumberOfTellSelfInvocation = 2; |
|||
int msgCount = 100; |
|||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|||
invokeTellSelf(wantedNumberOfTellSelfInvocation); |
|||
|
|||
config.setStrategy(DeDuplicateStrategy.LAST); |
|||
config.setInterval(deDuplicationInterval); |
|||
config.setMaxPendingMsgs(msgCount); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctx, nodeConfiguration); |
|||
|
|||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|||
long currentTimeMillis = System.currentTimeMillis(); |
|||
|
|||
List<TbMsg> inputMsgs = getTbMsgs(deviceId, msgCount, currentTimeMillis, 500); |
|||
|
|||
int indexOfLastMsgInArray = inputMsgs.size() - 1; |
|||
int indexToSetMaxTs = new Random().nextInt(indexOfLastMsgInArray) + 1; |
|||
TbMsg currentMaxTsMsg = inputMsgs.get(indexOfLastMsgInArray); |
|||
TbMsg newLastMsgOfArray = inputMsgs.get(indexToSetMaxTs); |
|||
inputMsgs.set(indexOfLastMsgInArray, newLastMsgOfArray); |
|||
inputMsgs.set(indexToSetMaxTs, currentMaxTsMsg); |
|||
|
|||
for (TbMsg msg : inputMsgs) { |
|||
node.onMsg(ctx, msg); |
|||
} |
|||
|
|||
TbMsg msgToReject = createMsg(deviceId, inputMsgs.get(indexOfLastMsgInArray).getMetaDataTs() + 2); |
|||
node.onMsg(ctx, msgToReject); |
|||
|
|||
awaitTellSelfLatch.await(); |
|||
|
|||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|||
|
|||
verify(ctx, times(msgCount)).ack(any()); |
|||
verify(ctx, times(1)).tellFailure(eq(msgToReject), any()); |
|||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation + 1)).onMsg(eq(ctx), any()); |
|||
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|||
Assertions.assertEquals(currentMaxTsMsg, newMsgCaptor.getValue()); |
|||
} |
|||
|
|||
@Test |
|||
public void given_100_messages_then_verifyOutputAll() throws TbNodeException, ExecutionException, InterruptedException { |
|||
int wantedNumberOfTellSelfInvocation = 2; |
|||
int msgCount = 100; |
|||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|||
invokeTellSelf(wantedNumberOfTellSelfInvocation); |
|||
|
|||
config.setInterval(deDuplicationInterval); |
|||
config.setMaxPendingMsgs(msgCount); |
|||
config.setStrategy(DeDuplicateStrategy.ALL); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctx, nodeConfiguration); |
|||
|
|||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|||
long currentTimeMillis = System.currentTimeMillis(); |
|||
|
|||
List<TbMsg> inputMsgs = getTbMsgs(deviceId, msgCount, currentTimeMillis, 500); |
|||
for (TbMsg msg : inputMsgs) { |
|||
node.onMsg(ctx, msg); |
|||
} |
|||
|
|||
awaitTellSelfLatch.await(); |
|||
|
|||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|||
|
|||
verify(ctx, times(msgCount)).ack(any()); |
|||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any()); |
|||
verify(ctx, times(1)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|||
|
|||
Assertions.assertEquals(1, newMsgCaptor.getAllValues().size()); |
|||
TbMsg outMessage = newMsgCaptor.getAllValues().get(0); |
|||
Assertions.assertEquals(getMergedData(inputMsgs), outMessage.getData()); |
|||
Assertions.assertEquals(deviceId, outMessage.getOriginator()); |
|||
} |
|||
|
|||
@Test |
|||
public void given_100_messages_then_verifyOutput_2_packs() throws TbNodeException, ExecutionException, InterruptedException { |
|||
int wantedNumberOfTellSelfInvocation = 2; |
|||
int msgCount = 100; |
|||
awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); |
|||
invokeTellSelf(wantedNumberOfTellSelfInvocation, true, 3); |
|||
|
|||
config.setInterval(deDuplicationInterval); |
|||
config.setMaxPendingMsgs(msgCount); |
|||
config.setStrategy(DeDuplicateStrategy.ALL); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctx, nodeConfiguration); |
|||
|
|||
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
|||
long currentTimeMillis = System.currentTimeMillis(); |
|||
|
|||
List<TbMsg> firstMsgPack = getTbMsgs(deviceId, msgCount / 2, currentTimeMillis, 500); |
|||
for (TbMsg msg : firstMsgPack) { |
|||
node.onMsg(ctx, msg); |
|||
} |
|||
TbMsg firstMsgFromFirstPack = firstMsgPack.get(0); |
|||
long firstPackDeDuplicationPackEndTs = firstMsgFromFirstPack.getMetaDataTs() + TimeUnit.SECONDS.toMillis(deDuplicationInterval); |
|||
|
|||
List<TbMsg> secondMsgPack = getTbMsgs(deviceId, msgCount / 2, firstPackDeDuplicationPackEndTs, 500); |
|||
for (TbMsg msg : secondMsgPack) { |
|||
node.onMsg(ctx, msg); |
|||
} |
|||
|
|||
awaitTellSelfLatch.await(); |
|||
|
|||
ArgumentCaptor<TbMsg> newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
ArgumentCaptor<Runnable> successCaptor = ArgumentCaptor.forClass(Runnable.class); |
|||
ArgumentCaptor<Consumer<Throwable>> failureCaptor = ArgumentCaptor.forClass(Consumer.class); |
|||
|
|||
verify(ctx, times(msgCount)).ack(any()); |
|||
verify(node, times(msgCount + wantedNumberOfTellSelfInvocation)).onMsg(eq(ctx), any()); |
|||
verify(ctx, times(2)).enqueueForTellNext(newMsgCaptor.capture(), eq(TbRelationTypes.SUCCESS), successCaptor.capture(), failureCaptor.capture()); |
|||
|
|||
Assertions.assertEquals(2, newMsgCaptor.getAllValues().size()); |
|||
Assertions.assertEquals(getMergedData(firstMsgPack), newMsgCaptor.getAllValues().get(0).getData()); |
|||
Assertions.assertEquals(getMergedData(secondMsgPack), newMsgCaptor.getAllValues().get(1).getData()); |
|||
} |
|||
|
|||
private List<TbMsg> getTbMsgs(DeviceId deviceId, int msgCount, long currentTimeMillis, int initTsStep) { |
|||
List<TbMsg> inputMsgs = new ArrayList<>(); |
|||
var ts = currentTimeMillis + initTsStep; |
|||
for (int i = 0; i < msgCount; i++) { |
|||
inputMsgs.add(createMsg(deviceId, ts)); |
|||
ts += 2; |
|||
} |
|||
return inputMsgs; |
|||
} |
|||
|
|||
private TbMsg createMsg(DeviceId deviceId, long ts) { |
|||
ObjectNode dataNode = JacksonUtil.newObjectNode(); |
|||
dataNode.put("deviceId", deviceId.getId().toString()); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("ts", String.valueOf(ts)); |
|||
return TbMsg.newMsg( |
|||
MAIN_QUEUE_NAME, |
|||
SessionMsgType.POST_TELEMETRY_REQUEST.name(), |
|||
deviceId, |
|||
metaData, |
|||
JacksonUtil.toString(dataNode)); |
|||
} |
|||
|
|||
private String getMergedData(List<TbMsg> msgs) { |
|||
ArrayNode mergedData = JacksonUtil.OBJECT_MAPPER.createArrayNode(); |
|||
msgs.forEach(msg -> { |
|||
ObjectNode msgNode = JacksonUtil.newObjectNode(); |
|||
msgNode.set("msg", JacksonUtil.toJsonNode(msg.getData())); |
|||
msgNode.set("metadata", JacksonUtil.valueToTree(msg.getMetaData().getData())); |
|||
mergedData.add(msgNode); |
|||
}); |
|||
return JacksonUtil.toString(mergedData); |
|||
} |
|||
|
|||
} |
|||
Loading…
Reference in new issue