|
|
|
@ -89,10 +89,8 @@ import java.util.stream.Collectors; |
|
|
|
|
|
|
|
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST; |
|
|
|
import static org.eclipse.leshan.core.attributes.Attribute.OBJECT_VERSION; |
|
|
|
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.DOWNLOADED; |
|
|
|
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.FAILED; |
|
|
|
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.INITIATED; |
|
|
|
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.UPDATING; |
|
|
|
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; |
|
|
|
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto; |
|
|
|
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.DEVICE_ATTRIBUTES_REQUEST; |
|
|
|
@ -125,7 +123,7 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.v |
|
|
|
@Slf4j |
|
|
|
@Service |
|
|
|
@TbLwM2mTransportComponent |
|
|
|
public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler { |
|
|
|
public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { |
|
|
|
|
|
|
|
private ExecutorService registrationExecutor; |
|
|
|
private ExecutorService updateRegistrationExecutor; |
|
|
|
@ -144,11 +142,11 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
private final Map<UUID, Long> rpcSubscriptions; |
|
|
|
public final Map<String, Integer> firmwareUpdateState; |
|
|
|
|
|
|
|
public DefaultLwM2MTransportMsgHandler(TransportService transportService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper, |
|
|
|
LwM2mClientContext clientContext, |
|
|
|
@Lazy LwM2mTransportRequest lwM2mTransportRequest, |
|
|
|
OtaPackageDataCache otaPackageDataCache, |
|
|
|
LwM2mTransportContext context, LwM2MJsonAdaptor adaptor, TbLwM2MDtlsSessionStore sessionStore) { |
|
|
|
public DefaultLwM2MUplinkMsgHandler(TransportService transportService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper, |
|
|
|
LwM2mClientContext clientContext, |
|
|
|
@Lazy LwM2mTransportRequest lwM2mTransportRequest, |
|
|
|
OtaPackageDataCache otaPackageDataCache, |
|
|
|
LwM2mTransportContext context, LwM2MJsonAdaptor adaptor, TbLwM2MDtlsSessionStore sessionStore) { |
|
|
|
this.transportService = transportService; |
|
|
|
this.config = config; |
|
|
|
this.helper = helper; |
|
|
|
@ -233,6 +231,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
updateRegistrationExecutor.submit(() -> { |
|
|
|
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint()); |
|
|
|
try { |
|
|
|
log.warn("[{}] [{{}] Client: update after Registration", registration.getEndpoint(), registration.getId()); |
|
|
|
clientContext.updateRegistration(lwM2MClient, registration); |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = lwM2MClient.getSession(); |
|
|
|
this.reportActivityAndRegister(sessionInfo); |
|
|
|
@ -267,8 +266,8 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
clientContext.unregister(client, registration); |
|
|
|
SessionInfoProto sessionInfo = client.getSession(); |
|
|
|
if (sessionInfo != null) { |
|
|
|
transportService.process(sessionInfo, DefaultTransportService.getSessionEventMsg(SessionEvent.CLOSED), null); |
|
|
|
transportService.deregisterSession(sessionInfo); |
|
|
|
this.doCloseSession(sessionInfo); |
|
|
|
sessionStore.remove(registration.getEndpoint()); |
|
|
|
log.info("Client close session: [{}] unReg [{}] name [{}] profile ", registration.getId(), registration.getEndpoint(), sessionInfo.getDeviceType()); |
|
|
|
} else { |
|
|
|
@ -552,19 +551,6 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
transportService.deregisterSession(sessionInfo); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* Session device in thingsboard is closed |
|
|
|
* |
|
|
|
* @param sessionInfo - lwm2m client |
|
|
|
*/ |
|
|
|
private void doCloseSession(SessionInfoProto sessionInfo) { |
|
|
|
TransportProtos.SessionEvent event = SessionEvent.CLOSED; |
|
|
|
TransportProtos.SessionEventMsg msg = TransportProtos.SessionEventMsg.newBuilder() |
|
|
|
.setSessionType(TransportProtos.SessionType.ASYNC) |
|
|
|
.setEvent(event).build(); |
|
|
|
transportService.process(sessionInfo, msg, null); |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* 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. |
|
|
|
@ -609,20 +595,20 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
* @param lwM2MClient - object with All parameters off client |
|
|
|
*/ |
|
|
|
private void initClientTelemetry(LwM2mClient lwM2MClient) { |
|
|
|
LwM2mClientProfile lwM2MClientProfile = clientContext.getProfile(lwM2MClient.getProfileId()); |
|
|
|
Set<String> clientObjects = clientContext.getSupportedIdVerInClient(lwM2MClient); |
|
|
|
if (clientObjects != null && clientObjects.size() > 0) { |
|
|
|
if (LwM2mTransportUtil.LwM2MClientStrategy.CLIENT_STRATEGY_2.code == lwM2MClientProfile.getClientStrategy()) { |
|
|
|
LwM2mClientProfile profile = clientContext.getProfile(lwM2MClient.getProfileId()); |
|
|
|
Set<String> supportedObjects = clientContext.getSupportedIdVerInClient(lwM2MClient); |
|
|
|
if (supportedObjects != null && supportedObjects.size() > 0) { |
|
|
|
if (LwM2mTransportUtil.LwM2MClientStrategy.CLIENT_STRATEGY_2.code == profile.getClientStrategy()) { |
|
|
|
// #2
|
|
|
|
lwM2MClient.getPendingReadRequests().addAll(clientObjects); |
|
|
|
clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(lwM2MClient, path, READ, |
|
|
|
lwM2MClient.getPendingReadRequests().addAll(supportedObjects); |
|
|
|
supportedObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(lwM2MClient, path, READ, |
|
|
|
null, this.config.getTimeout(), null)); |
|
|
|
} |
|
|
|
// #1
|
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, READ, clientObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, OBSERVE, clientObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, WRITE_ATTRIBUTES, clientObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, DISCOVER, clientObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, profile, READ, supportedObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, profile, OBSERVE, supportedObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, profile, WRITE_ATTRIBUTES, supportedObjects); |
|
|
|
this.initReadAttrTelemetryObserveToClient(lwM2MClient, profile, DISCOVER, supportedObjects); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -717,8 +703,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void initReadAttrTelemetryObserveToClient(LwM2mClient lwM2MClient, LwM2mTypeOper typeOper, Set<String> clientObjects) { |
|
|
|
LwM2mClientProfile lwM2MClientProfile = clientContext.getProfile(lwM2MClient.getProfileId()); |
|
|
|
private void initReadAttrTelemetryObserveToClient(LwM2mClient lwM2MClient, LwM2mClientProfile lwM2MClientProfile, LwM2mTypeOper typeOper, Set<String> supportedObjects) { |
|
|
|
Set<String> result = null; |
|
|
|
ConcurrentHashMap<String, Object> params = null; |
|
|
|
if (READ.equals(typeOper)) { |
|
|
|
@ -738,7 +723,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
params = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile()); |
|
|
|
result = params.keySet(); |
|
|
|
} |
|
|
|
sendRequestsToClient(lwM2MClient, typeOper, clientObjects, result, params); |
|
|
|
sendRequestsToClient(lwM2MClient, typeOper, supportedObjects, result, params); |
|
|
|
} |
|
|
|
|
|
|
|
private void sendRequestsToClient(LwM2mClient lwM2MClient, LwM2mTypeOper operationType, Set<String> supportedObjectIds, Set<String> desiredObjectIds, ConcurrentHashMap<String, Object> params) { |
|
|
|
@ -1325,7 +1310,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
if (lwM2MClient.getRegistration().getSupportedVersion(FW_5_ID) != null) { |
|
|
|
SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient); |
|
|
|
if (sessionInfo != null) { |
|
|
|
DefaultLwM2MTransportMsgHandler handler = this; |
|
|
|
DefaultLwM2MUplinkMsgHandler handler = this; |
|
|
|
this.transportService.process(sessionInfo, createOtaPackageRequestMsg(sessionInfo, OtaPackageType.FIRMWARE.name()), |
|
|
|
new TransportServiceCallback<>() { |
|
|
|
@Override |
|
|
|
@ -1375,7 +1360,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler |
|
|
|
if (lwM2MClient.getRegistration().getSupportedVersion(SW_ID) != null) { |
|
|
|
SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient); |
|
|
|
if (sessionInfo != null) { |
|
|
|
DefaultLwM2MTransportMsgHandler handler = this; |
|
|
|
DefaultLwM2MUplinkMsgHandler handler = this; |
|
|
|
transportService.process(sessionInfo, createOtaPackageRequestMsg(sessionInfo, OtaPackageType.SOFTWARE.name()), |
|
|
|
new TransportServiceCallback<>() { |
|
|
|
@Override |