Browse Source

Uplink notifications for PSM & eDRX for CoAP in MSA deployment

pull/4966/head
Andrii Shvaika 5 years ago
parent
commit
60c9e43ea5
  1. 24
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 6
      common/queue/src/main/proto/queue.proto
  3. 20
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java
  4. 8
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java
  5. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java
  6. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  7. 9
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

24
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -59,6 +59,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest;
import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg; import org.thingsboard.server.common.msg.timeout.DeviceActorServerSideRpcTimeoutMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry; import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry;
@ -202,7 +203,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
syncSessionSet.add(key); syncSessionSet.add(key);
} }
}); });
log.trace("46) Rpc syncSessionSet [{}] subscription after sent [{}]", syncSessionSet, rpcSubscriptions); log.trace("Rpc syncSessionSet [{}] subscription after sent [{}]", syncSessionSet, rpcSubscriptions);
syncSessionSet.forEach(rpcSubscriptions::remove); syncSessionSet.forEach(rpcSubscriptions::remove);
} }
@ -318,7 +319,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
.setOneway(request.isOneway()) .setOneway(request.isOneway())
.setPersisted(request.isPersisted()) .setPersisted(request.isPersisted())
.build(); .build();
sendToTransport(rpcRequest, sessionId, nodeId); sendToTransport(rpcRequest, sessionId, nodeId);
}; };
} }
@ -355,9 +355,26 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if (msg.hasPersistedRpcResponseMsg()) { if (msg.hasPersistedRpcResponseMsg()) {
processPersistedRpcResponses(context, sessionInfo, msg.getPersistedRpcResponseMsg()); processPersistedRpcResponses(context, sessionInfo, msg.getPersistedRpcResponseMsg());
} }
if (msg.hasUplinkNotificationMsg()) {
processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg());
}
callback.onSuccess(); callback.onSuccess();
} }
private void processUplinkNotificationMsg(TbActorCtx context, SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg uplinkNotificationMsg) {
String nodeId = sessionInfo.getNodeId();
sessions.entrySet().stream()
.filter(kv -> kv.getValue().getSessionInfo().getNodeId().equals(nodeId) && (kv.getValue().isSubscribedToAttributes() || kv.getValue().isSubscribedToRPC()))
.forEach(kv -> {
ToTransportMsg msg = ToTransportMsg.newBuilder()
.setSessionIdMSB(kv.getKey().getMostSignificantBits())
.setSessionIdLSB(kv.getKey().getLeastSignificantBits())
.setUplinkNotificationMsg(uplinkNotificationMsg)
.build();
systemContext.getTbCoreToTransportService().process(kv.getValue().getSessionInfo().getNodeId(), msg);
});
}
private void handleClaimDeviceMsg(TbActorCtx context, SessionInfoProto sessionInfo, ClaimDeviceMsg msg) { private void handleClaimDeviceMsg(TbActorCtx context, SessionInfoProto sessionInfo, ClaimDeviceMsg msg) {
DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB()));
systemContext.getClaimDevicesService().registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs()); systemContext.getClaimDevicesService().registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs());
@ -599,7 +616,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
void processCredentialsUpdate(TbActorMsg msg) { void processCredentialsUpdate(TbActorMsg msg) {
if (((DeviceCredentialsUpdateNotificationMsg) msg).getDeviceCredentials().getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) { if (((DeviceCredentialsUpdateNotificationMsg) msg).getDeviceCredentials().getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) {
log.info("1) LwM2Mtype: ");
sessions.forEach((k, v) -> { sessions.forEach((k, v) -> {
notifyTransportAboutProfileUpdate(k, v, ((DeviceCredentialsUpdateNotificationMsg) msg).getDeviceCredentials()); notifyTransportAboutProfileUpdate(k, v, ((DeviceCredentialsUpdateNotificationMsg) msg).getDeviceCredentials());
}); });
@ -616,7 +632,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
notifyTransportAboutClosedSession(sessionId, sessionMd, "max concurrent sessions limit reached per device!"); notifyTransportAboutClosedSession(sessionId, sessionMd, "max concurrent sessions limit reached per device!");
} }
private void notifyTransportAboutClosedSession(UUID sessionId, SessionInfoMetaData sessionMd, String message) { private void notifyTransportAboutClosedSession(UUID sessionId, SessionInfoMetaData sessionMd, String message) {
SessionCloseNotificationProto sessionCloseNotificationProto = SessionCloseNotificationProto SessionCloseNotificationProto sessionCloseNotificationProto = SessionCloseNotificationProto
.newBuilder() .newBuilder()
@ -630,7 +645,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
void notifyTransportAboutProfileUpdate(UUID sessionId, SessionInfoMetaData sessionMd, DeviceCredentials deviceCredentials) { void notifyTransportAboutProfileUpdate(UUID sessionId, SessionInfoMetaData sessionMd, DeviceCredentials deviceCredentials) {
log.info("2) LwM2Mtype: ");
ToTransportUpdateCredentialsProto.Builder notification = ToTransportUpdateCredentialsProto.newBuilder(); ToTransportUpdateCredentialsProto.Builder notification = ToTransportUpdateCredentialsProto.newBuilder();
notification.addCredentialsId(deviceCredentials.getCredentialsId()); notification.addCredentialsId(deviceCredentials.getCredentialsId());
notification.addCredentialsValue(deviceCredentials.getCredentialsValue()); notification.addCredentialsValue(deviceCredentials.getCredentialsValue());

6
common/queue/src/main/proto/queue.proto

@ -341,6 +341,10 @@ message ToDeviceRpcResponseMsg {
string payload = 2; string payload = 2;
} }
message UplinkNotificationMsg {
int64 uplinkTs = 1;
}
message ToDevicePersistedRpcResponseMsg { message ToDevicePersistedRpcResponseMsg {
int32 requestId = 1; int32 requestId = 1;
int64 requestIdMSB = 2; int64 requestIdMSB = 2;
@ -453,6 +457,7 @@ message TransportToDeviceActorMsg {
ProvisionDeviceRequestMsg provisionDevice = 9; ProvisionDeviceRequestMsg provisionDevice = 9;
ToDevicePersistedRpcResponseMsg persistedRpcResponseMsg = 10; ToDevicePersistedRpcResponseMsg persistedRpcResponseMsg = 10;
SendPendingRPCMsg sendPendingRPC = 11; SendPendingRPCMsg sendPendingRPC = 11;
UplinkNotificationMsg uplinkNotificationMsg = 12;
} }
message TransportToRuleEngineMsg { message TransportToRuleEngineMsg {
@ -713,6 +718,7 @@ message ToTransportMsg {
ToTransportUpdateCredentialsProto toTransportUpdateCredentialsNotification = 11; ToTransportUpdateCredentialsProto toTransportUpdateCredentialsNotification = 11;
ResourceUpdateMsg resourceUpdateMsg = 12; ResourceUpdateMsg resourceUpdateMsg = 12;
ResourceDeleteMsg resourceDeleteMsg = 13; ResourceDeleteMsg resourceDeleteMsg = 13;
UplinkNotificationMsg uplinkNotificationMsg = 14;
} }
message UsageStatsKVProto{ message UsageStatsKVProto{

20
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java

@ -161,7 +161,7 @@ public class DefaultCoapClientContext implements CoapClientContext {
} }
} }
private void onUplink(TbCoapClientState client) { private void onUplink(TbCoapClientState client, boolean notifyOtherServers, long uplinkTs) {
PowerMode powerMode = client.getPowerMode(); PowerMode powerMode = client.getPowerMode();
PowerSavingConfiguration profileSettings = null; PowerSavingConfiguration profileSettings = null;
if (powerMode == null) { if (powerMode == null) {
@ -174,12 +174,12 @@ public class DefaultCoapClientContext implements CoapClientContext {
} }
} }
if (powerMode == null || PowerMode.DRX.equals(powerMode)) { if (powerMode == null || PowerMode.DRX.equals(powerMode)) {
client.updateLastUplinkTime(); client.updateLastUplinkTime(uplinkTs);
return; return;
} }
client.lock(); client.lock();
try { try {
long uplinkTime = client.updateLastUplinkTime(); long uplinkTime = client.updateLastUplinkTime(uplinkTs);
long timeout; long timeout;
if (PowerMode.PSM.equals(powerMode)) { if (PowerMode.PSM.equals(powerMode)) {
Long psmActivityTimer = client.getPsmActivityTimer(); Long psmActivityTimer = client.getPsmActivityTimer();
@ -214,6 +214,9 @@ public class DefaultCoapClientContext implements CoapClientContext {
return null; return null;
}, timeout, TimeUnit.MILLISECONDS); }, timeout, TimeUnit.MILLISECONDS);
client.setSleepTask(task); client.setSleepTask(task);
if (notifyOtherServers) {
transportService.notifyAboutUplink(getNewSyncSession(client), TransportProtos.UplinkNotificationMsg.newBuilder().setUplinkTs(uplinkTime).build(), TransportServiceCallback.EMPTY);
}
} finally { } finally {
client.unlock(); client.unlock();
} }
@ -544,6 +547,11 @@ public class DefaultCoapClientContext implements CoapClientContext {
log.trace("[{}] Received server rpc response in the wrong session.", state.getSession()); log.trace("[{}] Received server rpc response in the wrong session.", state.getSession());
} }
@Override
public void onUplinkNotification(TransportProtos.UplinkNotificationMsg notificationMsg) {
awake(state, false, notificationMsg.getUplinkTs());
}
private void cancelObserveRelation(TbCoapObservationState attrs) { private void cancelObserveRelation(TbCoapObservationState attrs) {
if (attrs.getObserveRelation() != null) { if (attrs.getObserveRelation() != null) {
attrs.getObserveRelation().cancel(); attrs.getObserveRelation().cancel();
@ -562,7 +570,11 @@ public class DefaultCoapClientContext implements CoapClientContext {
@Override @Override
public boolean awake(TbCoapClientState client) { public boolean awake(TbCoapClientState client) {
onUplink(client); return awake(client, true, System.currentTimeMillis());
}
private boolean awake(TbCoapClientState client, boolean notifyOtherServers, long uplinkTs) {
onUplink(client, notifyOtherServers, uplinkTs);
boolean changed = compareAndSetSleepFlag(client, false); boolean changed = compareAndSetSleepFlag(client, false);
if (changed) { if (changed) {
log.debug("[{}] client is awake", client.getDeviceId()); log.debug("[{}] client is awake", client.getDeviceId());

8
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java

@ -97,9 +97,11 @@ public class TbCoapClientState {
lock.unlock(); lock.unlock();
} }
public long updateLastUplinkTime() { public long updateLastUplinkTime(long ts) {
this.lastUplinkTime = System.currentTimeMillis(); if (ts > lastUplinkTime) {
this.firstEdrxDownlink = true; this.lastUplinkTime = ts;
this.firstEdrxDownlink = true;
}
return lastUplinkTime; return lastUplinkTime;
} }

3
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java

@ -25,6 +25,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionCloseNotifica
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto;
import org.thingsboard.server.gen.transport.TransportProtos.UplinkNotificationMsg;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
@ -44,6 +45,8 @@ public interface SessionMsgListener {
void onToServerRpcResponse(ToServerRpcResponseMsg toServerResponse); void onToServerRpcResponse(ToServerRpcResponseMsg toServerResponse);
default void onUplinkNotification(UplinkNotificationMsg notificationMsg){};
default void onToTransportUpdateCredentials(ToTransportUpdateCredentialsProto toTransportUpdateCredentials){} default void onToTransportUpdateCredentials(ToTransportUpdateCredentialsProto toTransportUpdateCredentials){}
default void onDeviceProfileUpdate(TransportProtos.SessionInfoProto newSessionInfo, DeviceProfile deviceProfile) {} default void onDeviceProfileUpdate(TransportProtos.SessionInfoProto newSessionInfo, DeviceProfile deviceProfile) {}

3
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java

@ -20,6 +20,7 @@ import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
import org.thingsboard.server.common.transport.service.SessionMetaData; import org.thingsboard.server.common.transport.service.SessionMetaData;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetDeviceCredentialsRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetDeviceCredentialsRequestMsg;
@ -128,4 +129,6 @@ public interface TransportService {
void deregisterSession(SessionInfoProto sessionInfo); void deregisterSession(SessionInfoProto sessionInfo);
void log(SessionInfoProto sessionInfo, String msg); void log(SessionInfoProto sessionInfo, String msg);
void notifyAboutUplink(SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg build, TransportServiceCallback<Void> empty);
} }

9
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -571,6 +571,15 @@ public class DefaultTransportService implements TransportService {
} }
} }
@Override
public void notifyAboutUplink(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback);
}
}
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, boolean isFailedRpc, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, boolean isFailedRpc, TransportServiceCallback<Void> callback) {
if (msg.getPersisted()) { if (msg.getPersisted()) {

Loading…
Cancel
Save