|
|
|
@ -17,6 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server; |
|
|
|
|
|
|
|
import com.fasterxml.jackson.core.type.TypeReference; |
|
|
|
import com.fasterxml.jackson.databind.ObjectMapper; |
|
|
|
import com.google.common.collect.Sets; |
|
|
|
import com.google.gson.Gson; |
|
|
|
import com.google.gson.JsonArray; |
|
|
|
import com.google.gson.JsonElement; |
|
|
|
@ -35,7 +36,7 @@ import org.eclipse.leshan.core.response.ReadResponse; |
|
|
|
import org.eclipse.leshan.core.util.NamedThreadFactory; |
|
|
|
import org.eclipse.leshan.server.californium.LeshanServer; |
|
|
|
import org.eclipse.leshan.server.registration.Registration; |
|
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
|
import org.springframework.context.annotation.Lazy; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.thingsboard.server.common.data.Device; |
|
|
|
import org.thingsboard.server.common.data.DeviceProfile; |
|
|
|
@ -105,29 +106,32 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
protected final ReadWriteLock readWriteLock = new ReentrantReadWriteLock(); |
|
|
|
protected final Lock writeLock = readWriteLock.writeLock(); |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private TransportService transportService; |
|
|
|
private final TransportService transportService; |
|
|
|
|
|
|
|
@Autowired |
|
|
|
public LwM2mTransportContextServer context; |
|
|
|
public final LwM2mTransportContextServer lwM2mTransportContextServer; |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private LwM2mTransportRequest lwM2mTransportRequest; |
|
|
|
private final LwM2mClientContext lwM2mClientContext; |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private LwM2mClientContext lwM2mClientContext; |
|
|
|
private final LeshanServer leshanServer; |
|
|
|
|
|
|
|
@Autowired(required = false) |
|
|
|
private LeshanServer leshanServer; |
|
|
|
private final LwM2mTransportRequest lwM2mTransportRequest; |
|
|
|
|
|
|
|
public LwM2mTransportServiceImpl(TransportService transportService, LwM2mTransportContextServer lwM2mTransportContextServer, LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer, @Lazy LwM2mTransportRequest lwM2mTransportRequest) { |
|
|
|
this.transportService = transportService; |
|
|
|
this.lwM2mTransportContextServer = lwM2mTransportContextServer; |
|
|
|
this.lwM2mClientContext = lwM2mClientContext; |
|
|
|
this.leshanServer = leshanServer; |
|
|
|
this.lwM2mTransportRequest = lwM2mTransportRequest; |
|
|
|
} |
|
|
|
|
|
|
|
@PostConstruct |
|
|
|
public void init() { |
|
|
|
this.context.getScheduler().scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) context.getLwM2MTransportConfigServer().getSessionReportTimeout()), context.getLwM2MTransportConfigServer().getSessionReportTimeout(), TimeUnit.MILLISECONDS); |
|
|
|
this.executorRegistered = Executors.newFixedThreadPool(this.context.getLwM2MTransportConfigServer().getRegisteredPoolSize(), |
|
|
|
this.lwM2mTransportContextServer.getScheduler().scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) lwM2mTransportContextServer.getLwM2MTransportConfigServer().getSessionReportTimeout()), lwM2mTransportContextServer.getLwM2MTransportConfigServer().getSessionReportTimeout(), TimeUnit.MILLISECONDS); |
|
|
|
this.executorRegistered = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getRegisteredPoolSize(), |
|
|
|
new NamedThreadFactory(String.format("LwM2M %s channel registered", SERVICE_CHANNEL))); |
|
|
|
this.executorUpdateRegistered = Executors.newFixedThreadPool(this.context.getLwM2MTransportConfigServer().getUpdateRegisteredPoolSize(), |
|
|
|
this.executorUpdateRegistered = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getUpdateRegisteredPoolSize(), |
|
|
|
new NamedThreadFactory(String.format("LwM2M %s channel update registered", SERVICE_CHANNEL))); |
|
|
|
this.executorUnRegistered = Executors.newFixedThreadPool(this.context.getLwM2MTransportConfigServer().getUnRegisteredPoolSize(), |
|
|
|
this.executorUnRegistered = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getUnRegisteredPoolSize(), |
|
|
|
new NamedThreadFactory(String.format("LwM2M %s channel un registered", SERVICE_CHANNEL))); |
|
|
|
this.converter = LwM2mValueConverterImpl.getInstance(); |
|
|
|
} |
|
|
|
@ -143,17 +147,15 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* 1.2 Remove from sessions Model by enpPoint |
|
|
|
* Next -> Create new LwM2MClient for current session -> setModelClient... |
|
|
|
* |
|
|
|
* @param lwServer - LeshanServer |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
* @param previousObsersations - may be null |
|
|
|
*/ |
|
|
|
public void onRegistered(LeshanServer lwServer, Registration registration, Collection<Observation> previousObsersations) { |
|
|
|
public void onRegistered(Registration registration, Collection<Observation> previousObsersations) { |
|
|
|
executorRegistered.submit(() -> { |
|
|
|
try { |
|
|
|
log.warn("[{}] [{{}] Client: create after Registration", registration.getEndpoint(), registration.getId()); |
|
|
|
LwM2mClient lwM2MClient = this.lwM2mClientContext.updateInSessionsLwM2MClient(registration); |
|
|
|
if (lwM2MClient != null) { |
|
|
|
this.sentLogsToThingsboard(LOG_LW2M_INFO + ": Client Registered", registration); |
|
|
|
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registration); |
|
|
|
if (sessionInfo != null) { |
|
|
|
lwM2MClient.setDeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); |
|
|
|
@ -163,9 +165,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, sessionInfo)); |
|
|
|
transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.OPEN), null); |
|
|
|
transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null); |
|
|
|
this.sentLogsToThingsboard(LOG_LW2M_INFO + ": Client create after Registration", registration); |
|
|
|
this.initLwM2mFromClientValue(lwServer, registration, lwM2MClient); |
|
|
|
|
|
|
|
this.sentLogsToThingsboard(LOG_LW2M_INFO + ": Client create after Registration", registration); |
|
|
|
this.initLwM2mFromClientValue(registration, lwM2MClient); |
|
|
|
} else { |
|
|
|
log.error("Client: [{}] onRegistered [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), null); |
|
|
|
} |
|
|
|
@ -181,10 +182,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
/** |
|
|
|
* if sessionInfo removed from sessions, then new registerAsyncSession |
|
|
|
* |
|
|
|
* @param lwServer - LeshanServer |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
*/ |
|
|
|
public void updatedReg(LeshanServer lwServer, Registration registration) { |
|
|
|
public void updatedReg(Registration registration) { |
|
|
|
executorUpdateRegistered.submit(() -> { |
|
|
|
try { |
|
|
|
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registration); |
|
|
|
@ -205,10 +205,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* @param observations - All paths observations before unReg |
|
|
|
* !!! Warn: if have not finishing unReg, then this operation will be finished on next Client`s connect |
|
|
|
*/ |
|
|
|
public void unReg(LeshanServer lwServer, Registration registration, Collection<Observation> observations) { |
|
|
|
public void unReg(Registration registration, Collection<Observation> observations) { |
|
|
|
executorUnRegistered.submit(() -> { |
|
|
|
try { |
|
|
|
this.setCancelObservations(lwServer, registration); |
|
|
|
this.setCancelObservations(registration); |
|
|
|
this.sentLogsToThingsboard(LOG_LW2M_INFO + ": Client unRegistration", registration); |
|
|
|
this.closeClientSession(registration); |
|
|
|
} catch (Throwable t) { |
|
|
|
@ -238,10 +238,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void setCancelObservations(LeshanServer lwServer, Registration registration) { |
|
|
|
public void setCancelObservations(Registration registration) { |
|
|
|
if (registration != null) { |
|
|
|
Set<Observation> observations = lwServer.getObservationService().getObservations(registration); |
|
|
|
observations.forEach(observation -> this.setCancelObservationRecourse(lwServer, registration, observation.getPath().toString())); |
|
|
|
Set<Observation> observations = leshanServer.getObservationService().getObservations(registration); |
|
|
|
observations.forEach(observation -> this.setCancelObservationRecourse(registration, observation.getPath().toString())); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -251,8 +251,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* {@code ObservationService#cancelObservation()} |
|
|
|
*/ |
|
|
|
@Override |
|
|
|
public void setCancelObservationRecourse(LeshanServer lwServer, Registration registration, String path) { |
|
|
|
lwServer.getObservationService().cancelObservations(registration, path); |
|
|
|
public void setCancelObservationRecourse(Registration registration, String path) { |
|
|
|
leshanServer.getObservationService().cancelObservations(registration, path); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
@ -294,12 +294,12 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
String path = this.getPathAttributeUpdate(sessionInfo, de.getKey()); |
|
|
|
String value = de.getValue().getAsString(); |
|
|
|
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClient(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); |
|
|
|
LwM2mClientProfile profile = lwM2mClientContext.getProfile(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB())); |
|
|
|
ResourceModel resourceModel = context.getLwM2MTransportConfigServer().getResourceModel(lwM2MClient.getRegistration(), new LwM2mPath(path)); |
|
|
|
if (!path.isEmpty() && (this.validatePathInAttrProfile(profile, path) || this.validatePathInTelemetryProfile(profile, path))) { |
|
|
|
LwM2mClientProfile clientProfile = lwM2mClientContext.getProfile(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB())); |
|
|
|
if (path != null && !path.isEmpty() && (this.validatePathInAttrProfile(clientProfile, path) || this.validatePathInTelemetryProfile(clientProfile, path))) { |
|
|
|
ResourceModel resourceModel = lwM2mTransportContextServer.getLwM2MTransportConfigServer().getResourceModel(lwM2MClient.getRegistration(), new LwM2mPath(path)); |
|
|
|
if (resourceModel != null && resourceModel.operations.isWritable()) { |
|
|
|
lwM2mTransportRequest.sendAllRequest(leshanServer, lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE, |
|
|
|
ContentFormat.TLV.getName(), null, value, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE, |
|
|
|
ContentFormat.TLV.getName(), null, value, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
} else { |
|
|
|
log.error("Resource path - [{}] value - [{}] is not Writable and cannot be updated", path, value); |
|
|
|
String logMsg = String.format("%s: attributeUpdate: Resource path - %s value - %s is not Writable and cannot be updated", |
|
|
|
@ -353,9 +353,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* Trigger bootStrap path = "/1/0/9" - have to implemented on client |
|
|
|
*/ |
|
|
|
@Override |
|
|
|
public void doTrigger(LeshanServer lwServer, Registration registration, String path) { |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwServer, registration, path, POST_TYPE_OPER_EXECUTE, |
|
|
|
ContentFormat.TLV.getName(), null, null, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
public void doTrigger(Registration registration, String path) { |
|
|
|
lwM2mTransportRequest.sendAllRequest(registration, path, POST_TYPE_OPER_EXECUTE, |
|
|
|
ContentFormat.TLV.getName(), null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
@ -438,7 +438,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
public void updateParametersOnThingsboard(JsonElement msg, String topicName, Registration registration) { |
|
|
|
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registration); |
|
|
|
if (sessionInfo != null) { |
|
|
|
context.sentParametersOnThingsboard(msg, topicName, sessionInfo); |
|
|
|
lwM2mTransportContextServer.sentParametersOnThingsboard(msg, topicName, sessionInfo); |
|
|
|
} else { |
|
|
|
log.error("Client: [{}] updateParametersOnThingsboard [{}] sessionInfo ", registration, null); |
|
|
|
} |
|
|
|
@ -455,30 +455,29 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* - Request to the client after registration to read all resource values for all objects |
|
|
|
* - then Observe Request to the client marked as observe from the profile configuration. |
|
|
|
* |
|
|
|
* @param lwServer - LeshanServer |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
* @param lwM2MClient - object with All parameters off client |
|
|
|
*/ |
|
|
|
private void initLwM2mFromClientValue(LeshanServer lwServer, Registration registration, LwM2mClient lwM2MClient) { |
|
|
|
private void initLwM2mFromClientValue(Registration registration, LwM2mClient lwM2MClient) { |
|
|
|
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration); |
|
|
|
Set<String> clientObjects = this.getAllOjectsInClient(registration); |
|
|
|
if (clientObjects != null && !LwM2mTransportHandler.getClientOnlyObserveAfterConnect(lwM2MClientProfile)) { |
|
|
|
// #2
|
|
|
|
if (!LwM2mTransportHandler.getClientUpdateValueAfterConnect(lwM2MClientProfile)) { |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwServer, registration, lwM2MClient, GET_TYPE_OPER_READ); |
|
|
|
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_READ); |
|
|
|
|
|
|
|
} |
|
|
|
// #3
|
|
|
|
else { |
|
|
|
lwM2MClient.getPendingRequests().addAll(clientObjects); |
|
|
|
clientObjects.forEach(path -> { |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwServer, registration, path, GET_TYPE_OPER_READ, ContentFormat.TLV.getName(), |
|
|
|
null, null, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
lwM2mTransportRequest.sendAllRequest(registration, path, GET_TYPE_OPER_READ, ContentFormat.TLV.getName(), |
|
|
|
null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
}); |
|
|
|
} |
|
|
|
} |
|
|
|
// #1
|
|
|
|
this.initReadAttrTelemetryObserveToClient(lwServer, registration, lwM2MClient, GET_TYPE_OPER_OBSERVE); |
|
|
|
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_OBSERVE); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
@ -513,8 +512,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* #2 Update new Resources (replace old Resource Value on new Resource Value) |
|
|
|
* |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
* @param - LwM2mSingleResource response.getContent() |
|
|
|
* @param - LwM2mSingleResource response.getContent() |
|
|
|
* @param lwM2mResource - LwM2mSingleResource response.getContent() |
|
|
|
* @param path - resource |
|
|
|
*/ |
|
|
|
private void updateResourcesValue(Registration registration, LwM2mResource lwM2mResource, String path) { |
|
|
|
@ -522,31 +520,21 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
lwM2MClient.updateResourceValue(path, lwM2mResource); |
|
|
|
Set<String> paths = new HashSet<>(); |
|
|
|
paths.add(path); |
|
|
|
this.updateAttrTelemetry(registration, false, paths); |
|
|
|
this.updateAttrTelemetry(registration, paths); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* Sent Attribute and Telemetry to Thingsboard |
|
|
|
* #1 - get AttrName/TelemetryName with value: |
|
|
|
* #1.1 from Client |
|
|
|
* #1.2 from LwM2MClient: |
|
|
|
* #1 - get AttrName/TelemetryName with value from LwM2MClient: |
|
|
|
* -- resourceId == path from LwM2MClientProfile.postAttributeProfile/postTelemetryProfile/postObserveProfile |
|
|
|
* -- AttrName/TelemetryName == resourceName from ModelObject.objectModel, value from ModelObject.instance.resource(resourceId) |
|
|
|
* #2 - set Attribute/Telemetry |
|
|
|
* |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
*/ |
|
|
|
private void updateAttrTelemetry(Registration registration, boolean start, Set<String> paths) { |
|
|
|
private void updateAttrTelemetry(Registration registration, Set<String> paths) { |
|
|
|
JsonObject attributes = new JsonObject(); |
|
|
|
JsonObject telemetries = new JsonObject(); |
|
|
|
if (start) { |
|
|
|
// #1.1
|
|
|
|
JsonObject attributeClient = this.getAttributeClient(registration); |
|
|
|
if (attributeClient != null) { |
|
|
|
attributeClient.entrySet().forEach(p -> attributes.add(p.getKey(), p.getValue())); |
|
|
|
} |
|
|
|
} |
|
|
|
// #1.2
|
|
|
|
try { |
|
|
|
writeLock.lock(); |
|
|
|
this.getParametersFromProfile(attributes, telemetries, registration, paths); |
|
|
|
@ -562,25 +550,35 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* @param profile - |
|
|
|
* @param clientProfile - |
|
|
|
* @param path - |
|
|
|
* @return true if path isPresent in postAttributeProfile |
|
|
|
*/ |
|
|
|
private boolean validatePathInAttrProfile(LwM2mClientProfile profile, String path) { |
|
|
|
Set<String> attributesSet = new Gson().fromJson(profile.getPostAttributeProfile(), new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
return attributesSet.stream().filter(p -> p.equals(path)).findFirst().isPresent(); |
|
|
|
private boolean validatePathInAttrProfile(LwM2mClientProfile clientProfile, String path) { |
|
|
|
try { |
|
|
|
List<String> attributesSet = new Gson().fromJson(clientProfile.getPostAttributeProfile(), new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
return attributesSet.stream().anyMatch(p -> p.equals(path)); |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("Fail Validate Path [{}] ClientProfile.Attribute", path, e); |
|
|
|
return false; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* @param profile - |
|
|
|
* @param clientProfile - |
|
|
|
* @param path - |
|
|
|
* @return true if path isPresent in postAttributeProfile |
|
|
|
*/ |
|
|
|
private boolean validatePathInTelemetryProfile(LwM2mClientProfile profile, String path) { |
|
|
|
Set<String> telemetriesSet = new Gson().fromJson(profile.getPostTelemetryProfile(), new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
return telemetriesSet.stream().filter(p -> p.equals(path)).findFirst().isPresent(); |
|
|
|
private boolean validatePathInTelemetryProfile(LwM2mClientProfile clientProfile, String path) { |
|
|
|
try { |
|
|
|
List<String> telemetriesSet = new Gson().fromJson(clientProfile.getPostTelemetryProfile(), new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
return telemetriesSet.stream().anyMatch(p -> p.equals(path)); |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("Fail Validate Path [{}] ClientProfile.Telemetry", path, e); |
|
|
|
return false; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
@ -588,10 +586,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* #1 - Analyze: |
|
|
|
* #1.1 path in resource profile == client resource |
|
|
|
* |
|
|
|
* @param lwServer - |
|
|
|
* @param registration - |
|
|
|
*/ |
|
|
|
private void initReadAttrTelemetryObserveToClient(LeshanServer lwServer, Registration registration, LwM2mClient lwM2MClient, String typeOper) { |
|
|
|
private void initReadAttrTelemetryObserveToClient(Registration registration, LwM2mClient lwM2MClient, String typeOper) { |
|
|
|
try { |
|
|
|
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration); |
|
|
|
Set<String> clientInstances = this.getAllInstancesInClient(registration); |
|
|
|
@ -606,19 +603,18 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
}); |
|
|
|
} |
|
|
|
Set<String> pathSent = ConcurrentHashMap.newKeySet(); |
|
|
|
result.forEach(p -> { |
|
|
|
result.forEach(target -> { |
|
|
|
// #1.1
|
|
|
|
String target = p; |
|
|
|
String[] resPath = target.split("/"); |
|
|
|
String instance = "/" + resPath[1] + "/" + resPath[2]; |
|
|
|
if (clientInstances.contains(instance)) { |
|
|
|
if (clientInstances != null && clientInstances.size() > 0 && clientInstances.contains(instance)) { |
|
|
|
pathSent.add(target); |
|
|
|
} |
|
|
|
}); |
|
|
|
lwM2MClient.getPendingRequests().addAll(pathSent); |
|
|
|
pathSent.forEach(target -> { |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwServer, registration, target, typeOper, ContentFormat.TLV.getName(), |
|
|
|
null, null, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, ContentFormat.TLV.getName(), |
|
|
|
null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
}); |
|
|
|
if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) { |
|
|
|
lwM2MClient.initValue(this, null); |
|
|
|
@ -677,26 +673,26 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
return (clientInstances.size() > 0) ? clientInstances : null; |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* get AttrName/TelemetryName with value from Client |
|
|
|
* |
|
|
|
* @param registration - |
|
|
|
* @return - JsonObject, format: {name: value}} |
|
|
|
*/ |
|
|
|
private JsonObject getAttributeClient(Registration registration) { |
|
|
|
if (registration.getAdditionalRegistrationAttributes().size() > 0) { |
|
|
|
JsonObject resNameValues = new JsonObject(); |
|
|
|
registration.getAdditionalRegistrationAttributes().forEach(resNameValues::addProperty); |
|
|
|
return resNameValues; |
|
|
|
} |
|
|
|
return null; |
|
|
|
} |
|
|
|
// /**
|
|
|
|
// * get AttrName/TelemetryName with value from Client
|
|
|
|
// *
|
|
|
|
// * @param registration -
|
|
|
|
// * @return - JsonObject, format: {name: value}}
|
|
|
|
// */
|
|
|
|
// private JsonObject getAttributeClient(Registration registration) {
|
|
|
|
// if (registration.getAdditionalRegistrationAttributes().size() > 0) {
|
|
|
|
// JsonObject resNameValues = new JsonObject();
|
|
|
|
// registration.getAdditionalRegistrationAttributes().forEach(resNameValues::addProperty);
|
|
|
|
// return resNameValues;
|
|
|
|
// }
|
|
|
|
// return null;
|
|
|
|
// }
|
|
|
|
|
|
|
|
/** |
|
|
|
* @param attributes - new JsonObject |
|
|
|
* @param telemetry - new JsonObject |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
* @param path |
|
|
|
* @param path - |
|
|
|
*/ |
|
|
|
private void getParametersFromProfile(JsonObject attributes, JsonObject telemetry, Registration registration, Set<String> path) { |
|
|
|
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration); |
|
|
|
@ -746,7 +742,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
LwM2mPath pathIds = new LwM2mPath(path); |
|
|
|
ResourceValue resourceValue = this.returnResourceValueFromLwM2MClient(lwM2MClient, pathIds); |
|
|
|
return resourceValue == null ? null : |
|
|
|
this.converter.convertValue(resourceValue.getResourceValue(), this.context.getLwM2MTransportConfigServer().getResourceModelType(lwM2MClient.getRegistration(), pathIds), ResourceModel.Type.STRING, pathIds).toString(); |
|
|
|
this.converter.convertValue(resourceValue.getResourceValue(), this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getResourceModelType(lwM2MClient.getRegistration(), pathIds), ResourceModel.Type.STRING, pathIds).toString(); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
@ -796,25 +792,21 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* @param deviceProfile - |
|
|
|
*/ |
|
|
|
private void onDeviceUpdateChangeProfile(Set<String> registrationIds, DeviceProfile deviceProfile) { |
|
|
|
LwM2mClientProfile lwM2MClientProfileOld = lwM2mClientContext.getProfiles().get(deviceProfile.getUuidId()); |
|
|
|
LwM2mClientProfile lwM2MClientProfileOld = lwM2mClientContext.getProfiles().get(deviceProfile.getUuidId()).clone(); |
|
|
|
if (lwM2mClientContext.addUpdateProfileParameters(deviceProfile)) { |
|
|
|
// #1
|
|
|
|
JsonArray attributeOld = lwM2MClientProfileOld.getPostAttributeProfile(); |
|
|
|
Set<String> attributeSetOld = new Gson().fromJson(attributeOld, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
Set attributeSetOld = this.convertJsonArrayToSet (attributeOld); |
|
|
|
JsonArray telemetryOld = lwM2MClientProfileOld.getPostTelemetryProfile(); |
|
|
|
Set<String> telemetrySetOld = new Gson().fromJson(telemetryOld, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
Set telemetrySetOld = this.convertJsonArrayToSet (telemetryOld); |
|
|
|
JsonArray observeOld = lwM2MClientProfileOld.getPostObserveProfile(); |
|
|
|
JsonObject keyNameOld = lwM2MClientProfileOld.getPostKeyNameProfile(); |
|
|
|
|
|
|
|
LwM2mClientProfile lwM2MClientProfileNew = lwM2mClientContext.getProfiles().get(deviceProfile.getUuidId()); |
|
|
|
JsonArray attributeNew = lwM2MClientProfileNew.getPostAttributeProfile(); |
|
|
|
Set<String> attributeSetNew = new Gson().fromJson(attributeNew, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
Set<String> attributeSetNew = this.convertJsonArrayToSet (attributeNew); |
|
|
|
JsonArray telemetryNew = lwM2MClientProfileNew.getPostTelemetryProfile(); |
|
|
|
Set<String> telemetrySetNew = new Gson().fromJson(telemetryNew, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
Set telemetrySetNew = this.convertJsonArrayToSet (telemetryNew); |
|
|
|
JsonArray observeNew = lwM2MClientProfileNew.getPostObserveProfile(); |
|
|
|
JsonObject keyNameNew = lwM2MClientProfileNew.getPostKeyNameProfile(); |
|
|
|
|
|
|
|
@ -847,11 +839,11 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
if (sentAttrToThingsboard.getPathPostParametersAdd().size() > 0) { |
|
|
|
// update value in Resources
|
|
|
|
registrationIds.forEach(registrationId -> { |
|
|
|
LeshanServer lwServer = leshanServer; |
|
|
|
// LeshanServer lwServer = leshanServer;
|
|
|
|
Registration registration = lwM2mClientContext.getRegistration(registrationId); |
|
|
|
this.readResourceValueObserve(lwServer, registration, sentAttrToThingsboard.getPathPostParametersAdd(), GET_TYPE_OPER_READ); |
|
|
|
this.readResourceValueObserve(registration, sentAttrToThingsboard.getPathPostParametersAdd(), GET_TYPE_OPER_READ); |
|
|
|
// sent attr/telemetry to tingsboard for new path
|
|
|
|
this.updateAttrTelemetry(registration, false, sentAttrToThingsboard.getPathPostParametersAdd()); |
|
|
|
this.updateAttrTelemetry(registration, sentAttrToThingsboard.getPathPostParametersAdd()); |
|
|
|
}); |
|
|
|
} |
|
|
|
// #4.2 del
|
|
|
|
@ -862,10 +854,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
|
|
|
|
// #5.1
|
|
|
|
if (!observeOld.equals(observeNew)) { |
|
|
|
Set<String> observeSetOld = new Gson().fromJson(observeOld, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
Set<String> observeSetNew = new Gson().fromJson(observeNew, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
Set<String> observeSetOld = new Gson().fromJson(observeOld, new TypeToken<>() {}.getType()); |
|
|
|
Set<String> observeSetNew = new Gson().fromJson(observeNew, new TypeToken<>() {}.getType()); |
|
|
|
//#5.2 add
|
|
|
|
// path Attr/Telemetry includes newObserve
|
|
|
|
attributeSetOld.addAll(telemetrySetOld); |
|
|
|
@ -876,17 +866,22 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
ResultsAnalyzerParameters postObserveAnalyzer = this.getAnalyzerParameters(sentObserveToClientOld.getPathPostParametersAdd(), sentObserveToClientNew.getPathPostParametersAdd()); |
|
|
|
// sent Request observe to Client
|
|
|
|
registrationIds.forEach(registrationId -> { |
|
|
|
LeshanServer lwServer = leshanServer; |
|
|
|
Registration registration = lwM2mClientContext.getRegistration(registrationId); |
|
|
|
this.readResourceValueObserve(lwServer, registration, postObserveAnalyzer.getPathPostParametersAdd(), GET_TYPE_OPER_OBSERVE); |
|
|
|
this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), GET_TYPE_OPER_OBSERVE); |
|
|
|
// 5.3 del
|
|
|
|
// sent Request cancel observe to Client
|
|
|
|
this.cancelObserveIsValue(lwServer, registration, postObserveAnalyzer.getPathPostParametersDel()); |
|
|
|
this.cancelObserveIsValue(registration, postObserveAnalyzer.getPathPostParametersDel()); |
|
|
|
}); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private Set <String> convertJsonArrayToSet (JsonArray jsonArray) { |
|
|
|
List<String> attributeListOld = new Gson().fromJson(jsonArray, new TypeToken<>() { |
|
|
|
}.getType()); |
|
|
|
return Sets.newConcurrentHashSet(attributeListOld); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* Compare old list with new list after change AttrTelemetryObserve in config Profile |
|
|
|
* |
|
|
|
@ -917,20 +912,19 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
* Update Resource value after change RezAttrTelemetry in config Profile |
|
|
|
* sent response Read to Client and add path to pathResAttrTelemetry in LwM2MClient.getAttrTelemetryObserveValue() |
|
|
|
* |
|
|
|
* @param lwServer - LeshanServer |
|
|
|
* @param registration - Registration LwM2M Client |
|
|
|
* @param targets - path Resources == [ "/2/0/0", "/2/0/1"] |
|
|
|
*/ |
|
|
|
private void readResourceValueObserve(LeshanServer lwServer, Registration registration, Set<String> targets, String typeOper) { |
|
|
|
private void readResourceValueObserve(Registration registration, Set<String> targets, String typeOper) { |
|
|
|
targets.forEach(target -> { |
|
|
|
LwM2mPath pathIds = new LwM2mPath(target); |
|
|
|
if (pathIds.isResource()) { |
|
|
|
if (GET_TYPE_OPER_READ.equals(typeOper)) { |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwServer, registration, target, typeOper, |
|
|
|
ContentFormat.TLV.getName(), null, null, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, |
|
|
|
ContentFormat.TLV.getName(), null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
} else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) { |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwServer, registration, target, typeOper, |
|
|
|
null, null, null, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, |
|
|
|
null, null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
} |
|
|
|
} |
|
|
|
}); |
|
|
|
@ -946,11 +940,11 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
return analyzerParameters; |
|
|
|
} |
|
|
|
|
|
|
|
private void cancelObserveIsValue(LeshanServer lwServer, Registration registration, Set<String> paramAnallyzer) { |
|
|
|
private void cancelObserveIsValue(Registration registration, Set<String> paramAnallyzer) { |
|
|
|
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null); |
|
|
|
paramAnallyzer.forEach(p -> { |
|
|
|
if (this.returnResourceValueFromLwM2MClient(lwM2MClient, new LwM2mPath(p)) != null) { |
|
|
|
this.setCancelObservationRecourse(lwServer, registration, p); |
|
|
|
this.setCancelObservationRecourse(registration, p); |
|
|
|
} |
|
|
|
} |
|
|
|
); |
|
|
|
@ -958,8 +952,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
|
|
|
|
private void putDelayedUpdateResourcesClient(LwM2mClient lwM2MClient, Object valueOld, Object valueNew, String path) { |
|
|
|
if (valueNew != null && (valueOld == null || !valueNew.toString().equals(valueOld.toString()))) { |
|
|
|
lwM2mTransportRequest.sendAllRequest(leshanServer, lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE, |
|
|
|
ContentFormat.TLV.getName(), null, valueNew, this.context.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE, |
|
|
|
ContentFormat.TLV.getName(), null, valueNew, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); |
|
|
|
} else { |
|
|
|
log.error("05 delayError"); |
|
|
|
} |
|
|
|
@ -1037,7 +1031,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
|
|
|
|
/** |
|
|
|
* @param lwM2MClient - |
|
|
|
* @return |
|
|
|
* @return SessionInfoProto - |
|
|
|
*/ |
|
|
|
private SessionInfoProto getNewSessionInfoProto(LwM2mClient lwM2MClient) { |
|
|
|
if (lwM2MClient != null) { |
|
|
|
@ -1048,7 +1042,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
return null; |
|
|
|
} else { |
|
|
|
return SessionInfoProto.newBuilder() |
|
|
|
.setNodeId(this.context.getNodeId()) |
|
|
|
.setNodeId(this.lwM2mTransportContextServer.getNodeId()) |
|
|
|
.setSessionIdMSB(lwM2MClient.getSessionId().getMostSignificantBits()) |
|
|
|
.setSessionIdLSB(lwM2MClient.getSessionId().getLeastSignificantBits()) |
|
|
|
.setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB()) |
|
|
|
@ -1114,7 +1108,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
if (attrSharedNames.size() > 0) { |
|
|
|
//#2.1
|
|
|
|
try { |
|
|
|
TransportProtos.GetAttributeRequestMsg getAttributeMsg = context.getAdaptor().convertToGetAttributes(null, attrSharedNames); |
|
|
|
TransportProtos.GetAttributeRequestMsg getAttributeMsg = lwM2mTransportContextServer.getAdaptor().convertToGetAttributes(null, attrSharedNames); |
|
|
|
transportService.process(sessionInfo, getAttributeMsg, getAckCallback(lwM2MClient, getAttributeMsg.getRequestId(), DEVICE_ATTRIBUTES_REQUEST)); |
|
|
|
} catch (AdaptorException e) { |
|
|
|
log.warn("Failed to decode get attributes request", e); |
|
|
|
@ -1138,8 +1132,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { |
|
|
|
|
|
|
|
ConcurrentMap<String, String> keyNamesIsWritable = keyNamesMap.entrySet() |
|
|
|
.stream() |
|
|
|
.filter(e -> (attrSet.contains(e.getKey()) && context.getLwM2MTransportConfigServer().getResourceModel(lwM2MClient.getRegistration(), new LwM2mPath(e.getKey())) != null && |
|
|
|
context.getLwM2MTransportConfigServer().getResourceModel(lwM2MClient.getRegistration(), new LwM2mPath(e.getKey())).operations.isWritable())) |
|
|
|
.filter(e -> (attrSet.contains(e.getKey()) && lwM2mTransportContextServer.getLwM2MTransportConfigServer().getResourceModel(lwM2MClient.getRegistration(), new LwM2mPath(e.getKey())) != null && |
|
|
|
lwM2mTransportContextServer.getLwM2MTransportConfigServer().getResourceModel(lwM2MClient.getRegistration(), new LwM2mPath(e.getKey())).operations.isWritable())) |
|
|
|
.collect(Collectors.toConcurrentMap(Map.Entry::getKey, Map.Entry::getValue)); |
|
|
|
|
|
|
|
Set<String> namesIsWritable = ConcurrentHashMap.newKeySet(); |
|
|
|
|