From 60c9e43ea51ce4987be91ecad19a496320ec198a Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 23 Jul 2021 13:27:05 +0300 Subject: [PATCH] Uplink notifications for PSM & eDRX for CoAP in MSA deployment --- .../device/DeviceActorMessageProcessor.java | 24 +++++++++++++++---- common/queue/src/main/proto/queue.proto | 6 +++++ .../coap/client/DefaultCoapClientContext.java | 20 ++++++++++++---- .../coap/client/TbCoapClientState.java | 8 ++++--- .../common/transport/SessionMsgListener.java | 3 +++ .../common/transport/TransportService.java | 3 +++ .../service/DefaultTransportService.java | 9 +++++++ 7 files changed, 61 insertions(+), 12 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 30eb24a4d2..42e58ea9e5 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/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.rpc.ToDeviceRpcRequest; 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.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.DeviceSessionsCacheEntry; @@ -202,7 +203,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { 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); } @@ -318,7 +319,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setOneway(request.isOneway()) .setPersisted(request.isPersisted()) .build(); - sendToTransport(rpcRequest, sessionId, nodeId); }; } @@ -355,9 +355,26 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (msg.hasPersistedRpcResponseMsg()) { processPersistedRpcResponses(context, sessionInfo, msg.getPersistedRpcResponseMsg()); } + if (msg.hasUplinkNotificationMsg()) { + processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg()); + } 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) { DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB())); systemContext.getClaimDevicesService().registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs()); @@ -599,7 +616,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { void processCredentialsUpdate(TbActorMsg msg) { if (((DeviceCredentialsUpdateNotificationMsg) msg).getDeviceCredentials().getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) { - log.info("1) LwM2Mtype: "); sessions.forEach((k, v) -> { notifyTransportAboutProfileUpdate(k, v, ((DeviceCredentialsUpdateNotificationMsg) msg).getDeviceCredentials()); }); @@ -616,7 +632,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { notifyTransportAboutClosedSession(sessionId, sessionMd, "max concurrent sessions limit reached per device!"); } - private void notifyTransportAboutClosedSession(UUID sessionId, SessionInfoMetaData sessionMd, String message) { SessionCloseNotificationProto sessionCloseNotificationProto = SessionCloseNotificationProto .newBuilder() @@ -630,7 +645,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } void notifyTransportAboutProfileUpdate(UUID sessionId, SessionInfoMetaData sessionMd, DeviceCredentials deviceCredentials) { - log.info("2) LwM2Mtype: "); ToTransportUpdateCredentialsProto.Builder notification = ToTransportUpdateCredentialsProto.newBuilder(); notification.addCredentialsId(deviceCredentials.getCredentialsId()); notification.addCredentialsValue(deviceCredentials.getCredentialsValue()); diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto index 811704fec0..aebb212dae 100644 --- a/common/queue/src/main/proto/queue.proto +++ b/common/queue/src/main/proto/queue.proto @@ -341,6 +341,10 @@ message ToDeviceRpcResponseMsg { string payload = 2; } +message UplinkNotificationMsg { + int64 uplinkTs = 1; +} + message ToDevicePersistedRpcResponseMsg { int32 requestId = 1; int64 requestIdMSB = 2; @@ -453,6 +457,7 @@ message TransportToDeviceActorMsg { ProvisionDeviceRequestMsg provisionDevice = 9; ToDevicePersistedRpcResponseMsg persistedRpcResponseMsg = 10; SendPendingRPCMsg sendPendingRPC = 11; + UplinkNotificationMsg uplinkNotificationMsg = 12; } message TransportToRuleEngineMsg { @@ -713,6 +718,7 @@ message ToTransportMsg { ToTransportUpdateCredentialsProto toTransportUpdateCredentialsNotification = 11; ResourceUpdateMsg resourceUpdateMsg = 12; ResourceDeleteMsg resourceDeleteMsg = 13; + UplinkNotificationMsg uplinkNotificationMsg = 14; } message UsageStatsKVProto{ diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index 3bc53fefe5..0182badfce 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/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(); PowerSavingConfiguration profileSettings = null; if (powerMode == null) { @@ -174,12 +174,12 @@ public class DefaultCoapClientContext implements CoapClientContext { } } if (powerMode == null || PowerMode.DRX.equals(powerMode)) { - client.updateLastUplinkTime(); + client.updateLastUplinkTime(uplinkTs); return; } client.lock(); try { - long uplinkTime = client.updateLastUplinkTime(); + long uplinkTime = client.updateLastUplinkTime(uplinkTs); long timeout; if (PowerMode.PSM.equals(powerMode)) { Long psmActivityTimer = client.getPsmActivityTimer(); @@ -214,6 +214,9 @@ public class DefaultCoapClientContext implements CoapClientContext { return null; }, timeout, TimeUnit.MILLISECONDS); client.setSleepTask(task); + if (notifyOtherServers) { + transportService.notifyAboutUplink(getNewSyncSession(client), TransportProtos.UplinkNotificationMsg.newBuilder().setUplinkTs(uplinkTime).build(), TransportServiceCallback.EMPTY); + } } finally { client.unlock(); } @@ -544,6 +547,11 @@ public class DefaultCoapClientContext implements CoapClientContext { 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) { if (attrs.getObserveRelation() != null) { attrs.getObserveRelation().cancel(); @@ -562,7 +570,11 @@ public class DefaultCoapClientContext implements CoapClientContext { @Override 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); if (changed) { log.debug("[{}] client is awake", client.getDeviceId()); diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java index f106f57961..88393dbbbc 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapClientState.java @@ -97,9 +97,11 @@ public class TbCoapClientState { lock.unlock(); } - public long updateLastUplinkTime() { - this.lastUplinkTime = System.currentTimeMillis(); - this.firstEdrxDownlink = true; + public long updateLastUplinkTime(long ts) { + if (ts > lastUplinkTime) { + this.lastUplinkTime = ts; + this.firstEdrxDownlink = true; + } return lastUplinkTime; } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java index 644da7f4ec..156cff5f7d 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java +++ b/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.ToServerRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto; +import org.thingsboard.server.gen.transport.TransportProtos.UplinkNotificationMsg; import java.util.Optional; import java.util.UUID; @@ -44,6 +45,8 @@ public interface SessionMsgListener { void onToServerRpcResponse(ToServerRpcResponseMsg toServerResponse); + default void onUplinkNotification(UplinkNotificationMsg notificationMsg){}; + default void onToTransportUpdateCredentials(ToTransportUpdateCredentialsProto toTransportUpdateCredentials){} default void onDeviceProfileUpdate(TransportProtos.SessionInfoProto newSessionInfo, DeviceProfile deviceProfile) {} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index 9aff28d73f..5227f671bc 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/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.ValidateDeviceCredentialsResponse; 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.GetAttributeRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetDeviceCredentialsRequestMsg; @@ -128,4 +129,6 @@ public interface TransportService { void deregisterSession(SessionInfoProto sessionInfo); void log(SessionInfoProto sessionInfo, String msg); + + void notifyAboutUplink(SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg build, TransportServiceCallback empty); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index a1503e86be..93f96795cc 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/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 callback) { + + if (checkLimits(sessionInfo, msg, callback)) { + reportActivityInternal(sessionInfo); + sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback); + } + } + @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, boolean isFailedRpc, TransportServiceCallback callback) { if (msg.getPersisted()) {