122 changed files with 3642 additions and 669 deletions
@ -0,0 +1,92 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.common.data.kv; |
||||
|
|
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.EqualsAndHashCode; |
||||
|
import lombok.Getter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
|
||||
|
import java.time.DateTimeException; |
||||
|
import java.time.ZoneId; |
||||
|
import java.util.Map; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
@AllArgsConstructor |
||||
|
@EqualsAndHashCode |
||||
|
@Slf4j |
||||
|
public class AggregationParams { |
||||
|
private static final Map<String, String> TZ_LINKS = Map.of("EST", "America/New_York", "GMT+0", "GMT", "GMT-0", "GMT", "HST", "US/Hawaii", "MST", "America/Phoenix", "ROC", "Asia/Taipei"); |
||||
|
@Getter |
||||
|
private final Aggregation aggregation; |
||||
|
@Getter |
||||
|
private final IntervalType intervalType; |
||||
|
@Getter |
||||
|
private final ZoneId tzId; |
||||
|
|
||||
|
private final long interval; |
||||
|
|
||||
|
public static AggregationParams none() { |
||||
|
return new AggregationParams(Aggregation.NONE, null, null, 0L); |
||||
|
} |
||||
|
|
||||
|
public static AggregationParams milliseconds(Aggregation aggregationType, long aggregationIntervalMs) { |
||||
|
return new AggregationParams(aggregationType, IntervalType.MILLISECONDS, null, aggregationIntervalMs); |
||||
|
} |
||||
|
|
||||
|
public static AggregationParams calendar(Aggregation aggregationType, IntervalType intervalType, String tzIdStr) { |
||||
|
return calendar(aggregationType, intervalType, getZoneId(tzIdStr)); |
||||
|
} |
||||
|
|
||||
|
public static AggregationParams calendar(Aggregation aggregationType, IntervalType intervalType, ZoneId tzId) { |
||||
|
return new AggregationParams(aggregationType, intervalType, tzId, 0L); |
||||
|
} |
||||
|
|
||||
|
public static AggregationParams of(Aggregation aggregation, IntervalType intervalType, ZoneId tzId, long interval) { |
||||
|
return new AggregationParams(aggregation, intervalType, tzId, interval); |
||||
|
} |
||||
|
|
||||
|
public long getInterval() { |
||||
|
if (intervalType == null) { |
||||
|
return 0L; |
||||
|
} else { |
||||
|
switch (intervalType) { |
||||
|
case WEEK: |
||||
|
case WEEK_ISO: |
||||
|
return TimeUnit.DAYS.toMillis(7); |
||||
|
case MONTH: |
||||
|
return TimeUnit.DAYS.toMillis(30); |
||||
|
case QUARTER: |
||||
|
return TimeUnit.DAYS.toMillis(90); |
||||
|
default: |
||||
|
return interval; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private static ZoneId getZoneId(String tzIdStr) { |
||||
|
if (StringUtils.isEmpty(tzIdStr)) { |
||||
|
return ZoneId.systemDefault(); |
||||
|
} |
||||
|
try { |
||||
|
return ZoneId.of(tzIdStr, TZ_LINKS); |
||||
|
} catch (DateTimeException e) { |
||||
|
log.warn("[{}] Failed to convert the time zone. Fallback to default.", tzIdStr); |
||||
|
return ZoneId.systemDefault(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,22 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.common.data.kv; |
||||
|
|
||||
|
public enum IntervalType { |
||||
|
|
||||
|
MILLISECONDS, WEEK/*Sunday-Saturday*/, WEEK_ISO/*Monday-Sunday*/, MONTH, QUARTER |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,187 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy; |
||||
|
import org.thingsboard.server.queue.scheduler.SchedulerComponent; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.Random; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.ConcurrentMap; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.concurrent.atomic.AtomicBoolean; |
||||
|
import java.util.concurrent.atomic.AtomicLong; |
||||
|
import java.util.concurrent.atomic.AtomicReference; |
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class AbstractActivityManager<Key, Metadata> implements ActivityManager<Key> { |
||||
|
|
||||
|
private final ConcurrentMap<Key, ActivityStateWrapper> states = new ConcurrentHashMap<>(); |
||||
|
|
||||
|
@Autowired |
||||
|
protected SchedulerComponent scheduler; |
||||
|
|
||||
|
@Data |
||||
|
private class ActivityStateWrapper { |
||||
|
|
||||
|
private volatile ActivityState<Metadata> state; |
||||
|
private volatile long lastReportedTime; |
||||
|
private volatile ActivityStrategy strategy; |
||||
|
|
||||
|
} |
||||
|
|
||||
|
protected void init() { |
||||
|
var reportingPeriodMillis = getReportingPeriodMillis(); |
||||
|
scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodMillis), reportingPeriodMillis, TimeUnit.MILLISECONDS); |
||||
|
} |
||||
|
|
||||
|
protected abstract long getReportingPeriodMillis(); |
||||
|
|
||||
|
protected abstract ActivityState<Metadata> createNewState(Key key); |
||||
|
|
||||
|
protected abstract ActivityStrategy getStrategy(); |
||||
|
|
||||
|
protected abstract ActivityState<Metadata> updateState(Key key, ActivityState<Metadata> state); |
||||
|
|
||||
|
protected abstract boolean hasExpired(long lastRecordedTime); |
||||
|
|
||||
|
protected abstract void onStateExpiry(Key key, Metadata metadata); |
||||
|
|
||||
|
protected abstract void reportActivity(Key key, Metadata metadata, long timeToReport, ActivityReportCallback<Key> callback); |
||||
|
|
||||
|
@Override |
||||
|
public void onActivity(Key key, long newLastRecordedTime) { |
||||
|
if (key == null) { |
||||
|
log.error("Failed to process activity event: provided activity key is null."); |
||||
|
return; |
||||
|
} |
||||
|
log.debug("Received activity event for key: [{}]", key); |
||||
|
|
||||
|
var shouldReport = new AtomicBoolean(false); |
||||
|
var lastRecordedTime = new AtomicLong(); |
||||
|
var lastReportedTime = new AtomicLong(); |
||||
|
var metadata = new AtomicReference<Metadata>(); |
||||
|
|
||||
|
var activityStateWrapper = states.compute(key, (__, stateWrapper) -> { |
||||
|
if (stateWrapper == null) { |
||||
|
var newState = createNewState(key); |
||||
|
if (newState == null) { |
||||
|
return null; |
||||
|
} |
||||
|
stateWrapper = new ActivityStateWrapper(); |
||||
|
stateWrapper.setState(newState); |
||||
|
stateWrapper.setStrategy(getStrategy()); |
||||
|
} |
||||
|
var state = stateWrapper.getState(); |
||||
|
if (state.getLastRecordedTime() < newLastRecordedTime) { |
||||
|
state.setLastRecordedTime(newLastRecordedTime); |
||||
|
} |
||||
|
shouldReport.set(stateWrapper.getStrategy().onActivity()); |
||||
|
lastRecordedTime.set(state.getLastRecordedTime()); |
||||
|
lastReportedTime.set(stateWrapper.getLastReportedTime()); |
||||
|
metadata.set(state.getMetadata()); |
||||
|
return stateWrapper; |
||||
|
}); |
||||
|
|
||||
|
if (activityStateWrapper == null) { |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
if (shouldReport.get() && lastReportedTime.get() < lastRecordedTime.get()) { |
||||
|
log.debug("Going to report first activity event for key: [{}].", key); |
||||
|
reportActivity(key, metadata.get(), lastRecordedTime.get(), new ActivityReportCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(Key key, long reportedTime) { |
||||
|
updateLastReportedTime(key, reportedTime); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Key key, Throwable t) { |
||||
|
log.debug("Failed to report first activity event for key: [{}].", key, t); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onReportingPeriodEnd() { |
||||
|
log.debug("Going to end reporting period."); |
||||
|
for (Map.Entry<Key, ActivityStateWrapper> entry : states.entrySet()) { |
||||
|
var key = entry.getKey(); |
||||
|
var stateWrapper = entry.getValue(); |
||||
|
var currentState = stateWrapper.getState(); |
||||
|
|
||||
|
long lastRecordedTime = currentState.getLastRecordedTime(); |
||||
|
long lastReportedTime = stateWrapper.getLastReportedTime(); |
||||
|
var metadata = currentState.getMetadata(); |
||||
|
|
||||
|
boolean hasExpired; |
||||
|
boolean shouldReport; |
||||
|
|
||||
|
var updatedState = updateState(key, currentState); |
||||
|
if (updatedState != null) { |
||||
|
stateWrapper.setState(updatedState); |
||||
|
lastRecordedTime = updatedState.getLastRecordedTime(); |
||||
|
metadata = updatedState.getMetadata(); |
||||
|
hasExpired = hasExpired(lastRecordedTime); |
||||
|
shouldReport = stateWrapper.getStrategy().onReportingPeriodEnd(); |
||||
|
} else { |
||||
|
states.remove(key); |
||||
|
hasExpired = false; |
||||
|
shouldReport = true; |
||||
|
} |
||||
|
|
||||
|
if (hasExpired) { |
||||
|
states.remove(key); |
||||
|
onStateExpiry(key, metadata); |
||||
|
shouldReport = true; |
||||
|
} |
||||
|
|
||||
|
if (shouldReport && lastReportedTime < lastRecordedTime) { |
||||
|
log.debug("Going to report last activity event for key: [{}].", key); |
||||
|
reportActivity(key, metadata, lastRecordedTime, new ActivityReportCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(Key key, long reportedTime) { |
||||
|
updateLastReportedTime(key, reportedTime); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(Key key, Throwable t) { |
||||
|
log.debug("Failed to report last activity event for key: [{}].", key, t); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public long getLastRecordedTime(Key key) { |
||||
|
ActivityStateWrapper stateWrapper = states.get(key); |
||||
|
return stateWrapper == null ? 0L : stateWrapper.getState().getLastRecordedTime(); |
||||
|
} |
||||
|
|
||||
|
private void updateLastReportedTime(Key key, long newLastReportedTime) { |
||||
|
states.computeIfPresent(key, (__, stateWrapper) -> { |
||||
|
stateWrapper.setLastReportedTime(Math.max(stateWrapper.getLastReportedTime(), newLastReportedTime)); |
||||
|
return stateWrapper; |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,26 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity; |
||||
|
|
||||
|
public interface ActivityManager<Key> { |
||||
|
|
||||
|
void onActivity(Key key, long activityTimeMillis); |
||||
|
|
||||
|
void onReportingPeriodEnd(); |
||||
|
|
||||
|
long getLastRecordedTime(Key key); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,24 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity; |
||||
|
|
||||
|
public interface ActivityReportCallback<Key> { |
||||
|
|
||||
|
void onSuccess(Key key, long reportedTime); |
||||
|
|
||||
|
void onFailure(Key key, Throwable t); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,24 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
public interface ActivityStrategy { |
||||
|
|
||||
|
boolean onActivity(); |
||||
|
|
||||
|
boolean onReportingPeriodEnd(); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,47 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
public enum ActivityStrategyType { |
||||
|
|
||||
|
ALL { |
||||
|
@Override |
||||
|
public ActivityStrategy toStrategy() { |
||||
|
return AllEventsActivityStrategy.getInstance(); |
||||
|
} |
||||
|
}, |
||||
|
FIRST { |
||||
|
@Override |
||||
|
public ActivityStrategy toStrategy() { |
||||
|
return new FirstEventActivityStrategy(); |
||||
|
} |
||||
|
}, |
||||
|
LAST { |
||||
|
@Override |
||||
|
public ActivityStrategy toStrategy() { |
||||
|
return LastEventActivityStrategy.getInstance(); |
||||
|
} |
||||
|
}, |
||||
|
FIRST_AND_LAST { |
||||
|
@Override |
||||
|
public ActivityStrategy toStrategy() { |
||||
|
return new FirstAndLastEventActivityStrategy(); |
||||
|
} |
||||
|
}; |
||||
|
|
||||
|
public abstract ActivityStrategy toStrategy(); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import lombok.EqualsAndHashCode; |
||||
|
|
||||
|
@EqualsAndHashCode |
||||
|
public final class FirstAndLastEventActivityStrategy implements ActivityStrategy { |
||||
|
|
||||
|
private boolean firstEventReceived; |
||||
|
|
||||
|
@Override |
||||
|
public synchronized boolean onActivity() { |
||||
|
if (!firstEventReceived) { |
||||
|
firstEventReceived = true; |
||||
|
return true; |
||||
|
} |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public synchronized boolean onReportingPeriodEnd() { |
||||
|
firstEventReceived = false; |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import lombok.EqualsAndHashCode; |
||||
|
|
||||
|
@EqualsAndHashCode |
||||
|
public final class FirstEventActivityStrategy implements ActivityStrategy { |
||||
|
|
||||
|
private boolean firstEventReceived; |
||||
|
|
||||
|
@Override |
||||
|
public synchronized boolean onActivity() { |
||||
|
if (!firstEventReceived) { |
||||
|
firstEventReceived = true; |
||||
|
return true; |
||||
|
} |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public synchronized boolean onReportingPeriodEnd() { |
||||
|
firstEventReceived = false; |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,39 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
public final class LastEventActivityStrategy implements ActivityStrategy { |
||||
|
|
||||
|
private static final LastEventActivityStrategy INSTANCE = new LastEventActivityStrategy(); |
||||
|
|
||||
|
private LastEventActivityStrategy() { |
||||
|
} |
||||
|
|
||||
|
public static LastEventActivityStrategy getInstance() { |
||||
|
return INSTANCE; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean onActivity() { |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean onReportingPeriodEnd() { |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,149 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.service; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.thingsboard.server.common.transport.TransportService; |
||||
|
import org.thingsboard.server.common.transport.TransportServiceCallback; |
||||
|
import org.thingsboard.server.common.transport.activity.AbstractActivityManager; |
||||
|
import org.thingsboard.server.common.transport.activity.ActivityReportCallback; |
||||
|
import org.thingsboard.server.common.transport.activity.ActivityState; |
||||
|
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy; |
||||
|
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategyType; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.ConcurrentMap; |
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class TransportActivityManager extends AbstractActivityManager<UUID, TransportProtos.SessionInfoProto> implements TransportService { |
||||
|
|
||||
|
public static final String SESSION_EXPIRED_MESSAGE = "Session has expired due to last activity time!"; |
||||
|
|
||||
|
public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_CLOSED = TransportProtos.SessionEventMsg.newBuilder() |
||||
|
.setSessionType(TransportProtos.SessionType.ASYNC) |
||||
|
.setEvent(TransportProtos.SessionEvent.CLOSED).build(); |
||||
|
public static final TransportProtos.SessionCloseNotificationProto SESSION_EXPIRED_NOTIFICATION_PROTO = TransportProtos.SessionCloseNotificationProto.newBuilder() |
||||
|
.setMessage(SESSION_EXPIRED_MESSAGE).build(); |
||||
|
|
||||
|
public final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>(); |
||||
|
|
||||
|
@Value("${transport.sessions.report_timeout}") |
||||
|
protected long sessionReportTimeout; |
||||
|
|
||||
|
@Value("${transport.sessions.inactivity_timeout}") |
||||
|
protected long sessionInactivityTimeout; |
||||
|
|
||||
|
@Value("${transport.activity.reporting_strategy:LAST}") |
||||
|
private ActivityStrategyType reportingStrategyType; |
||||
|
|
||||
|
@Override |
||||
|
protected long getReportingPeriodMillis() { |
||||
|
return sessionReportTimeout; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ActivityState<TransportProtos.SessionInfoProto> createNewState(UUID sessionId) { |
||||
|
SessionMetaData session = sessions.get(sessionId); |
||||
|
if (session == null) { |
||||
|
return null; |
||||
|
} |
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setMetadata(session.getSessionInfo()); |
||||
|
return state; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ActivityStrategy getStrategy() { |
||||
|
return reportingStrategyType.toStrategy(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ActivityState<TransportProtos.SessionInfoProto> updateState(UUID sessionId, ActivityState<TransportProtos.SessionInfoProto> state) { |
||||
|
SessionMetaData session = sessions.get(sessionId); |
||||
|
if (session == null) { |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
state.setMetadata(session.getSessionInfo()); |
||||
|
var sessionInfo = state.getMetadata(); |
||||
|
|
||||
|
if (sessionInfo.getGwSessionIdMSB() == 0L || sessionInfo.getGwSessionIdLSB() == 0L) { |
||||
|
return state; |
||||
|
} |
||||
|
|
||||
|
var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB()); |
||||
|
SessionMetaData gwSession = sessions.get(gwSessionId); |
||||
|
if (gwSession == null || !gwSession.isOverwriteActivityTime()) { |
||||
|
return state; |
||||
|
} |
||||
|
|
||||
|
long lastRecordedTime = state.getLastRecordedTime(); |
||||
|
long gwLastRecordedTime = getLastRecordedTime(gwSessionId); |
||||
|
log.debug("Session with id: [{}] has gateway session with id: [{}] with overwrite activity time enabled. " + |
||||
|
"Updating last activity time. Session last recorded time: [{}], gateway session last recorded time: [{}].", |
||||
|
sessionId, gwSessionId, lastRecordedTime, gwLastRecordedTime); |
||||
|
state.setLastRecordedTime(Math.max(lastRecordedTime, gwLastRecordedTime)); |
||||
|
return state; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected boolean hasExpired(long lastRecordedTime) { |
||||
|
return (getCurrentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void onStateExpiry(UUID sessionId, TransportProtos.SessionInfoProto sessionInfo) { |
||||
|
log.debug("Session with id: [{}] has expired due to last activity time.", sessionId); |
||||
|
SessionMetaData expiredSession = sessions.remove(sessionId); |
||||
|
if (expiredSession != null) { |
||||
|
deregisterSession(sessionInfo); |
||||
|
process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); |
||||
|
expiredSession.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void reportActivity(UUID sessionId, TransportProtos.SessionInfoProto currentSessionInfo, long timeToReport, ActivityReportCallback<UUID> callback) { |
||||
|
log.debug("Reporting activity state for session with id: [{}]. Time to report: [{}].", sessionId, timeToReport); |
||||
|
SessionMetaData session = sessions.get(sessionId); |
||||
|
TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() |
||||
|
.setAttributeSubscription(session != null && session.isSubscribedToAttributes()) |
||||
|
.setRpcSubscription(session != null && session.isSubscribedToRPC()) |
||||
|
.setLastActivityTime(timeToReport) |
||||
|
.build(); |
||||
|
TransportProtos.SessionInfoProto sessionInfo = session != null ? session.getSessionInfo() : currentSessionInfo; |
||||
|
process(sessionInfo, subscriptionInfo, new TransportServiceCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(Void msgAcknowledged) { |
||||
|
callback.onSuccess(sessionId, timeToReport); |
||||
|
|
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onError(Throwable e) { |
||||
|
callback.onFailure(sessionId, e); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
protected long getCurrentTimeMillis() { |
||||
|
return System.currentTimeMillis(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,44 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import org.junit.jupiter.api.Test; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
|
||||
|
public class ActivityStrategyTypeTest { |
||||
|
|
||||
|
@Test |
||||
|
public void testCreateAllEventsStrategy() { |
||||
|
assertThat(ActivityStrategyType.ALL.toStrategy()).isEqualTo(AllEventsActivityStrategy.getInstance()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testCreateFirstEventStrategy() { |
||||
|
assertThat(ActivityStrategyType.FIRST.toStrategy()).isEqualTo(new FirstEventActivityStrategy()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testCreateLastEventStrategy() { |
||||
|
assertThat(ActivityStrategyType.LAST.toStrategy()).isEqualTo(LastEventActivityStrategy.getInstance()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testCreateFirstAndLastEventStrategy() { |
||||
|
assertThat(ActivityStrategyType.FIRST_AND_LAST.toStrategy()).isEqualTo(new FirstAndLastEventActivityStrategy()); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,42 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
|
||||
|
import static org.junit.jupiter.api.Assertions.assertTrue; |
||||
|
|
||||
|
public class AllEventsActivityStrategyTest { |
||||
|
|
||||
|
private AllEventsActivityStrategy strategy; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
strategy = AllEventsActivityStrategy.getInstance(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnActivity() { |
||||
|
assertTrue(strategy.onActivity(), "onActivity() should always return true."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnReportingPeriodEnd() { |
||||
|
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,53 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
|
||||
|
import static org.junit.jupiter.api.Assertions.assertFalse; |
||||
|
import static org.junit.jupiter.api.Assertions.assertTrue; |
||||
|
|
||||
|
public class FirstAndLastEventActivityStrategyTest { |
||||
|
|
||||
|
private FirstAndLastEventActivityStrategy strategy; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
strategy = new FirstAndLastEventActivityStrategy(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnActivity_FirstCall() { |
||||
|
assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnActivity_SubsequentCalls() { |
||||
|
assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); |
||||
|
assertFalse(strategy.onActivity(), "Subsequent calls of onActivity() should return false."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnReportingPeriodEnd() { |
||||
|
assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); |
||||
|
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); |
||||
|
assertTrue(strategy.onActivity(), "onActivity() should return true after onReportingPeriodEnd() for the next reporting period"); |
||||
|
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); |
||||
|
} |
||||
|
|
||||
|
|
||||
|
} |
||||
@ -0,0 +1,52 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
|
||||
|
import static org.junit.jupiter.api.Assertions.assertFalse; |
||||
|
import static org.junit.jupiter.api.Assertions.assertTrue; |
||||
|
|
||||
|
public class FirstEventActivityStrategyTest { |
||||
|
|
||||
|
private FirstEventActivityStrategy strategy; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
strategy = new FirstEventActivityStrategy(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnActivity_FirstCall() { |
||||
|
assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnActivity_SubsequentCalls() { |
||||
|
assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); |
||||
|
assertFalse(strategy.onActivity(), "Subsequent calls of onActivity() should return false."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnReportingPeriodEnd() { |
||||
|
assertTrue(strategy.onActivity(), "First call of onActivity() should return true."); |
||||
|
assertFalse(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return false."); |
||||
|
assertTrue(strategy.onActivity(), "onActivity() should return true after onReportingPeriodEnd()."); |
||||
|
assertFalse(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return false."); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,43 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.activity.strategy; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
|
||||
|
import static org.junit.jupiter.api.Assertions.assertFalse; |
||||
|
import static org.junit.jupiter.api.Assertions.assertTrue; |
||||
|
|
||||
|
public class LastEventActivityStrategyTest { |
||||
|
|
||||
|
private LastEventActivityStrategy strategy; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
strategy = LastEventActivityStrategy.getInstance(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnActivity() { |
||||
|
assertFalse(strategy.onActivity(), "onActivity() should always return false."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testOnReportingPeriodEnd() { |
||||
|
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true."); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,508 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.transport.service; |
||||
|
|
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.junit.jupiter.params.ParameterizedTest; |
||||
|
import org.junit.jupiter.params.provider.Arguments; |
||||
|
import org.junit.jupiter.params.provider.EnumSource; |
||||
|
import org.junit.jupiter.params.provider.MethodSource; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.common.transport.SessionMsgListener; |
||||
|
import org.thingsboard.server.common.transport.TransportServiceCallback; |
||||
|
import org.thingsboard.server.common.transport.activity.ActivityReportCallback; |
||||
|
import org.thingsboard.server.common.transport.activity.ActivityState; |
||||
|
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy; |
||||
|
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategyType; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.ConcurrentMap; |
||||
|
import java.util.stream.Stream; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.Mockito.doCallRealMethod; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.never; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_CLOSED; |
||||
|
import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EXPIRED_NOTIFICATION_PROTO; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class TransportActivityManagerTest { |
||||
|
|
||||
|
private final UUID SESSION_ID = UUID.fromString("1306648a-9b26-11ee-b9d1-0242ac120002"); |
||||
|
|
||||
|
@Mock |
||||
|
private DefaultTransportService transportServiceMock; |
||||
|
private ConcurrentMap<UUID, SessionMetaData> sessions; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setup() { |
||||
|
sessions = new ConcurrentHashMap<>(); |
||||
|
ReflectionTestUtils.setField(transportServiceMock, "sessions", sessions); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenKeyAndTimeToReportAndSessionExists_whenReportingActivity_thenShouldReportActivityWithSubscriptionsAndSessionInfoFromSession() { |
||||
|
// GIVEN
|
||||
|
long expectedTime = 123L; |
||||
|
boolean expectedAttributesSubscription = true; |
||||
|
boolean expectedRPCSubscription = true; |
||||
|
TransportProtos.SessionInfoProto expectedSessionInfo = TransportProtos.SessionInfoProto.getDefaultInstance(); |
||||
|
|
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
SessionMetaData session = new SessionMetaData(expectedSessionInfo, TransportProtos.SessionType.ASYNC, listenerMock); |
||||
|
session.setSubscribedToAttributes(expectedAttributesSubscription); |
||||
|
session.setSubscribedToRPC(expectedRPCSubscription); |
||||
|
sessions.put(SESSION_ID, session); |
||||
|
|
||||
|
ActivityReportCallback<UUID> callbackMock = mock(ActivityReportCallback.class); |
||||
|
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
|
||||
|
doCallRealMethod().when(transportServiceMock).reportActivity(SESSION_ID, sessionInfo, expectedTime, callbackMock); |
||||
|
|
||||
|
// WHEN
|
||||
|
transportServiceMock.reportActivity(SESSION_ID, sessionInfo, expectedTime, callbackMock); |
||||
|
|
||||
|
// THEN
|
||||
|
ArgumentCaptor<TransportProtos.SessionInfoProto> sessionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SessionInfoProto.class); |
||||
|
ArgumentCaptor<TransportProtos.SubscriptionInfoProto> subscriptionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SubscriptionInfoProto.class); |
||||
|
ArgumentCaptor<TransportServiceCallback<Void>> callbackCaptor = ArgumentCaptor.forClass(TransportServiceCallback.class); |
||||
|
|
||||
|
verify(transportServiceMock).process(sessionInfoCaptor.capture(), subscriptionInfoCaptor.capture(), callbackCaptor.capture()); |
||||
|
|
||||
|
assertThat(sessionInfoCaptor.getValue()).isEqualTo(expectedSessionInfo); |
||||
|
|
||||
|
TransportProtos.SubscriptionInfoProto expectedSubscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() |
||||
|
.setAttributeSubscription(expectedAttributesSubscription) |
||||
|
.setRpcSubscription(expectedRPCSubscription) |
||||
|
.setLastActivityTime(expectedTime) |
||||
|
.build(); |
||||
|
assertThat(subscriptionInfoCaptor.getValue()).isEqualTo(expectedSubscriptionInfo); |
||||
|
|
||||
|
TransportServiceCallback<Void> queueCallback = callbackCaptor.getValue(); |
||||
|
|
||||
|
queueCallback.onSuccess(null); |
||||
|
verify(callbackMock).onSuccess(SESSION_ID, expectedTime); |
||||
|
|
||||
|
var throwable = new Throwable(); |
||||
|
queueCallback.onError(throwable); |
||||
|
verify(callbackMock).onFailure(SESSION_ID, throwable); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenKeyAndTimeToReportAndSessionDoesNotExist_whenReportingActivity_thenShouldReportActivityWithNoSubscriptionsAndPreviousSessionInfo() { |
||||
|
// GIVEN
|
||||
|
long expectedTime = 123L; |
||||
|
boolean expectedAttributesSubscription = false; |
||||
|
boolean expectedRPCSubscription = false; |
||||
|
TransportProtos.SessionInfoProto expectedSessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
|
||||
|
ActivityReportCallback<UUID> callbackMock = mock(ActivityReportCallback.class); |
||||
|
|
||||
|
doCallRealMethod().when(transportServiceMock).reportActivity(SESSION_ID, expectedSessionInfo, expectedTime, callbackMock); |
||||
|
|
||||
|
// WHEN
|
||||
|
transportServiceMock.reportActivity(SESSION_ID, expectedSessionInfo, expectedTime, callbackMock); |
||||
|
|
||||
|
// THEN
|
||||
|
ArgumentCaptor<TransportProtos.SessionInfoProto> sessionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SessionInfoProto.class); |
||||
|
ArgumentCaptor<TransportProtos.SubscriptionInfoProto> subscriptionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SubscriptionInfoProto.class); |
||||
|
ArgumentCaptor<TransportServiceCallback<Void>> callbackCaptor = ArgumentCaptor.forClass(TransportServiceCallback.class); |
||||
|
|
||||
|
verify(transportServiceMock).process(sessionInfoCaptor.capture(), subscriptionInfoCaptor.capture(), callbackCaptor.capture()); |
||||
|
|
||||
|
assertThat(sessionInfoCaptor.getValue()).isEqualTo(expectedSessionInfo); |
||||
|
|
||||
|
TransportProtos.SubscriptionInfoProto expectedSubscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() |
||||
|
.setAttributeSubscription(expectedAttributesSubscription) |
||||
|
.setRpcSubscription(expectedRPCSubscription) |
||||
|
.setLastActivityTime(expectedTime) |
||||
|
.build(); |
||||
|
assertThat(subscriptionInfoCaptor.getValue()).isEqualTo(expectedSubscriptionInfo); |
||||
|
|
||||
|
TransportServiceCallback<Void> queueCallback = callbackCaptor.getValue(); |
||||
|
|
||||
|
queueCallback.onSuccess(null); |
||||
|
verify(callbackMock).onSuccess(SESSION_ID, expectedTime); |
||||
|
|
||||
|
var throwable = new Throwable(); |
||||
|
queueCallback.onError(throwable); |
||||
|
verify(callbackMock).onFailure(SESSION_ID, throwable); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenActivityHappened_whenRecordActivity_thenShouldDelegateToOnActivity() { |
||||
|
// GIVEN
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
doCallRealMethod().when(transportServiceMock).recordActivity(sessionInfo); |
||||
|
when(transportServiceMock.toSessionId(sessionInfo)).thenReturn(SESSION_ID); |
||||
|
long expectedTime = 123L; |
||||
|
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(expectedTime); |
||||
|
|
||||
|
// WHEN
|
||||
|
transportServiceMock.recordActivity(sessionInfo); |
||||
|
|
||||
|
// THEN
|
||||
|
verify(transportServiceMock).onActivity(SESSION_ID, expectedTime); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenKey_whenCreatingNewState_thenShouldCorrectlyCreateNewEmptyState() { |
||||
|
// GIVEN
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, null)); |
||||
|
|
||||
|
when(transportServiceMock.createNewState(SESSION_ID)).thenCallRealMethod(); |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> expectedState = new ActivityState<>(); |
||||
|
expectedState.setMetadata(sessionInfo); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> actualState = transportServiceMock.createNewState(SESSION_ID); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(actualState).isEqualTo(expectedState); |
||||
|
} |
||||
|
|
||||
|
@ParameterizedTest |
||||
|
@EnumSource(ActivityStrategyType.class) |
||||
|
void givenDifferentReportingStrategies_whenGettingStrategy_thenShouldReturnCorrectStrategy(ActivityStrategyType reportingStrategyType) { |
||||
|
// GIVEN
|
||||
|
doCallRealMethod().when(transportServiceMock).getStrategy(); |
||||
|
ReflectionTestUtils.setField(transportServiceMock, "reportingStrategyType", reportingStrategyType); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityStrategy actualStrategy = transportServiceMock.getStrategy(); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(actualStrategy).isEqualTo(reportingStrategyType.toStrategy()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenSessionDoesNotExist_whenUpdatingActivityState_thenShouldReturnNull() { |
||||
|
// GIVEN
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setLastRecordedTime(123L); |
||||
|
state.setMetadata(sessionInfo); |
||||
|
|
||||
|
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(updatedState).isNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenNoGwSessionId_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() { |
||||
|
// GIVEN
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); |
||||
|
|
||||
|
long lastRecordedTime = 123L; |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setLastRecordedTime(lastRecordedTime); |
||||
|
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); |
||||
|
|
||||
|
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(updatedState).isSameAs(state); |
||||
|
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); |
||||
|
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenHasGwSessionIdButGwSessionIsNotNull_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() { |
||||
|
// GIVEN
|
||||
|
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); |
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); |
||||
|
|
||||
|
long lastRecordedTime = 123L; |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setLastRecordedTime(lastRecordedTime); |
||||
|
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); |
||||
|
|
||||
|
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(updatedState).isSameAs(state); |
||||
|
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); |
||||
|
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); |
||||
|
|
||||
|
verify(transportServiceMock, never()).getLastRecordedTime(gwSessionId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenHasGwSessionWithoutOverwriteEnabled_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() { |
||||
|
// GIVEN
|
||||
|
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); |
||||
|
TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener gwListenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(gwSessionId, new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock)); |
||||
|
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); |
||||
|
|
||||
|
long lastRecordedTime = 123L; |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setLastRecordedTime(lastRecordedTime); |
||||
|
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); |
||||
|
|
||||
|
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(updatedState).isSameAs(state); |
||||
|
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); |
||||
|
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); |
||||
|
|
||||
|
verify(transportServiceMock, never()).getLastRecordedTime(gwSessionId); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenHasGwSessionWithOverwriteEnabledAndGwLastRecordedTimeIsGreater_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfoAndLastRecordedTime() { |
||||
|
// GIVEN
|
||||
|
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); |
||||
|
TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener gwListenerMock = mock(SessionMsgListener.class); |
||||
|
SessionMetaData gwSession = new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock); |
||||
|
gwSession.setOverwriteActivityTime(true); |
||||
|
sessions.put(gwSessionId, gwSession); |
||||
|
|
||||
|
long gwLastRecordedTime = 500L; |
||||
|
when(transportServiceMock.getLastRecordedTime(gwSessionId)).thenReturn(gwLastRecordedTime); |
||||
|
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); |
||||
|
|
||||
|
long lastRecordedTime = 123L; |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setLastRecordedTime(lastRecordedTime); |
||||
|
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); |
||||
|
|
||||
|
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(updatedState).isSameAs(state); |
||||
|
assertThat(updatedState.getLastRecordedTime()).isEqualTo(gwLastRecordedTime); |
||||
|
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenHasGwSessionWithOverwriteEnabledAndGwLastRecordedTimeIsLess_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfoOnly() { |
||||
|
// GIVEN
|
||||
|
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002"); |
||||
|
TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener gwListenerMock = mock(SessionMsgListener.class); |
||||
|
SessionMetaData gwSession = new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock); |
||||
|
gwSession.setOverwriteActivityTime(true); |
||||
|
sessions.put(gwSessionId, gwSession); |
||||
|
|
||||
|
long gwLastRecordedTime = 100L; |
||||
|
when(transportServiceMock.getLastRecordedTime(gwSessionId)).thenReturn(gwLastRecordedTime); |
||||
|
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits()) |
||||
|
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); |
||||
|
|
||||
|
long lastRecordedTime = 123L; |
||||
|
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>(); |
||||
|
state.setLastRecordedTime(lastRecordedTime); |
||||
|
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance()); |
||||
|
|
||||
|
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(updatedState).isSameAs(state); |
||||
|
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime); |
||||
|
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo); |
||||
|
} |
||||
|
|
||||
|
@ParameterizedTest |
||||
|
@MethodSource("provideTestParamsForHasExpiredTrue") |
||||
|
public void givenExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnTrue(long currentTimeMillis, long lastRecordedTime, long sessionInactivityTimeout) { |
||||
|
// GIVEN
|
||||
|
ReflectionTestUtils.setField(transportServiceMock, "sessionInactivityTimeout", sessionInactivityTimeout); |
||||
|
|
||||
|
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis); |
||||
|
when(transportServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
boolean hasExpired = transportServiceMock.hasExpired(lastRecordedTime); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(hasExpired).isTrue(); |
||||
|
} |
||||
|
|
||||
|
private static Stream<Arguments> provideTestParamsForHasExpiredTrue() { |
||||
|
return Stream.of( |
||||
|
Arguments.of(10L, 0L, 9L), |
||||
|
Arguments.of(10L, 7L, 2L), |
||||
|
Arguments.of(10L, 8L, 1L), |
||||
|
Arguments.of(10000L, 5000L, 3000L) |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@ParameterizedTest |
||||
|
@MethodSource("provideTestParamsForHasExpiredFalse") |
||||
|
public void givenNotExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnFalse(long currentTimeMillis, long lastRecordedTime, long sessionInactivityTimeout) { |
||||
|
// GIVEN
|
||||
|
ReflectionTestUtils.setField(transportServiceMock, "sessionInactivityTimeout", sessionInactivityTimeout); |
||||
|
|
||||
|
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis); |
||||
|
when(transportServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod(); |
||||
|
|
||||
|
// WHEN
|
||||
|
boolean hasExpired = transportServiceMock.hasExpired(lastRecordedTime); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(hasExpired).isFalse(); |
||||
|
} |
||||
|
|
||||
|
private static Stream<Arguments> provideTestParamsForHasExpiredFalse() { |
||||
|
return Stream.of( |
||||
|
Arguments.of(10L, 9L, 2L), |
||||
|
Arguments.of(10L, 0L, 11L), |
||||
|
Arguments.of(10L, 8L, 3L), |
||||
|
Arguments.of(10000L, 8000L, 3000L) |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenSessionExists_whenOnStateExpiryCalled_thenShouldPerformExpirationActions() { |
||||
|
// GIVEN
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
SessionMsgListener listenerMock = mock(SessionMsgListener.class); |
||||
|
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock)); |
||||
|
doCallRealMethod().when(transportServiceMock).onStateExpiry(SESSION_ID, sessionInfo); |
||||
|
|
||||
|
// WHEN
|
||||
|
transportServiceMock.onStateExpiry(SESSION_ID, sessionInfo); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(sessions.containsKey(SESSION_ID)).isFalse(); |
||||
|
verify(transportServiceMock).deregisterSession(sessionInfo); |
||||
|
verify(transportServiceMock).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); |
||||
|
verify(listenerMock).onRemoteSessionCloseCommand(SESSION_ID, SESSION_EXPIRED_NOTIFICATION_PROTO); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void givenSessionDoesNotExist_whenOnStateExpiryCalled_thenShouldNotPerformExpirationActions() { |
||||
|
// GIVEN
|
||||
|
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder() |
||||
|
.setSessionIdMSB(SESSION_ID.getMostSignificantBits()) |
||||
|
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits()) |
||||
|
.build(); |
||||
|
doCallRealMethod().when(transportServiceMock).onStateExpiry(SESSION_ID, sessionInfo); |
||||
|
|
||||
|
// WHEN
|
||||
|
transportServiceMock.onStateExpiry(SESSION_ID, sessionInfo); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(sessions.containsKey(SESSION_ID)).isFalse(); |
||||
|
verify(transportServiceMock, never()).deregisterSession(sessionInfo); |
||||
|
verify(transportServiceMock, never()).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,45 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.dao.util; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.kv.IntervalType; |
||||
|
|
||||
|
import java.time.Instant; |
||||
|
import java.time.ZoneId; |
||||
|
import java.time.ZonedDateTime; |
||||
|
import java.time.temporal.ChronoUnit; |
||||
|
import java.time.temporal.IsoFields; |
||||
|
import java.time.temporal.WeekFields; |
||||
|
|
||||
|
public class TimeUtils { |
||||
|
|
||||
|
public static long calculateIntervalEnd(long startTs, IntervalType intervalType, ZoneId tzId) { |
||||
|
var startTime = ZonedDateTime.ofInstant(Instant.ofEpochMilli(startTs), tzId); |
||||
|
switch (intervalType) { |
||||
|
case WEEK: |
||||
|
return startTime.truncatedTo(ChronoUnit.DAYS).with(WeekFields.SUNDAY_START.dayOfWeek(), 1).plusDays(7).toInstant().toEpochMilli(); |
||||
|
case WEEK_ISO: |
||||
|
return startTime.truncatedTo(ChronoUnit.DAYS).with(WeekFields.ISO.dayOfWeek(), 1).plusDays(7).toInstant().toEpochMilli(); |
||||
|
case MONTH: |
||||
|
return startTime.truncatedTo(ChronoUnit.DAYS).withDayOfMonth(1).plusMonths(1).toInstant().toEpochMilli(); |
||||
|
case QUARTER: |
||||
|
return startTime.truncatedTo(ChronoUnit.DAYS).with(IsoFields.DAY_OF_QUARTER, 1).plusMonths(3).toInstant().toEpochMilli(); |
||||
|
default: |
||||
|
throw new RuntimeException("Not supported!"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,60 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.dao.util; |
||||
|
|
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.thingsboard.server.common.data.kv.IntervalType; |
||||
|
|
||||
|
import java.time.ZoneId; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
|
||||
|
class TimeUtilsTest { |
||||
|
|
||||
|
@Test |
||||
|
void testWeekEnd() { |
||||
|
long ts = 1704899727000L; // Wednesday, January 10 15:15:27 GMT
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.WEEK, ZoneId.of("Europe/Kyiv"))).isEqualTo(1705183200000L); // Sunday, January 14, 2024 0:00:00 GMT+02:00
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.WEEK_ISO, ZoneId.of("Europe/Kyiv"))).isEqualTo(1705269600000L); // Monday, January 15, 2024 0:00:00 GMT+02:00
|
||||
|
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.WEEK, ZoneId.of("Europe/Amsterdam"))).isEqualTo(1705186800000L); // Sunday, January 14, 2024 0:00:00 GMT+01:00
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.WEEK_ISO, ZoneId.of("Europe/Amsterdam"))).isEqualTo(1705273200000L); // Monday, January 15, 2024 0:00:00 GMT+01:00
|
||||
|
|
||||
|
ts = 1704621600000L; // Sunday, January 7, 2024 12:00:00 GMT+02:00
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.WEEK, ZoneId.of("Europe/Kyiv"))).isEqualTo(1705183200000L); // Sunday, January 14, 2024 0:00:00 GMT+02:00
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.WEEK_ISO, ZoneId.of("Europe/Kyiv"))).isEqualTo(1704664800000L); // Monday, January 8, 2024 0:00:00 GMT+02:00
|
||||
|
} |
||||
|
|
||||
|
|
||||
|
@Test |
||||
|
void testMonthEnd() { |
||||
|
long ts = 1704899727000L; // Wednesday, January 10 15:15:27 GMT
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.MONTH, ZoneId.of("Europe/Kyiv"))).isEqualTo(1706738400000L); // Thursday, February 1, 2024 0:00:00 GMT+02:00
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.MONTH, ZoneId.of("Europe/Amsterdam"))).isEqualTo(1706742000000L); // Monday, January 15, 2024 0:00:00 GMT+02:00
|
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testQuarterEnd() { |
||||
|
long ts = 1704899727000L; // Wednesday, January 10 15:15:27 GMT
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.QUARTER, ZoneId.of("Europe/Kyiv"))).isEqualTo(1711918800000L); // Monday, April 1, 2024 0:00:00 GMT+03:00 DST
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.QUARTER, ZoneId.of("Europe/Amsterdam"))).isEqualTo(1711922400000L); // Monday, April 1, 2024 1:00:00 GMT+03:00 DST
|
||||
|
|
||||
|
ts = 1711929600000L; // Monday, April 1, 2024 3:00:00 GMT+03:00
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.QUARTER, ZoneId.of("Europe/Kyiv"))).isEqualTo(1719781200000L); // Monday, July 1, 2024 0:00:00 GMT+03:00 DST
|
||||
|
assertThat(TimeUtils.calculateIntervalEnd(ts, IntervalType.QUARTER, ZoneId.of("America/New_York"))).isEqualTo(1711944000000L); // Monday, April 1, 2024 7:00:00 GMT+03:00 DST
|
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,204 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.msa.rule.node; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import lombok.Data; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.awaitility.Awaitility; |
||||
|
import org.eclipse.paho.client.mqttv3.IMqttMessageListener; |
||||
|
import org.eclipse.paho.client.mqttv3.MqttClient; |
||||
|
import org.eclipse.paho.client.mqttv3.MqttConnectOptions; |
||||
|
import org.eclipse.paho.client.mqttv3.MqttMessage; |
||||
|
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; |
||||
|
import org.testng.annotations.AfterMethod; |
||||
|
import org.testng.annotations.BeforeMethod; |
||||
|
import org.testng.annotations.Test; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.EventInfo; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.event.EventType; |
||||
|
import org.thingsboard.server.common.data.id.RuleChainId; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.common.data.page.TimePageLink; |
||||
|
import org.thingsboard.server.common.data.rule.NodeConnectionInfo; |
||||
|
import org.thingsboard.server.common.data.rule.RuleChain; |
||||
|
import org.thingsboard.server.common.data.rule.RuleChainMetaData; |
||||
|
import org.thingsboard.server.common.data.rule.RuleNode; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentials; |
||||
|
import org.thingsboard.server.msa.AbstractContainerTest; |
||||
|
import org.thingsboard.server.msa.DisableUIListeners; |
||||
|
import org.thingsboard.server.msa.TestProperties; |
||||
|
import org.thingsboard.server.msa.WsClient; |
||||
|
import org.thingsboard.server.msa.mapper.WsTelemetryResponse; |
||||
|
|
||||
|
import java.util.Arrays; |
||||
|
import java.util.List; |
||||
|
import java.util.Objects; |
||||
|
import java.util.Optional; |
||||
|
import java.util.concurrent.ArrayBlockingQueue; |
||||
|
import java.util.concurrent.BlockingQueue; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.testng.Assert.fail; |
||||
|
import static org.thingsboard.server.msa.prototypes.DevicePrototypes.defaultDevicePrototype; |
||||
|
|
||||
|
@DisableUIListeners |
||||
|
@Slf4j |
||||
|
public class MqttNodeTest extends AbstractContainerTest { |
||||
|
|
||||
|
private static final String TOPIC = "tb/mqtt/device"; |
||||
|
|
||||
|
private Device device; |
||||
|
|
||||
|
@BeforeMethod |
||||
|
public void setUp() { |
||||
|
testRestClient.login("tenant@thingsboard.org", "tenant"); |
||||
|
device = testRestClient.postDevice("", defaultDevicePrototype("mqtt_")); |
||||
|
} |
||||
|
|
||||
|
@AfterMethod |
||||
|
public void tearDown() { |
||||
|
testRestClient.deleteDeviceIfExists(device.getId()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void telemetryUpload() throws Exception { |
||||
|
RuleChainId defaultRuleChainId = getDefaultRuleChainId(); |
||||
|
|
||||
|
createRootRuleChainWithTestNode("MqttRuleNodeTestMetadata.json", "org.thingsboard.rule.engine.mqtt.TbMqttNode", 2); |
||||
|
|
||||
|
DeviceCredentials deviceCredentials = testRestClient.getDeviceCredentialsByDeviceId(device.getId()); |
||||
|
|
||||
|
WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS); |
||||
|
|
||||
|
MqttMessageListener messageListener = new MqttMessageListener(); |
||||
|
MqttClient responseClient = new MqttClient(TestProperties.getMqttBrokerUrl(), StringUtils.randomAlphanumeric(10), new MemoryPersistence()); |
||||
|
responseClient.connect(); |
||||
|
responseClient.subscribe(TOPIC, messageListener); |
||||
|
|
||||
|
MqttClient mqttClient = new MqttClient("tcp://localhost:1883", StringUtils.randomAlphanumeric(10), new MemoryPersistence()); |
||||
|
MqttConnectOptions mqttConnectOptions = new MqttConnectOptions(); |
||||
|
mqttConnectOptions.setUserName(deviceCredentials.getCredentialsId()); |
||||
|
mqttClient.connect(mqttConnectOptions); |
||||
|
mqttClient.publish("v1/devices/me/telemetry", new MqttMessage(createPayload().toString().getBytes())); |
||||
|
|
||||
|
WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); |
||||
|
log.info("Received telemetry: {}", actualLatestTelemetry); |
||||
|
wsClient.closeBlocking(); |
||||
|
|
||||
|
assertThat(actualLatestTelemetry.getData()).hasSize(4); |
||||
|
assertThat(actualLatestTelemetry.getLatestValues().keySet()).containsOnlyOnceElementsOf(Arrays.asList("booleanKey", "stringKey", "doubleKey", "longKey")); |
||||
|
|
||||
|
assertThat(actualLatestTelemetry.getDataValuesByKey("booleanKey").get(1)).isEqualTo(Boolean.TRUE.toString()); |
||||
|
assertThat(actualLatestTelemetry.getDataValuesByKey("stringKey").get(1)).isEqualTo("value1"); |
||||
|
assertThat(actualLatestTelemetry.getDataValuesByKey("doubleKey").get(1)).isEqualTo(Double.toString(42.0)); |
||||
|
assertThat(actualLatestTelemetry.getDataValuesByKey("longKey").get(1)).isEqualTo(Long.toString(73)); |
||||
|
|
||||
|
Awaitility |
||||
|
.await() |
||||
|
.alias("Get integration events") |
||||
|
.atMost(10, TimeUnit.SECONDS) |
||||
|
.until(() -> messageListener.getEvents().size() > 0); |
||||
|
|
||||
|
BlockingQueue<MqttEvent> events = messageListener.getEvents(); |
||||
|
JsonNode actual = JacksonUtil.toJsonNode(Objects.requireNonNull(events.poll()).message); |
||||
|
|
||||
|
assertThat(actual.get("stringKey").asText()).isEqualTo("value1"); |
||||
|
assertThat(actual.get("booleanKey").asBoolean()).isEqualTo(Boolean.TRUE); |
||||
|
assertThat(actual.get("doubleKey").asDouble()).isEqualTo(42.0); |
||||
|
assertThat(actual.get("longKey").asLong()).isEqualTo(73); |
||||
|
|
||||
|
testRestClient.setRootRuleChain(defaultRuleChainId); |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
private class MqttMessageListener implements IMqttMessageListener { |
||||
|
private final BlockingQueue<MqttEvent> events; |
||||
|
|
||||
|
private MqttMessageListener() { |
||||
|
events = new ArrayBlockingQueue<>(100); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void messageArrived(String s, MqttMessage mqttMessage) { |
||||
|
log.info("MQTT message [{}], topic [{}]", mqttMessage.toString(), s); |
||||
|
events.add(new MqttEvent(s, mqttMessage.toString())); |
||||
|
} |
||||
|
|
||||
|
public BlockingQueue<MqttEvent> getEvents() { |
||||
|
return events; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
private class MqttEvent { |
||||
|
private final String topic; |
||||
|
private final String message; |
||||
|
} |
||||
|
|
||||
|
private RuleChainId getDefaultRuleChainId() { |
||||
|
PageData<RuleChain> ruleChains = testRestClient.getRuleChains(new PageLink(40, 0)); |
||||
|
|
||||
|
Optional<RuleChain> defaultRuleChain = ruleChains.getData() |
||||
|
.stream() |
||||
|
.filter(RuleChain::isRoot) |
||||
|
.findFirst(); |
||||
|
if (!defaultRuleChain.isPresent()) { |
||||
|
fail("Root rule chain wasn't found"); |
||||
|
} |
||||
|
return defaultRuleChain.get().getId(); |
||||
|
} |
||||
|
|
||||
|
protected RuleChainId createRootRuleChainWithTestNode(String ruleChainMetadataFile, String ruleNodeType, int eventsCount) throws Exception { |
||||
|
RuleChain newRuleChain = new RuleChain(); |
||||
|
newRuleChain.setName("testRuleChain"); |
||||
|
RuleChain ruleChain = testRestClient.postRuleChain(newRuleChain); |
||||
|
|
||||
|
JsonNode configuration = JacksonUtil.OBJECT_MAPPER.readTree(this.getClass().getClassLoader().getResourceAsStream(ruleChainMetadataFile)); |
||||
|
RuleChainMetaData ruleChainMetaData = new RuleChainMetaData(); |
||||
|
ruleChainMetaData.setRuleChainId(ruleChain.getId()); |
||||
|
ruleChainMetaData.setFirstNodeIndex(configuration.get("firstNodeIndex").asInt()); |
||||
|
ruleChainMetaData.setNodes(Arrays.asList(JacksonUtil.OBJECT_MAPPER.treeToValue(configuration.get("nodes"), RuleNode[].class))); |
||||
|
ruleChainMetaData.setConnections(Arrays.asList(JacksonUtil.OBJECT_MAPPER.treeToValue(configuration.get("connections"), NodeConnectionInfo[].class))); |
||||
|
|
||||
|
ruleChainMetaData = testRestClient.postRuleChainMetadata(ruleChainMetaData); |
||||
|
|
||||
|
testRestClient.setRootRuleChain(ruleChain.getId()); |
||||
|
|
||||
|
RuleNode node = ruleChainMetaData.getNodes().stream().filter(ruleNode -> ruleNode.getType().equals(ruleNodeType)).findFirst().get(); |
||||
|
|
||||
|
Awaitility |
||||
|
.await() |
||||
|
.alias("Get events from rule chain") |
||||
|
.atMost(10, TimeUnit.SECONDS) |
||||
|
.until(() -> { |
||||
|
PageData<EventInfo> events = testRestClient.getEvents(node.getId(), EventType.LC_EVENT, ruleChain.getTenantId(), new TimePageLink(1024)); |
||||
|
List<EventInfo> eventInfos = events.getData().stream().filter(eventInfo -> |
||||
|
"STARTED".equals(eventInfo.getBody().get("event").asText()) && |
||||
|
"true".equals(eventInfo.getBody().get("success").asText())) |
||||
|
.collect(Collectors.toList()); |
||||
|
|
||||
|
return eventInfos.size() == eventsCount; |
||||
|
}); |
||||
|
|
||||
|
return ruleChain.getId(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,79 @@ |
|||||
|
{ |
||||
|
"firstNodeIndex": 2, |
||||
|
"nodes": [ |
||||
|
{ |
||||
|
"additionalInfo": { |
||||
|
"description": "", |
||||
|
"layoutX": 626, |
||||
|
"layoutY": 152 |
||||
|
}, |
||||
|
"type": "org.thingsboard.rule.engine.mqtt.TbMqttNode", |
||||
|
"name": "test mqtt", |
||||
|
"debugMode": true, |
||||
|
"singletonMode": true, |
||||
|
"queueName": "HighPriority", |
||||
|
"configurationVersion": 0, |
||||
|
"configuration": { |
||||
|
"topicPattern": "tb/mqtt/device", |
||||
|
"host": "broker", |
||||
|
"port": 1883, |
||||
|
"connectTimeoutSec": 10, |
||||
|
"clientId": null, |
||||
|
"cleanSession": true, |
||||
|
"retainedMessage": false, |
||||
|
"ssl": false, |
||||
|
"credentials": { |
||||
|
"type": "anonymous" |
||||
|
} |
||||
|
}, |
||||
|
"externalId": null |
||||
|
}, |
||||
|
{ |
||||
|
"additionalInfo": { |
||||
|
"description": "", |
||||
|
"layoutX": 949, |
||||
|
"layoutY": 153 |
||||
|
}, |
||||
|
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode", |
||||
|
"name": "save timeseries", |
||||
|
"debugMode": true, |
||||
|
"singletonMode": false, |
||||
|
"configurationVersion": 0, |
||||
|
"configuration": { |
||||
|
"defaultTTL": 0, |
||||
|
"skipLatestPersistence": false, |
||||
|
"useServerTs": false |
||||
|
}, |
||||
|
"externalId": null |
||||
|
}, |
||||
|
{ |
||||
|
"additionalInfo": { |
||||
|
"description": "", |
||||
|
"layoutX": 305, |
||||
|
"layoutY": 151 |
||||
|
}, |
||||
|
"type": "org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode", |
||||
|
"name": "switch", |
||||
|
"debugMode": false, |
||||
|
"singletonMode": false, |
||||
|
"configurationVersion": 0, |
||||
|
"configuration": { |
||||
|
"version": 0 |
||||
|
}, |
||||
|
"externalId": null |
||||
|
} |
||||
|
], |
||||
|
"connections": [ |
||||
|
{ |
||||
|
"fromIndex": 0, |
||||
|
"toIndex": 1, |
||||
|
"type": "Success" |
||||
|
}, |
||||
|
{ |
||||
|
"fromIndex": 2, |
||||
|
"toIndex": 0, |
||||
|
"type": "Post telemetry" |
||||
|
} |
||||
|
], |
||||
|
"ruleChainConnections": null |
||||
|
} |
||||
@ -0,0 +1,25 @@ |
|||||
|
# |
||||
|
# Copyright © 2016-2024 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. |
||||
|
# |
||||
|
|
||||
|
version: '3.0' |
||||
|
services: |
||||
|
broker: |
||||
|
image: eclipse-mosquitto |
||||
|
volumes: |
||||
|
- ./mosquitto/mosquitto.conf:/mosquitto/config/mosquitto.conf |
||||
|
ports: |
||||
|
- "1883" |
||||
|
restart: always |
||||
@ -0,0 +1,18 @@ |
|||||
|
# |
||||
|
# Copyright © 2016-2024 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. |
||||
|
# |
||||
|
|
||||
|
listener 1883 |
||||
|
allow_anonymous true |
||||
File diff suppressed because one or more lines are too long
@ -0,0 +1,52 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
||||
|
import org.junit.jupiter.params.ParameterizedTest; |
||||
|
import org.junit.jupiter.params.provider.MethodSource; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.rule.engine.api.TbNode; |
||||
|
import org.thingsboard.rule.engine.api.TbNodeException; |
||||
|
import org.thingsboard.server.common.data.util.TbPair; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.anyInt; |
||||
|
import static org.mockito.BDDMockito.willCallRealMethod; |
||||
|
|
||||
|
public abstract class AbstractRuleNodeUpgradeTest { |
||||
|
|
||||
|
protected abstract TbNode getTestNode(); |
||||
|
|
||||
|
@ParameterizedTest |
||||
|
@MethodSource |
||||
|
public void givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig(int givenVersion, String givenConfigStr, boolean hasChanges, String expectedConfigStr) throws TbNodeException { |
||||
|
// GIVEN
|
||||
|
willCallRealMethod().given(getTestNode()).upgrade(anyInt(), any()); |
||||
|
JsonNode givenConfig = JacksonUtil.toJsonNode(givenConfigStr); |
||||
|
JsonNode expectedConfig = JacksonUtil.toJsonNode(expectedConfigStr); |
||||
|
|
||||
|
// WHEN
|
||||
|
TbPair<Boolean, JsonNode> upgradeResult = getTestNode().upgrade(givenVersion, givenConfig); |
||||
|
|
||||
|
// THEN
|
||||
|
assertThat(upgradeResult.getFirst()).isEqualTo(hasChanges); |
||||
|
ObjectNode upgradedConfig = (ObjectNode) upgradeResult.getSecond(); |
||||
|
assertThat(upgradedConfig).isEqualTo(expectedConfig); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,53 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.debug; |
||||
|
|
||||
|
import org.junit.jupiter.params.provider.Arguments; |
||||
|
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest; |
||||
|
import org.thingsboard.rule.engine.api.TbNode; |
||||
|
|
||||
|
import java.util.stream.Stream; |
||||
|
|
||||
|
import static org.mockito.Mockito.spy; |
||||
|
|
||||
|
public class TbMsgGeneratorNodeTest extends AbstractRuleNodeUpgradeTest { |
||||
|
|
||||
|
// Rule nodes upgrade
|
||||
|
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { |
||||
|
return Stream.of( |
||||
|
// default config for version 0
|
||||
|
Arguments.of(0, |
||||
|
"{\"msgCount\":0,\"periodInSeconds\":1,\"originatorId\":null,\"originatorType\":null, \"queueName\":null, \"scriptLang\":\"TBEL\",\"jsScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\",\"tbelScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\"}", |
||||
|
true, |
||||
|
"{\"msgCount\":0,\"periodInSeconds\":1,\"originatorId\":null,\"originatorType\":null, \"scriptLang\":\"TBEL\",\"jsScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\",\"tbelScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\"}"), |
||||
|
// default config for version 0 with queueName
|
||||
|
Arguments.of(0, |
||||
|
"{\"msgCount\":0,\"periodInSeconds\":1,\"originatorId\":null,\"originatorType\":null, \"queueName\":\"Main\", \"scriptLang\":\"TBEL\",\"jsScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\",\"tbelScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\"}", |
||||
|
true, |
||||
|
"{\"msgCount\":0,\"periodInSeconds\":1,\"originatorId\":null,\"originatorType\":null, \"scriptLang\":\"TBEL\",\"jsScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\",\"tbelScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\"}"), |
||||
|
// default config for version 1 with upgrade from version 0
|
||||
|
Arguments.of(0, |
||||
|
"{\"msgCount\":0,\"periodInSeconds\":1,\"originatorId\":null,\"originatorType\":null, \"scriptLang\":\"TBEL\",\"jsScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\",\"tbelScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\"}", |
||||
|
false, |
||||
|
"{\"msgCount\":0,\"periodInSeconds\":1,\"originatorId\":null,\"originatorType\":null, \"scriptLang\":\"TBEL\",\"jsScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\",\"tbelScript\":\"var msg = { temp: 42, humidity: 77 };\\nvar metadata = { data: 40 };\\nvar msgType = \\\"POST_TELEMETRY_REQUEST\\\";\\n\\nreturn { msg: msg, metadata: metadata, msgType: msgType };\"}") |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected TbNode getTestNode() { |
||||
|
return spy(TbMsgGeneratorNode.class); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,55 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.flow; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.junit.jupiter.params.provider.Arguments; |
||||
|
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest; |
||||
|
import org.thingsboard.rule.engine.api.TbNode; |
||||
|
|
||||
|
import java.util.stream.Stream; |
||||
|
|
||||
|
import static org.mockito.Mockito.spy; |
||||
|
|
||||
|
@Slf4j |
||||
|
public class TbCheckpointNodeTest extends AbstractRuleNodeUpgradeTest { |
||||
|
|
||||
|
// Rule nodes upgrade
|
||||
|
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { |
||||
|
return Stream.of( |
||||
|
// default config for version 0
|
||||
|
Arguments.of(0, |
||||
|
"{\"queueName\":null}", |
||||
|
true, |
||||
|
"{}"), |
||||
|
// default config for version 0 with queueName
|
||||
|
Arguments.of(0, |
||||
|
"{\"queueName\":\"Main\"}", |
||||
|
true, |
||||
|
"{}"), |
||||
|
// default config for version 1 with upgrade from version 0
|
||||
|
Arguments.of(0, |
||||
|
"{}", |
||||
|
false, |
||||
|
"{}") |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected TbNode getTestNode() { |
||||
|
return spy(TbCheckpointNode.class); |
||||
|
} |
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue