diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java index 4282531237..59df37de48 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java @@ -87,7 +87,7 @@ public class LwM2mTransportUtil { public static final String LWM2M_VERSION_DEFAULT = "1.0"; - public static final String LOG_LWM2M_TELEMETRY = "logLwm2m"; + public static final String LOG_LWM2M_TELEMETRY = "transportLog"; public static final String LOG_LWM2M_INFO = "info"; public static final String LOG_LWM2M_ERROR = "error"; public static final String LOG_LWM2M_WARN = "warn"; @@ -169,19 +169,6 @@ public class LwM2mTransportUtil { return lwM2mOtaConvert; } - public static LwM2mNode getLvM2mNodeToObject(LwM2mNode content) { - if (content instanceof LwM2mObject) { - return (LwM2mObject) content; - } else if (content instanceof LwM2mObjectInstance) { - return (LwM2mObjectInstance) content; - } else if (content instanceof LwM2mSingleResource) { - return (LwM2mSingleResource) content; - } else if (content instanceof LwM2mMultipleResource) { - return (LwM2mMultipleResource) content; - } - return null; - } - public static Lwm2mDeviceProfileTransportConfiguration toLwM2MClientProfile(DeviceProfile deviceProfile) { DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); if (transportConfiguration.getType().equals(DeviceTransportType.LWM2M)) { @@ -196,62 +183,6 @@ public class LwM2mTransportUtil { return toLwM2MClientProfile(deviceProfile).getBootstrap(); } - public static JsonObject validateJson(String jsonStr) { - JsonObject object = null; - if (jsonStr != null && !jsonStr.isEmpty()) { - String jsonValidFlesh = jsonStr.replaceAll("\\\\", ""); - jsonValidFlesh = jsonValidFlesh.replaceAll("\n", ""); - jsonValidFlesh = jsonValidFlesh.replaceAll("\t", ""); - jsonValidFlesh = jsonValidFlesh.replaceAll(" ", ""); - String jsonValid = (jsonValidFlesh.charAt(0) == '"' && jsonValidFlesh.charAt(jsonValidFlesh.length() - 1) == '"') ? jsonValidFlesh.substring(1, jsonValidFlesh.length() - 1) : jsonValidFlesh; - try { - object = new JsonParser().parse(jsonValid).getAsJsonObject(); - } catch (JsonSyntaxException e) { - log.error("[{}] Fail validateJson [{}]", jsonStr, e.getMessage()); - } - } - return object; - } - - @SuppressWarnings("unchecked") - public static Optional decode(byte[] byteArray) { - try { - FSTConfiguration config = FSTConfiguration.createDefaultConfiguration(); - T msg = (T) config.asObject(byteArray); - return Optional.ofNullable(msg); - } catch (IllegalArgumentException e) { - log.error("Error during deserialization message, [{}]", e.getMessage()); - return Optional.empty(); - } - } - - public static String splitCamelCaseString(String s) { - LinkedList linkedListOut = new LinkedList<>(); - LinkedList linkedList = new LinkedList((Arrays.asList(s.split(" ")))); - linkedList.forEach(str -> { - String strOut = str.replaceAll("\\W", "").replaceAll("_", "").toUpperCase(); - if (strOut.length() > 1) linkedListOut.add(strOut.charAt(0) + strOut.substring(1).toLowerCase()); - else linkedListOut.add(strOut); - }); - linkedListOut.set(0, (linkedListOut.get(0).substring(0, 1).toLowerCase() + linkedListOut.get(0).substring(1))); - return StringUtils.join(linkedListOut, ""); - } - - public static TransportServiceCallback getAckCallback(LwM2mClient lwM2MClient, - int requestId, String typeTopic) { - return new TransportServiceCallback() { - @Override - public void onSuccess(Void dummy) { - log.trace("[{}] [{}] - requestId [{}] - EndPoint , Access AckCallback", typeTopic, requestId, lwM2MClient.getEndpoint()); - } - - @Override - public void onError(Throwable e) { - log.trace("[{}] Failed to publish msg", e.toString()); - } - }; - } - public static String fromVersionedIdToObjectId(String pathIdVer) { try { if (pathIdVer == null) { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java index ee723dbaa6..0f0d95610b 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java @@ -181,6 +181,7 @@ public class DefaultLwM2MAttributesService implements LwM2MAttributesService { } } }); + clientContext.update(lwM2MClient); // #2.1 lwM2MClient.getSharedAttributes().forEach((pathIdVer, tsKvProto) -> { this.pushUpdateToClientIfNeeded(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java index 5fa0858359..15e023404c 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java @@ -29,7 +29,6 @@ import org.eclipse.leshan.core.node.codec.LwM2mValueConverter; import org.eclipse.leshan.core.request.ContentFormat; import org.eclipse.leshan.server.model.LwM2mModelProvider; import org.eclipse.leshan.server.registration.Registration; -import org.eclipse.leshan.server.security.SecurityInfo; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.device.data.Lwm2mDeviceTransportConfiguration; @@ -38,16 +37,16 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; -import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest; +import java.io.IOException; +import java.io.ObjectInputStream; +import java.io.Serializable; import java.util.Collection; import java.util.Map; import java.util.Optional; -import java.util.Queue; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; @@ -60,48 +59,38 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.f import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.getVerFromPathIdVerOrId; @Slf4j -public class LwM2mClient implements Cloneable { +public class LwM2mClient implements Serializable { + + private static final long serialVersionUID = 8793482946289222623L; private final String nodeId; @Getter private final String endpoint; - private final Lock lock; - @Getter - @Setter - private LwM2MClientState state; + + private transient Lock lock; + @Getter private final Map resources; @Getter private final Map sharedAttributes; - @Getter - private final Queue queuedRequests; - @Getter - private String deviceName; - @Getter - private String deviceProfileName; - - @Getter - private PowerMode powerMode; - - @Getter - private String identity; - @Getter - private SecurityInfo securityInfo; @Getter private TenantId tenantId; @Getter + private UUID profileId; + @Getter private UUID deviceId; @Getter + @Setter + private LwM2MClientState state; + @Getter private SessionInfoProto session; @Getter - private UUID profileId; + private PowerMode powerMode; @Getter @Setter private Registration registration; - private ValidateDeviceCredentialsResponse credentials; - public Object clone() throws CloneNotSupportedException { return super.clone(); } @@ -109,23 +98,17 @@ public class LwM2mClient implements Cloneable { public LwM2mClient(String nodeId, String endpoint) { this.nodeId = nodeId; this.endpoint = endpoint; - this.lock = new ReentrantLock(); this.sharedAttributes = new ConcurrentHashMap<>(); this.resources = new ConcurrentHashMap<>(); - this.queuedRequests = new ConcurrentLinkedQueue<>(); this.state = LwM2MClientState.CREATED; + this.lock = new ReentrantLock(); } - public void init(String identity, SecurityInfo securityInfo, ValidateDeviceCredentialsResponse credentials, UUID sessionId) { - this.identity = identity; - this.securityInfo = securityInfo; - this.credentials = credentials; + public void init(ValidateDeviceCredentialsResponse credentials, UUID sessionId) { this.session = createSession(nodeId, sessionId, credentials); this.tenantId = new TenantId(new UUID(session.getTenantIdMSB(), session.getTenantIdLSB())); this.deviceId = new UUID(session.getDeviceIdMSB(), session.getDeviceIdLSB()); this.profileId = new UUID(session.getDeviceProfileIdMSB(), session.getDeviceProfileIdLSB()); - this.deviceName = session.getDeviceName(); - this.deviceProfileName = session.getDeviceType(); this.powerMode = credentials.getDeviceInfo().getPowerMode(); } @@ -140,10 +123,9 @@ public class LwM2mClient implements Cloneable { public void onDeviceUpdate(Device device, Optional deviceProfileOpt) { SessionInfoProto.Builder builder = SessionInfoProto.newBuilder().mergeFrom(session); this.deviceId = device.getUuidId(); - this.deviceName = device.getName(); builder.setDeviceIdMSB(deviceId.getMostSignificantBits()); builder.setDeviceIdLSB(deviceId.getLeastSignificantBits()); - builder.setDeviceName(deviceName); + builder.setDeviceName(device.getName()); deviceProfileOpt.ifPresent(deviceProfile -> updateSession(deviceProfile, builder)); this.session = builder.build(); this.powerMode = ((Lwm2mDeviceTransportConfiguration) device.getDeviceData().getTransportConfiguration()).getPowerMode(); @@ -156,13 +138,22 @@ public class LwM2mClient implements Cloneable { } private void updateSession(DeviceProfile deviceProfile, SessionInfoProto.Builder builder) { - this.deviceProfileName = deviceProfile.getName(); this.profileId = deviceProfile.getUuidId(); builder.setDeviceProfileIdMSB(profileId.getMostSignificantBits()); builder.setDeviceProfileIdLSB(profileId.getLeastSignificantBits()); - builder.setDeviceType(this.deviceProfileName); + builder.setDeviceType(deviceProfile.getName()); } + public void refreshSessionId(String nodeId) { + UUID newId = UUID.randomUUID(); + SessionInfoProto.Builder builder = SessionInfoProto.newBuilder().mergeFrom(session); + builder.setNodeId(nodeId); + builder.setSessionIdMSB(newId.getMostSignificantBits()); + builder.setSessionIdLSB(newId.getLeastSignificantBits()); + this.session = builder.build(); + } + + private SessionInfoProto createSession(String nodeId, UUID sessionId, ValidateDeviceCredentialsResponse msg) { return SessionInfoProto.newBuilder() .setNodeId(nodeId) @@ -364,5 +355,10 @@ public class LwM2mClient implements Cloneable { } } + private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundException { + in.defaultReadObject(); + this.lock = new ReentrantLock(); + } + } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java index 006953af58..792946e1c3 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java @@ -28,8 +28,6 @@ import java.util.UUID; public interface LwM2mClientContext { - LwM2mClient getClientByRegistrationId(String registrationId); - LwM2mClient getClientByEndpoint(String endpoint); LwM2mClient getClientBySessionInfo(TransportProtos.SessionInfoProto sessionInfo); @@ -53,10 +51,9 @@ public interface LwM2mClientContext { LwM2mClient getClientByDeviceId(UUID deviceId); - String getObjectIdByKeyNameFromProfile(TransportProtos.SessionInfoProto sessionInfo, String keyName); - String getObjectIdByKeyNameFromProfile(LwM2mClient lwM2mClient, String keyName); void registerClient(Registration registration, ValidateDeviceCredentialsResponse credentials); + void update(LwM2mClient lwM2MClient); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java index 700a7aec78..b4b73a0625 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java @@ -24,6 +24,9 @@ import org.eclipse.leshan.server.registration.Registration; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; +import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.transport.TransportDeviceProfileCache; +import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; @@ -31,6 +34,8 @@ import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MSecurityInfo; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.session.LwM2MSessionManager; +import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MClientStore; import org.thingsboard.server.transport.lwm2m.server.store.TbMainSecurityStore; import java.util.Arrays; @@ -56,25 +61,50 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { private final LwM2mTransportContext context; private final LwM2MTransportServerConfig config; private final TbMainSecurityStore securityStore; + private final TbLwM2MClientStore clientStore; + private final LwM2MSessionManager sessionManager; + private final TransportDeviceProfileCache deviceProfileCache; private final Map lwM2mClientsByEndpoint = new ConcurrentHashMap<>(); private final Map lwM2mClientsByRegistrationId = new ConcurrentHashMap<>(); private final Map profiles = new ConcurrentHashMap<>(); @Override public LwM2mClient getClientByEndpoint(String endpoint) { - return lwM2mClientsByEndpoint.computeIfAbsent(endpoint, ep -> new LwM2mClient(context.getNodeId(), ep)); + return lwM2mClientsByEndpoint.computeIfAbsent(endpoint, ep -> { + LwM2mClient client = clientStore.get(ep); + String nodeId = context.getNodeId(); + if (client == null) { + log.info("[{}] initialized new client.", endpoint); + client = new LwM2mClient(nodeId, ep); + } else { + log.debug("[{}] fetched client from store: {}", endpoint, client); + boolean updated = false; + if (client.getRegistration() != null) { + lwM2mClientsByRegistrationId.put(client.getRegistration().getId(), client); + } + if (client.getSession() != null) { + client.refreshSessionId(nodeId); + sessionManager.register(client.getSession()); + updated = true; + } + if (updated) { + clientStore.put(client); + } + } + return client; + }); } @Override - public Optional register(LwM2mClient lwM2MClient, Registration registration) throws LwM2MClientStateException { + public Optional register(LwM2mClient client, Registration registration) throws LwM2MClientStateException { TransportProtos.SessionInfoProto oldSession = null; - lwM2MClient.lock(); + client.lock(); try { - if (LwM2MClientState.UNREGISTERED.equals(lwM2MClient.getState())) { - throw new LwM2MClientStateException(lwM2MClient.getState(), "Client is in invalid state."); + if (LwM2MClientState.UNREGISTERED.equals(client.getState())) { + throw new LwM2MClientStateException(client.getState(), "Client is in invalid state."); } - oldSession = lwM2MClient.getSession(); - TbLwM2MSecurityInfo securityInfo = securityStore.getTbLwM2MSecurityInfoByEndpoint(lwM2MClient.getEndpoint()); + oldSession = client.getSession(); + TbLwM2MSecurityInfo securityInfo = securityStore.getTbLwM2MSecurityInfoByEndpoint(client.getEndpoint()); if (securityInfo.getSecurityMode() != null) { if (SecurityMode.X509.equals(securityInfo.getSecurityMode())) { securityStore.registerX509(registration.getEndpoint(), registration.getId()); @@ -82,54 +112,57 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { if (securityInfo.getDeviceProfile() != null) { profileUpdate(securityInfo.getDeviceProfile()); if (securityInfo.getSecurityInfo() != null) { - lwM2MClient.init(securityInfo.getSecurityInfo().getIdentity(), securityInfo.getSecurityInfo(), securityInfo.getMsg(), UUID.randomUUID()); + client.init(securityInfo.getMsg(), UUID.randomUUID()); } else if (NO_SEC.equals(securityInfo.getSecurityMode())) { - lwM2MClient.init(null, null, securityInfo.getMsg(), UUID.randomUUID()); + client.init(securityInfo.getMsg(), UUID.randomUUID()); } else { - throw new RuntimeException(String.format("Registration failed: device %s not found.", lwM2MClient.getEndpoint())); + throw new RuntimeException(String.format("Registration failed: device %s not found.", client.getEndpoint())); } } else { - throw new RuntimeException(String.format("Registration failed: device %s not found.", lwM2MClient.getEndpoint())); + throw new RuntimeException(String.format("Registration failed: device %s not found.", client.getEndpoint())); } } else { - throw new RuntimeException(String.format("Registration failed: FORBIDDEN, endpointId: %s", lwM2MClient.getEndpoint())); + throw new RuntimeException(String.format("Registration failed: FORBIDDEN, endpointId: %s", client.getEndpoint())); } - lwM2MClient.setRegistration(registration); - this.lwM2mClientsByRegistrationId.put(registration.getId(), lwM2MClient); - lwM2MClient.setState(LwM2MClientState.REGISTERED); + client.setRegistration(registration); + this.lwM2mClientsByRegistrationId.put(registration.getId(), client); + client.setState(LwM2MClientState.REGISTERED); + clientStore.put(client); } finally { - lwM2MClient.unlock(); + client.unlock(); } return Optional.ofNullable(oldSession); } @Override - public void updateRegistration(LwM2mClient lwM2MClient, Registration registration) throws LwM2MClientStateException { - lwM2MClient.lock(); + public void updateRegistration(LwM2mClient client, Registration registration) throws LwM2MClientStateException { + client.lock(); try { - if (!LwM2MClientState.REGISTERED.equals(lwM2MClient.getState())) { - throw new LwM2MClientStateException(lwM2MClient.getState(), "Client is in invalid state."); + if (!LwM2MClientState.REGISTERED.equals(client.getState())) { + throw new LwM2MClientStateException(client.getState(), "Client is in invalid state."); } - lwM2MClient.setRegistration(registration); + client.setRegistration(registration); + clientStore.put(client); } finally { - lwM2MClient.unlock(); + client.unlock(); } } @Override - public void unregister(LwM2mClient lwM2MClient, Registration registration) throws LwM2MClientStateException { - lwM2MClient.lock(); + public void unregister(LwM2mClient client, Registration registration) throws LwM2MClientStateException { + client.lock(); try { - if (!LwM2MClientState.REGISTERED.equals(lwM2MClient.getState())) { - throw new LwM2MClientStateException(lwM2MClient.getState(), "Client is in invalid state."); + if (!LwM2MClientState.REGISTERED.equals(client.getState())) { + throw new LwM2MClientStateException(client.getState(), "Client is in invalid state."); } lwM2mClientsByRegistrationId.remove(registration.getId()); - Registration currentRegistration = lwM2MClient.getRegistration(); + Registration currentRegistration = client.getRegistration(); if (currentRegistration.getId().equals(registration.getId())) { - lwM2MClient.setState(LwM2MClientState.UNREGISTERED); - lwM2mClientsByEndpoint.remove(lwM2MClient.getEndpoint()); - this.securityStore.remove(lwM2MClient.getEndpoint(), registration.getId()); - UUID profileId = lwM2MClient.getProfileId(); + client.setState(LwM2MClientState.UNREGISTERED); + lwM2mClientsByEndpoint.remove(client.getEndpoint()); + this.securityStore.remove(client.getEndpoint(), registration.getId()); + clientStore.remove(client.getEndpoint()); + UUID profileId = client.getProfileId(); if (profileId != null) { Optional otherClients = lwM2mClientsByRegistrationId.values().stream().filter(e -> e.getProfileId().equals(profileId)).findFirst(); if (otherClients.isEmpty()) { @@ -137,24 +170,19 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { } } } else { - throw new LwM2MClientStateException(lwM2MClient.getState(), "Client has different registration."); + throw new LwM2MClientStateException(client.getState(), "Client has different registration."); } } finally { - lwM2MClient.unlock(); + client.unlock(); } } - @Override - public LwM2mClient getClientByRegistrationId(String registrationId) { - return lwM2mClientsByRegistrationId.get(registrationId); - } - @Override public LwM2mClient getClientBySessionInfo(TransportProtos.SessionInfoProto sessionInfo) { LwM2mClient lwM2mClient = null; + UUID sessionId = new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()); Predicate isClientFilter = c -> - (new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())) - .equals((new UUID(c.getSession().getSessionIdMSB(), c.getSession().getSessionIdLSB()))); + sessionId.equals((new UUID(c.getSession().getSessionIdMSB(), c.getSession().getSessionIdLSB()))); if (this.lwM2mClientsByEndpoint.size() > 0) { lwM2mClient = this.lwM2mClientsByEndpoint.values().stream().filter(isClientFilter).findAny().orElse(null); } @@ -162,31 +190,17 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { lwM2mClient = this.lwM2mClientsByRegistrationId.values().stream().filter(isClientFilter).findAny().orElse(null); } if (lwM2mClient == null) { - log.warn("Device TimeOut? lwM2mClient is null."); - log.warn("SessionInfo input [{}], lwM2mClientsByEndpoint size: [{}] lwM2mClientsByRegistrationId: [{}]", sessionInfo, lwM2mClientsByEndpoint.values(), lwM2mClientsByRegistrationId.values()); - log.error("", new RuntimeException()); + log.error("[{}] Failed to lookup client by session id.", sessionId); } return lwM2mClient; } - /** - * Get path to resource from profile equal keyName - * - * @param sessionInfo - - * @param keyName - - * @return - - */ - @Override - public String getObjectIdByKeyNameFromProfile(TransportProtos.SessionInfoProto sessionInfo, String keyName) { - return getObjectIdByKeyNameFromProfile(getClientBySessionInfo(sessionInfo), keyName); - } - @Override - public String getObjectIdByKeyNameFromProfile(LwM2mClient lwM2mClient, String keyName) { - Lwm2mDeviceProfileTransportConfiguration profile = getProfile(lwM2mClient.getProfileId()); + public String getObjectIdByKeyNameFromProfile(LwM2mClient client, String keyName) { + Lwm2mDeviceProfileTransportConfiguration profile = getProfile(client.getProfileId()); return profile.getObserveAttr().getKeyName().entrySet().stream() - .filter(e -> e.getValue().equals(keyName) && validateResourceInModel(lwM2mClient, e.getKey(), false)).findFirst().orElseThrow( + .filter(e -> e.getValue().equals(keyName) && validateResourceInModel(client, e.getKey(), false)).findFirst().orElseThrow( () -> new IllegalArgumentException(keyName + " is not configured in the device profile!") ).getKey(); } @@ -198,11 +212,21 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { @Override public void registerClient(Registration registration, ValidateDeviceCredentialsResponse credentials) { LwM2mClient client = getClientByEndpoint(registration.getEndpoint()); - client.init(null, null, credentials, UUID.randomUUID()); + client.init(credentials, UUID.randomUUID()); lwM2mClientsByRegistrationId.put(registration.getId(), client); profileUpdate(credentials.getDeviceProfile()); } + @Override + public void update(LwM2mClient client) { + client.lock(); + try { + clientStore.put(client); + } finally { + client.unlock(); + } + } + @Override public Collection getLwM2mClients() { return lwM2mClientsByEndpoint.values(); @@ -210,19 +234,33 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { @Override public Lwm2mDeviceProfileTransportConfiguration getProfile(UUID profileId) { - return profiles.get(profileId); + return doGetAndCache(profileId); } @Override public Lwm2mDeviceProfileTransportConfiguration getProfile(Registration registration) { - return profiles.get(getClientByEndpoint(registration.getEndpoint()).getProfileId()); + UUID profileId = getClientByEndpoint(registration.getEndpoint()).getProfileId(); + Lwm2mDeviceProfileTransportConfiguration result = doGetAndCache(profileId); + if (result == null) { + log.debug("[{}] Fetching profile [{}]", registration.getEndpoint(), profileId); + DeviceProfile deviceProfile = deviceProfileCache.get(new DeviceProfileId(profileId)); + if (deviceProfile != null) { + profileUpdate(deviceProfile); + result = doGetAndCache(profileId); + } + } + return result; + } + + private Lwm2mDeviceProfileTransportConfiguration doGetAndCache(UUID profileId) { + return profiles.get(profileId); } @Override public Lwm2mDeviceProfileTransportConfiguration profileUpdate(DeviceProfile deviceProfile) { - Lwm2mDeviceProfileTransportConfiguration lwM2MClientProfile = LwM2mTransportUtil.toLwM2MClientProfile(deviceProfile); - profiles.put(deviceProfile.getUuidId(), lwM2MClientProfile); - return lwM2MClientProfile; + Lwm2mDeviceProfileTransportConfiguration clientProfile = LwM2mTransportUtil.toLwM2MClientProfile(deviceProfile); + profiles.put(deviceProfile.getUuidId(), clientProfile); + return clientProfile; } @Override diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceValue.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceValue.java index cbaf60ca77..38a1815b06 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceValue.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceValue.java @@ -18,14 +18,46 @@ package org.thingsboard.server.transport.lwm2m.server.client; import lombok.Data; import org.eclipse.leshan.core.model.ResourceModel; import org.eclipse.leshan.core.node.LwM2mResource; +import org.eclipse.leshan.core.node.LwM2mResourceInstance; + +import java.io.Serializable; @Data -public class ResourceValue { - private LwM2mResource lwM2mResource; - private ResourceModel resourceModel; +public class ResourceValue implements Serializable { + + private static final long serialVersionUID = -228268906779089402L; + + private TbLwM2MResource lwM2mResource; + private TbResourceModel resourceModel; public ResourceValue(LwM2mResource lwM2mResource, ResourceModel resourceModel) { - this.lwM2mResource = lwM2mResource; - this.resourceModel = resourceModel; + this.lwM2mResource = toTbLwM2MResource(lwM2mResource); + this.resourceModel = toTbResourceModel(resourceModel); + } + + public void setLwM2mResource(LwM2mResource lwM2mResource) { + this.lwM2mResource = toTbLwM2MResource(lwM2mResource); + } + + public void setResourceModel(ResourceModel resourceModel) { + this.resourceModel = toTbResourceModel(resourceModel); + } + + private static TbLwM2MResource toTbLwM2MResource(LwM2mResource lwM2mResource) { + if (lwM2mResource.isMultiInstances()) { + TbLwM2MResourceInstance[] instances = (TbLwM2MResourceInstance[]) lwM2mResource.getInstances().values().stream().map(ResourceValue::toTbLwM2MResourceInstance).toArray(); + return new TbLwM2MMultipleResource(lwM2mResource.getId(), lwM2mResource.getType(), instances); + } else { + return new TbLwM2MSingleResource(lwM2mResource.getId(), lwM2mResource.getValue(), lwM2mResource.getType()); + } + } + + private static TbLwM2MResourceInstance toTbLwM2MResourceInstance(LwM2mResourceInstance instance) { + return new TbLwM2MResourceInstance(instance.getId(), instance.getValue(), instance.getType()); + } + + private static TbResourceModel toTbResourceModel(ResourceModel resourceModel) { + return new TbResourceModel(resourceModel.id, resourceModel.name, resourceModel.operations, resourceModel.multiple, + resourceModel.mandatory, resourceModel.type, resourceModel.rangeEnumeration, resourceModel.units, resourceModel.description); } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MMultipleResource.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MMultipleResource.java new file mode 100644 index 0000000000..ad5431b383 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MMultipleResource.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.client; + +import org.eclipse.leshan.core.model.ResourceModel; +import org.eclipse.leshan.core.node.LwM2mMultipleResource; +import org.eclipse.leshan.core.node.LwM2mResourceInstance; + +import java.io.Serializable; +import java.util.Collection; + +public class TbLwM2MMultipleResource extends LwM2mMultipleResource implements TbLwM2MResource, Serializable { + + private static final long serialVersionUID = 4658477128628087186L; + + public TbLwM2MMultipleResource(int id, ResourceModel.Type type, TbLwM2MResourceInstance... instances) { + super(id, type, instances); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MResource.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MResource.java new file mode 100644 index 0000000000..7f48a3e5c3 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MResource.java @@ -0,0 +1,21 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.client; + +import org.eclipse.leshan.core.node.LwM2mResource; + +public interface TbLwM2MResource extends LwM2mResource { +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MResourceInstance.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MResourceInstance.java new file mode 100644 index 0000000000..04cd91f4c2 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MResourceInstance.java @@ -0,0 +1,30 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.client; + +import org.eclipse.leshan.core.model.ResourceModel; +import org.eclipse.leshan.core.node.LwM2mResourceInstance; + +import java.io.Serializable; + +public class TbLwM2MResourceInstance extends LwM2mResourceInstance implements Serializable { + + private static final long serialVersionUID = -8322290426892538345L; + + protected TbLwM2MResourceInstance(int id, Object value, ResourceModel.Type type) { + super(id, value, type); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MSingleResource.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MSingleResource.java new file mode 100644 index 0000000000..6835481e36 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbLwM2MSingleResource.java @@ -0,0 +1,30 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.client; + +import org.eclipse.leshan.core.model.ResourceModel; +import org.eclipse.leshan.core.node.LwM2mSingleResource; + +import java.io.Serializable; + +public class TbLwM2MSingleResource extends LwM2mSingleResource implements TbLwM2MResource, Serializable { + + private static final long serialVersionUID = -878078368245340809L; + + public TbLwM2MSingleResource(int id, Object value, ResourceModel.Type type) { + super(id, value, type); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbResourceModel.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbResourceModel.java new file mode 100644 index 0000000000..6363f49cc8 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/TbResourceModel.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.client; + +import org.eclipse.leshan.core.model.ResourceModel; + +import java.io.Serializable; + +public class TbResourceModel extends ResourceModel implements Serializable { + + private static final long serialVersionUID = -2082846558899793932L; + + public TbResourceModel(Integer id, String name, Operations operations, Boolean multiple, Boolean mandatory, Type type, String rangeEnumeration, String units, String description) { + super(id, name, operations, multiple, mandatory, type, rangeEnumeration, units, description); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/DefaultLwM2MTelemetryLogService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/DefaultLwM2MTelemetryLogService.java index 421f4cf3a6..d84019f5e8 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/DefaultLwM2MTelemetryLogService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/DefaultLwM2MTelemetryLogService.java @@ -31,18 +31,8 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.L @RequiredArgsConstructor public class DefaultLwM2MTelemetryLogService implements LwM2MTelemetryLogService { - private final LwM2mClientContext clientContext; private final LwM2mTransportServerHelper helper; - /** - * @param logMsg - text msg - * @param registrationId - Id of Registration LwM2M Client - */ - @Override - public void log(String registrationId, String logMsg) { - log(clientContext.getClientByRegistrationId(registrationId), logMsg); - } - @Override public void log(LwM2mClient client, String logMsg) { if (logMsg != null && client != null && client.getSession() != null) { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/LwM2MTelemetryLogService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/LwM2MTelemetryLogService.java index ff8543303e..f2276fa9af 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/LwM2MTelemetryLogService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/log/LwM2MTelemetryLogService.java @@ -21,6 +21,4 @@ public interface LwM2MTelemetryLogService { void log(LwM2mClient client, String msg); - void log(String registrationId, String msg); - } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java index 374633dcef..c6c2f4ea63 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java @@ -52,6 +52,7 @@ import org.thingsboard.server.transport.lwm2m.server.ota.firmware.FirmwareUpdate import org.thingsboard.server.transport.lwm2m.server.ota.software.LwM2MSoftwareUpdateStrategy; import org.thingsboard.server.transport.lwm2m.server.ota.software.SoftwareUpdateResult; import org.thingsboard.server.transport.lwm2m.server.ota.software.SoftwareUpdateState; +import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MClientOtaInfoStore; import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import javax.annotation.PostConstruct; @@ -123,6 +124,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl private final OtaPackageDataCache otaPackageDataCache; private final LwM2MTelemetryLogService logService; private final LwM2mTransportServerHelper helper; + private final TbLwM2MClientOtaInfoStore otaInfoStore; @Autowired @Lazy @@ -174,6 +176,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl }, throwable -> { if (fwInfo.isSupported()) { fwInfo.setTargetFetchFailure(true); + update(fwInfo); } }, executor); } @@ -191,6 +194,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl public void onTargetFirmwareUpdate(LwM2mClient client, String newFirmwareTitle, String newFirmwareVersion, Optional newFirmwareUrl) { LwM2MClientOtaInfo fwInfo = getOrInitFwInfo(client); fwInfo.updateTarget(newFirmwareTitle, newFirmwareVersion, newFirmwareUrl); + update(fwInfo); startFirmwareUpdateIfNeeded(client, fwInfo); } @@ -202,7 +206,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl } @Override - public void onCurrentFirmwareStrategyUpdate(LwM2mClient client, OtherConfiguration configuration) { + public void onFirmwareStrategyUpdate(LwM2mClient client, OtherConfiguration configuration) { log.debug("[{}] Current fw strategy: {}", client.getEndpoint(), configuration.getFwUpdateStrategy()); LwM2MClientOtaInfo fwInfo = getOrInitFwInfo(client); fwInfo.setFwStrategy(LwM2MFirmwareUpdateStrategy.fromStrategyFwByCode(configuration.getFwUpdateStrategy())); @@ -242,9 +246,10 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl executeFwUpdate(client); } fwInfo.setUpdateState(state); - Optional status = this.toOtaPackageUpdateStatus(state); + Optional status = toOtaPackageUpdateStatus(state); status.ifPresent(otaStatus -> sendStateUpdateToTelemetry(client, fwInfo, otaStatus, "Firmware Update State: " + state.name())); + update(fwInfo); } @Override @@ -252,15 +257,16 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl log.debug("[{}] Current fw result: {}", client.getEndpoint(), code); LwM2MClientOtaInfo fwInfo = getOrInitFwInfo(client); FirmwareUpdateResult result = FirmwareUpdateResult.fromUpdateResultFwByCode(code.intValue()); - Optional status = this.toOtaPackageUpdateStatus(result); + Optional status = toOtaPackageUpdateStatus(result); status.ifPresent(otaStatus -> sendStateUpdateToTelemetry(client, fwInfo, otaStatus, "Firmware Update Result: " + result.name())); if (result.isAgain() && fwInfo.getRetryAttempts() <= 2) { fwInfo.setRetryAttempts(fwInfo.getRetryAttempts() + 1); startFirmwareUpdateIfNeeded(client, fwInfo); } else { - fwInfo.setUpdateResult(result); + fwInfo.update(result); } + update(fwInfo); } @Override @@ -378,23 +384,38 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl } public LwM2MClientOtaInfo getOrInitFwInfo(LwM2mClient client) { - //TODO: fetch state from the cache or DB. return this.fwStates.computeIfAbsent(client.getEndpoint(), endpoint -> { - var profile = clientContext.getProfile(client.getProfileId()); - return new LwM2MClientOtaInfo(endpoint, OtaPackageType.FIRMWARE, profile.getClientLwM2mSettings().getFwUpdateStrategy(), - profile.getClientLwM2mSettings().getFwUpdateResource()); + LwM2MClientOtaInfo info = otaInfoStore.get(OtaPackageType.FIRMWARE, endpoint); + if (info == null) { + var profile = clientContext.getProfile(client.getProfileId()); + info = new LwM2MClientOtaInfo(endpoint, OtaPackageType.FIRMWARE, + LwM2MFirmwareUpdateStrategy.fromStrategyFwByCode(profile.getClientLwM2mSettings().getFwUpdateStrategy()), + profile.getClientLwM2mSettings().getFwUpdateResource()); + update(info); + } + return info; }); } private LwM2MClientOtaInfo getOrInitSwInfo(LwM2mClient client) { - //TODO: fetch state from the cache or DB. - return swStates.computeIfAbsent(client.getEndpoint(), endpoint -> { - var profile = clientContext.getProfile(client.getProfileId()); - return new LwM2MClientOtaInfo(endpoint, OtaPackageType.SOFTWARE, profile.getClientLwM2mSettings().getSwUpdateStrategy(), profile.getClientLwM2mSettings().getSwUpdateResource()); + return this.fwStates.computeIfAbsent(client.getEndpoint(), endpoint -> { + LwM2MClientOtaInfo info = otaInfoStore.get(OtaPackageType.SOFTWARE, endpoint); + if (info == null) { + var profile = clientContext.getProfile(client.getProfileId()); + info = new LwM2MClientOtaInfo(endpoint, OtaPackageType.SOFTWARE, + LwM2MSoftwareUpdateStrategy.fromStrategySwByCode(profile.getClientLwM2mSettings().getFwUpdateStrategy()), + profile.getClientLwM2mSettings().getSwUpdateResource()); + update(info); + } + return info; }); } + private void update(LwM2MClientOtaInfo info) { + otaInfoStore.put(info); + } + private void sendStateUpdateToTelemetry(LwM2mClient client, LwM2MClientOtaInfo fwInfo, OtaPackageUpdateStatus status, String log) { List result = new ArrayList<>(); TransportProtos.KeyValueProto.Builder kvProto = TransportProtos.KeyValueProto.newBuilder().setKey(getAttributeKey(fwInfo.getType(), STATE)); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java index 33b8bdfe81..260c266906 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaInfo.java @@ -15,7 +15,9 @@ */ package org.thingsboard.server.transport.lwm2m.server.ota; +import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Data; +import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.transport.lwm2m.server.ota.firmware.LwM2MFirmwareUpdateStrategy; @@ -26,10 +28,11 @@ import org.thingsboard.server.transport.lwm2m.server.ota.software.LwM2MSoftwareU import java.util.Optional; @Data +@NoArgsConstructor public class LwM2MClientOtaInfo { - private final String endpoint; - private final OtaPackageType type; + private String endpoint; + private OtaPackageType type; private String baseUrl; @@ -53,10 +56,17 @@ public class LwM2MClientOtaInfo { private String failedPackageId; private int retryAttempts; - public LwM2MClientOtaInfo(String endpoint, OtaPackageType type, Integer strategyCode, String baseUrl) { + public LwM2MClientOtaInfo(String endpoint, OtaPackageType type, LwM2MFirmwareUpdateStrategy fwStrategy, String baseUrl) { this.endpoint = endpoint; this.type = type; - this.fwStrategy = strategyCode != null ? LwM2MFirmwareUpdateStrategy.fromStrategyFwByCode(strategyCode) : LwM2MFirmwareUpdateStrategy.OBJ_5_BINARY; + this.fwStrategy = fwStrategy; + this.baseUrl = baseUrl; + } + + public LwM2MClientOtaInfo(String endpoint, OtaPackageType type, LwM2MSoftwareUpdateStrategy swStrategy, String baseUrl) { + this.endpoint = endpoint; + this.type = type; + this.swStrategy = swStrategy; this.baseUrl = baseUrl; } @@ -66,6 +76,7 @@ public class LwM2MClientOtaInfo { this.targetUrl = newFirmwareUrl.orElse(null); } + @JsonIgnore public boolean isUpdateRequired() { if (StringUtils.isEmpty(targetName) || StringUtils.isEmpty(targetVersion) || !isSupported()) { return false; @@ -86,11 +97,12 @@ public class LwM2MClientOtaInfo { } } + @JsonIgnore public boolean isSupported() { return StringUtils.isNotEmpty(currentName) || StringUtils.isNotEmpty(currentVersion5) || StringUtils.isNotEmpty(currentVersion3); } - public void setUpdateResult(FirmwareUpdateResult updateResult) { + public void update(FirmwareUpdateResult updateResult) { this.updateResult = updateResult; switch (updateResult) { case INITIAL: diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java index 3ba1b1500a..2aa7f39872 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java @@ -32,7 +32,7 @@ public interface LwM2MOtaUpdateService { void onCurrentFirmwareNameUpdate(LwM2mClient client, String name); - void onCurrentFirmwareStrategyUpdate(LwM2mClient client, OtherConfiguration configuration); + void onFirmwareStrategyUpdate(LwM2mClient client, OtherConfiguration configuration); void onCurrentSoftwareStrategyUpdate(LwM2mClient client, OtherConfiguration configuration); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/session/DefaultLwM2MSessionManager.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/session/DefaultLwM2MSessionManager.java new file mode 100644 index 0000000000..ee6e4a779b --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/session/DefaultLwM2MSessionManager.java @@ -0,0 +1,67 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.session; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Lazy; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.common.transport.service.DefaultTransportService; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; +import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener; +import org.thingsboard.server.transport.lwm2m.server.attributes.LwM2MAttributesService; +import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; + +@Slf4j +@Service +@TbLwM2mTransportComponent +public class DefaultLwM2MSessionManager implements LwM2MSessionManager { + + private final TransportService transportService; + private final LwM2MAttributesService attributesService; + private final LwM2MRpcRequestHandler rpcHandler; + private final LwM2mUplinkMsgHandler uplinkHandler; + + public DefaultLwM2MSessionManager(TransportService transportService, + @Lazy LwM2MAttributesService attributesService, + @Lazy LwM2MRpcRequestHandler rpcHandler, + @Lazy LwM2mUplinkMsgHandler uplinkHandler) { + this.transportService = transportService; + this.attributesService = attributesService; + this.rpcHandler = rpcHandler; + this.uplinkHandler = uplinkHandler; + } + + @Override + public void register(TransportProtos.SessionInfoProto sessionInfo) { + transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(uplinkHandler, attributesService, rpcHandler, sessionInfo, transportService)); + TransportProtos.TransportToDeviceActorMsg msg = TransportProtos.TransportToDeviceActorMsg.newBuilder() + .setSessionInfo(sessionInfo) + .setSessionEvent(DefaultTransportService.getSessionEventMsg(TransportProtos.SessionEvent.OPEN)) + .setSubscribeToAttributes(TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setSessionType(TransportProtos.SessionType.ASYNC).build()) + .setSubscribeToRPC(TransportProtos.SubscribeToRPCMsg.newBuilder().setSessionType(TransportProtos.SessionType.ASYNC).build()) + .build(); + transportService.process(msg, null); + } + + @Override + public void deregister(TransportProtos.SessionInfoProto sessionInfo) { + transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null); + transportService.deregisterSession(sessionInfo); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/session/LwM2MSessionManager.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/session/LwM2MSessionManager.java new file mode 100644 index 0000000000..eb85271492 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/session/LwM2MSessionManager.java @@ -0,0 +1,27 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.session; + +import org.thingsboard.server.gen.transport.TransportProtos; + +public interface LwM2MSessionManager { + + void register(TransportProtos.SessionInfoProto sessionInfo); + + void deregister(TransportProtos.SessionInfoProto sessionInfo); + + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbDummyLwM2MClientOtaInfoStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbDummyLwM2MClientOtaInfoStore.java new file mode 100644 index 0000000000..ee46adcdf8 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbDummyLwM2MClientOtaInfoStore.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.store; + +import org.thingsboard.server.common.data.ota.OtaPackageType; +import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MClientOtaInfo; + +public class TbDummyLwM2MClientOtaInfoStore implements TbLwM2MClientOtaInfoStore { + + @Override + public LwM2MClientOtaInfo get(OtaPackageType type, String endpoint) { + return null; + } + + @Override + public void put(LwM2MClientOtaInfo info) { + + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbDummyLwM2MClientStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbDummyLwM2MClientStore.java new file mode 100644 index 0000000000..3708bf4b0e --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbDummyLwM2MClientStore.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.store; + +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +public class TbDummyLwM2MClientStore implements TbLwM2MClientStore { + @Override + public LwM2mClient get(String endpoint) { + return null; + } + + @Override + public void put(LwM2mClient client) { + + } + + @Override + public void remove(String endpoint) { + + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MClientOtaInfoStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MClientOtaInfoStore.java new file mode 100644 index 0000000000..b80321ccfb --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MClientOtaInfoStore.java @@ -0,0 +1,26 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.store; + +import org.thingsboard.server.common.data.ota.OtaPackageType; +import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MClientOtaInfo; + +public interface TbLwM2MClientOtaInfoStore { + + LwM2MClientOtaInfo get(OtaPackageType type, String endpoint); + + void put(LwM2MClientOtaInfo info); +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MClientStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MClientStore.java new file mode 100644 index 0000000000..16e4e9e90a --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MClientStore.java @@ -0,0 +1,27 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.store; + +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +public interface TbLwM2MClientStore { + + LwM2mClient get(String endpoint); + + void put(LwM2mClient client); + + void remove(String endpoint); +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisClientOtaInfoStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisClientOtaInfoStore.java new file mode 100644 index 0000000000..36283ef3ff --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisClientOtaInfoStore.java @@ -0,0 +1,54 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.store; + +import org.eclipse.leshan.server.security.NonUniqueSecurityInfoException; +import org.eclipse.leshan.server.security.SecurityInfo; +import org.nustaq.serialization.FSTConfiguration; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.integration.redis.util.RedisLockRegistry; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.ota.OtaPackageType; +import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MSecurityInfo; +import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MClientOtaInfo; + +import java.util.concurrent.locks.Lock; + +public class TbLwM2mRedisClientOtaInfoStore implements TbLwM2MClientOtaInfoStore { + private static final String OTA_EP = "OTA#EP#"; + + private final RedisConnectionFactory connectionFactory; + + public TbLwM2mRedisClientOtaInfoStore(RedisConnectionFactory connectionFactory) { + this.connectionFactory = connectionFactory; + } + + @Override + public LwM2MClientOtaInfo get(OtaPackageType type, String endpoint) { + try (var connection = connectionFactory.getConnection()) { + byte[] data = connection.get((OTA_EP + type + endpoint).getBytes()); + return JacksonUtil.fromBytes(data, LwM2MClientOtaInfo.class); + } + } + + @Override + public void put(LwM2MClientOtaInfo info) { + try (var connection = connectionFactory.getConnection()) { + connection.set((OTA_EP + info.getType() + info.getEndpoint()).getBytes(), JacksonUtil.toString(info).getBytes()); + } + } + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java index b9eb865df5..9156d73181 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java @@ -20,6 +20,7 @@ import org.eclipse.leshan.server.californium.registration.InMemoryRegistrationSt import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; +import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.stereotype.Component; import org.thingsboard.server.cache.TBRedisCacheConfiguration; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; @@ -46,20 +47,37 @@ public class TbLwM2mStoreFactory { @Bean private CaliforniumRegistrationStore registrationStore() { - return redisConfiguration.isPresent() && useRedis ? - new TbLwM2mRedisRegistrationStore(redisConfiguration.get().redisConnectionFactory()) : new InMemoryRegistrationStore(config.getCleanPeriodInSec()); + return isRedis() ? + new TbLwM2mRedisRegistrationStore(getConnectionFactory()) : new InMemoryRegistrationStore(config.getCleanPeriodInSec()); } @Bean private TbMainSecurityStore securityStore() { - return new TbLwM2mSecurityStore(redisConfiguration.isPresent() && useRedis ? - new TbLwM2mRedisSecurityStore(redisConfiguration.get().redisConnectionFactory()) : new TbInMemorySecurityStore(), validator); + return new TbLwM2mSecurityStore(isRedis() ? + new TbLwM2mRedisSecurityStore(getConnectionFactory()) : new TbInMemorySecurityStore(), validator); + } + + @Bean + private TbLwM2MClientStore clientStore() { + return isRedis() ? new TbRedisLwM2MClientStore(getConnectionFactory()) : new TbDummyLwM2MClientStore(); + } + + @Bean + private TbLwM2MClientOtaInfoStore otaStore() { + return isRedis() ? new TbLwM2mRedisClientOtaInfoStore(getConnectionFactory()) : new TbDummyLwM2MClientOtaInfoStore(); } @Bean private TbLwM2MDtlsSessionStore sessionStore() { - return redisConfiguration.isPresent() && useRedis ? - new TbLwM2MDtlsSessionRedisStore(redisConfiguration.get().redisConnectionFactory()) : new TbL2M2MDtlsSessionInMemoryStore(); + return isRedis() ? new TbLwM2MDtlsSessionRedisStore(getConnectionFactory()) : new TbL2M2MDtlsSessionInMemoryStore(); + } + + private RedisConnectionFactory getConnectionFactory() { + return redisConfiguration.get().redisConnectionFactory(); + } + + private boolean isRedis() { + return redisConfiguration.isPresent() && useRedis; } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java new file mode 100644 index 0000000000..a408cc22c4 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbRedisLwM2MClientStore.java @@ -0,0 +1,63 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.store; + +import org.nustaq.serialization.FSTConfiguration; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +public class TbRedisLwM2MClientStore implements TbLwM2MClientStore { + + private static final String CLIENT_EP = "CLIENT#EP#"; + private final RedisConnectionFactory connectionFactory; + private final FSTConfiguration serializer; + + public TbRedisLwM2MClientStore(RedisConnectionFactory redisConnectionFactory) { + this.connectionFactory = redisConnectionFactory; + this.serializer = FSTConfiguration.createDefaultConfiguration(); + } + + @Override + public LwM2mClient get(String endpoint) { + try (var connection = connectionFactory.getConnection()) { + byte[] data = connection.get(getKey(endpoint)); + if (data == null) { + return null; + } else { + return (LwM2mClient) serializer.asObject(data); + } + } + } + + @Override + public void put(LwM2mClient client) { + byte[] clientSerialized = serializer.asByteArray(client); + try (var connection = connectionFactory.getConnection()) { + connection.getSet(getKey(client.getEndpoint()), clientSerialized); + } + } + + @Override + public void remove(String endpoint) { + try (var connection = connectionFactory.getConnection()) { + connection.del(getKey(endpoint)); + } + } + + private byte[] getKey(String endpoint) { + return (CLIENT_EP + endpoint).getBytes(); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java index 56dad43db7..5f8ccad17e 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java @@ -48,15 +48,11 @@ import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTrans import org.thingsboard.server.common.data.ota.OtaPackageUtil; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; -import org.thingsboard.server.common.transport.service.DefaultTransportService; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent; import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.server.LwM2mOtaConvert; -import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest; -import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; @@ -84,6 +80,7 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttrib import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MOtaUpdateService; import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler; +import org.thingsboard.server.transport.lwm2m.server.session.LwM2MSessionManager; import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; @@ -129,6 +126,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl private final TransportService transportService; private final LwM2mTransportContext context; private final LwM2MAttributesService attributesService; + private final LwM2MSessionManager sessionManager; private final LwM2MOtaUpdateService otaService; private final LwM2MTransportServerConfig config; private final LwM2MTelemetryLogService logService; @@ -143,12 +141,14 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl LwM2mTransportServerHelper helper, LwM2mClientContext clientContext, LwM2MTelemetryLogService logService, + LwM2MSessionManager sessionManager, @Lazy LwM2MOtaUpdateService otaService, @Lazy LwM2MAttributesService attributesService, @Lazy LwM2MRpcRequestHandler rpcHandler, @Lazy LwM2mDownlinkMsgHandler defaultLwM2MDownlinkMsgHandler, LwM2mTransportContext context, TbLwM2MDtlsSessionStore sessionStore) { this.transportService = transportService; + this.sessionManager = sessionManager; this.attributesService = attributesService; this.otaService = otaService; this.config = config; @@ -205,18 +205,10 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl Optional oldSessionInfo = this.clientContext.register(lwM2MClient, registration); if (oldSessionInfo.isPresent()) { log.info("[{}] Closing old session: {}", registration.getEndpoint(), new UUID(oldSessionInfo.get().getSessionIdMSB(), oldSessionInfo.get().getSessionIdLSB())); - closeSession(oldSessionInfo.get()); + sessionManager.deregister(oldSessionInfo.get()); } logService.log(lwM2MClient, LOG_LWM2M_INFO + ": Client registered with registration id: " + registration.getId()); - SessionInfoProto sessionInfo = lwM2MClient.getSession(); - transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, attributesService, rpcHandler, sessionInfo, transportService)); - TransportProtos.TransportToDeviceActorMsg msg = TransportProtos.TransportToDeviceActorMsg.newBuilder() - .setSessionInfo(sessionInfo) - .setSessionEvent(DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN)) - .setSubscribeToAttributes(TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setSessionType(TransportProtos.SessionType.ASYNC).build()) - .setSubscribeToRPC(TransportProtos.SubscribeToRPCMsg.newBuilder().setSessionType(TransportProtos.SessionType.ASYNC).build()) - .build(); - transportService.process(msg, null); + sessionManager.register(lwM2MClient.getSession()); this.initClientTelemetry(lwM2MClient); this.initAttributes(lwM2MClient); otaService.init(lwM2MClient); @@ -247,14 +239,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl log.warn("[{}] [{{}] Client: update after Registration", registration.getEndpoint(), registration.getId()); logService.log(lwM2MClient, String.format("[%s][%s] Updated registration.", registration.getId(), registration.getSocketAddress())); clientContext.updateRegistration(lwM2MClient, registration); - TransportProtos.SessionInfoProto sessionInfo = lwM2MClient.getSession(); - this.reportActivityAndRegister(sessionInfo); - if (registration.usesQueueMode()) { - LwM2mQueuedRequest request; - while ((request = lwM2MClient.getQueuedRequests().poll()) != null) { - request.send(); - } - } + this.reportActivityAndRegister(lwM2MClient.getSession()); } catch (LwM2MClientStateException stateException) { if (LwM2MClientState.REGISTERED.equals(stateException.getState())) { log.info("[{}] update registration failed because client has different registration id: [{}] {}.", registration.getEndpoint(), stateException.getState(), stateException.getMessage()); @@ -280,7 +265,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl clientContext.unregister(client, registration); SessionInfoProto sessionInfo = client.getSession(); if (sessionInfo != null) { - closeSession(sessionInfo); + sessionManager.deregister(sessionInfo); sessionStore.remove(registration.getEndpoint()); log.info("Client close session: [{}] unReg [{}] name [{}] profile ", registration.getId(), registration.getEndpoint(), sessionInfo.getDeviceType()); } else { @@ -295,11 +280,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl }); } - public void closeSession(SessionInfoProto sessionInfo) { - transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.CLOSED), null); - transportService.deregisterSession(sessionInfo); - } - @Override public void onSleepingDev(Registration registration) { log.info("[{}] [{}] Received endpoint Sleeping version event", registration.getId(), registration.getEndpoint()); @@ -307,19 +287,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl //TODO: associate endpointId with device information. } -// /** -// * Cancel observation for All objects for this registration -// */ -// @Override -// public void setCancelObservationsAll(Registration registration) { -// if (registration != null) { -// LwM2mClient client = clientContext.getClientByEndpoint(registration.getEndpoint()); -// if (client != null && client.getRegistration() != null && client.getRegistration().getId().equals(registration.getId())) { -// defaultLwM2MDownlinkMsgHandler.sendCancelAllRequest(client, TbLwM2MCancelAllRequest.builder().build(), new TbLwM2MCancelAllObserveRequestCallback(this, client)); -// } -// } -// } - /** * Sending observe value to thingsboard from ObservationListener.onResponse: object, instance, SingleResource or MultipleResource * @@ -344,6 +311,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl this.updateResourcesValue(lwM2MClient, lwM2mResource, path); } } + clientContext.update(lwM2MClient); } } @@ -390,16 +358,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl clientContext.getLwM2mClients().forEach(e -> e.deleteResources(pathIdVer, this.config.getModelProvider())); } - /** - * Deregister session in transport - * - * @param sessionInfo - lwm2m client - */ - @Override - public void doDisconnect(SessionInfoProto sessionInfo) { - closeSession(sessionInfo); - } - /** * Those methods are called by the protocol stage thread pool, this means that execution MUST be done in a short delay, * * if you need to do long time processing use a dedicated thread pool. @@ -494,14 +452,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl attributesMap.forEach((targetId, params) -> sendWriteAttributesRequest(lwM2MClient, targetId, params)); } - private void sendDiscoverRequests(LwM2mClient lwM2MClient, Lwm2mDeviceProfileTransportConfiguration profile, Set supportedObjects) { - Set targetIds = profile.getObserveAttr().getAttributeLwm2m().keySet(); - targetIds = targetIds.stream().filter(target -> isSupportedTargetId(supportedObjects, target)).collect(Collectors.toSet()); -// TODO: why do we need to put observe into pending read requests? -// lwM2MClient.getPendingReadRequests().addAll(targetIds); - targetIds.forEach(targetId -> sendDiscoverRequest(lwM2MClient, targetId)); - } - private void sendDiscoverRequest(LwM2mClient lwM2MClient, String targetId) { TbLwM2MDiscoverRequest request = TbLwM2MDiscoverRequest.builder().versionedId(targetId).timeout(this.config.getTimeout()).build(); defaultLwM2MDownlinkMsgHandler.sendDiscoverRequest(lwM2MClient, request, new TbLwM2MDiscoverCallback(logService, lwM2MClient, targetId)); @@ -652,7 +602,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl List resultAttributes = new ArrayList<>(); profile.getObserveAttr().getAttribute().forEach(pathIdVer -> { if (path.contains(pathIdVer)) { - TransportProtos.KeyValueProto kvAttr = this.getKvToThingsboard(pathIdVer, registration); + TransportProtos.KeyValueProto kvAttr = this.getKvToThingsBoard(pathIdVer, registration); if (kvAttr != null) { resultAttributes.add(kvAttr); } @@ -661,7 +611,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl List resultTelemetries = new ArrayList<>(); profile.getObserveAttr().getTelemetry().forEach(pathIdVer -> { if (path.contains(pathIdVer)) { - TransportProtos.KeyValueProto kvAttr = this.getKvToThingsboard(pathIdVer, registration); + TransportProtos.KeyValueProto kvAttr = this.getKvToThingsBoard(pathIdVer, registration); if (kvAttr != null) { resultTelemetries.add(kvAttr); } @@ -678,7 +628,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl return null; } - private TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) { + private TransportProtos.KeyValueProto getKvToThingsBoard(String pathIdVer, Registration registration) { LwM2mClient lwM2MClient = this.clientContext.getClientByEndpoint(registration.getEndpoint()); Map names = clientContext.getProfile(lwM2MClient.getProfileId()).getObserveAttr().getKeyName(); if (names != null && names.containsKey(pathIdVer)) { @@ -725,10 +675,12 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl public void onWriteResponseOk(LwM2mClient client, String path, WriteRequest request) { if (request.getNode() instanceof LwM2mResource) { this.updateResourcesValue(client, ((LwM2mResource) request.getNode()), path); + clientContext.update(client); } else if (request.getNode() instanceof LwM2mObjectInstance) { ((LwM2mObjectInstance) request.getNode()).getResources().forEach((resId, resource) -> { this.updateResourcesValue(client, resource, path + "/" + resId); }); + clientContext.update(client); } } @@ -803,7 +755,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl if (!newLwM2mSettings.getFwUpdateStrategy().equals(oldLwM2mSettings.getFwUpdateStrategy()) || (StringUtils.isNotEmpty(newLwM2mSettings.getFwUpdateResource()) && !newLwM2mSettings.getFwUpdateResource().equals(oldLwM2mSettings.getFwUpdateResource()))) { - clients.forEach(lwM2MClient -> otaService.onCurrentFirmwareStrategyUpdate(lwM2MClient, newLwM2mSettings)); + clients.forEach(lwM2MClient -> otaService.onFirmwareStrategyUpdate(lwM2MClient, newLwM2mSettings)); } if (!newLwM2mSettings.getSwUpdateStrategy().equals(oldLwM2mSettings.getSwUpdateStrategy()) @@ -908,7 +860,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl */ private void reportActivityAndRegister(SessionInfoProto sessionInfo) { if (sessionInfo != null && transportService.reportActivity(sessionInfo) == null) { - transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, attributesService, rpcHandler, sessionInfo, transportService)); + sessionManager.register(sessionInfo); this.reportActivitySubscription(sessionInfo); } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java index 372daa7052..dd693450ec 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java @@ -48,8 +48,6 @@ public interface LwM2mUplinkMsgHandler { void onResourceDelete(Optional resourceDeleteMsgOpt); - void doDisconnect(TransportProtos.SessionInfoProto sessionInfo); - void onAwakeDev(Registration registration); void onWriteResponseOk(LwM2mClient client, String path, WriteRequest request); diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index aa250a7039..18b7abb67e 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -47,7 +47,7 @@ public class JacksonUtil { throw new IllegalArgumentException("The given object value: " + fromValue + " cannot be converted to " + toValueTypeRef, e); } - } + } public static T fromString(String string, Class clazz) { try { @@ -67,6 +67,15 @@ public class JacksonUtil { } } + public static T fromBytes(byte[] bytes, Class clazz) { + try { + return bytes != null ? OBJECT_MAPPER.readValue(bytes, clazz) : null; + } catch (IOException e) { + throw new IllegalArgumentException("The given string value: " + + Arrays.toString(bytes) + " cannot be transformed to Json object", e); + } + } + public static JsonNode fromBytes(byte[] bytes) { try { return OBJECT_MAPPER.readTree(bytes); @@ -96,7 +105,7 @@ public class JacksonUtil { } } - public static ObjectNode newObjectNode(){ + public static ObjectNode newObjectNode() { return OBJECT_MAPPER.createObjectNode(); }