|
|
|
@ -26,24 +26,31 @@ import org.eclipse.leshan.core.node.LwM2mResource; |
|
|
|
import org.eclipse.leshan.core.node.ObjectLink; |
|
|
|
import org.eclipse.leshan.core.node.codec.CodecException; |
|
|
|
import org.eclipse.leshan.core.observation.Observation; |
|
|
|
import org.eclipse.leshan.core.request.CompositeDownlinkRequest; |
|
|
|
import org.eclipse.leshan.core.request.ContentFormat; |
|
|
|
import org.eclipse.leshan.core.request.DeleteRequest; |
|
|
|
import org.eclipse.leshan.core.request.DiscoverRequest; |
|
|
|
import org.eclipse.leshan.core.request.DownlinkRequest; |
|
|
|
import org.eclipse.leshan.core.request.ExecuteRequest; |
|
|
|
import org.eclipse.leshan.core.request.ObserveRequest; |
|
|
|
import org.eclipse.leshan.core.request.ReadCompositeRequest; |
|
|
|
import org.eclipse.leshan.core.request.ReadRequest; |
|
|
|
import org.eclipse.leshan.core.request.SimpleDownlinkRequest; |
|
|
|
import org.eclipse.leshan.core.request.WriteAttributesRequest; |
|
|
|
import org.eclipse.leshan.core.request.WriteCompositeRequest; |
|
|
|
import org.eclipse.leshan.core.request.WriteRequest; |
|
|
|
import org.eclipse.leshan.core.response.DeleteResponse; |
|
|
|
import org.eclipse.leshan.core.response.DiscoverResponse; |
|
|
|
import org.eclipse.leshan.core.response.ExecuteResponse; |
|
|
|
import org.eclipse.leshan.core.response.LwM2mResponse; |
|
|
|
import org.eclipse.leshan.core.response.ObserveResponse; |
|
|
|
import org.eclipse.leshan.core.response.ReadCompositeResponse; |
|
|
|
import org.eclipse.leshan.core.response.ReadResponse; |
|
|
|
import org.eclipse.leshan.core.response.WriteAttributesResponse; |
|
|
|
import org.eclipse.leshan.core.response.WriteCompositeResponse; |
|
|
|
import org.eclipse.leshan.core.response.WriteResponse; |
|
|
|
import org.eclipse.leshan.core.util.Hex; |
|
|
|
import org.eclipse.leshan.server.model.LwM2mModelProvider; |
|
|
|
import org.eclipse.leshan.server.registration.Registration; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
@ -53,7 +60,10 @@ import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.common.LwM2MExecutorAwareService; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MReadCompositeRequest; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MWriteResponseCompositeCallback; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; |
|
|
|
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; |
|
|
|
|
|
|
|
import javax.annotation.PostConstruct; |
|
|
|
@ -63,6 +73,7 @@ import java.util.Collection; |
|
|
|
import java.util.Date; |
|
|
|
import java.util.LinkedList; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.function.Function; |
|
|
|
import java.util.function.Predicate; |
|
|
|
@ -73,6 +84,7 @@ import static org.eclipse.leshan.core.attributes.Attribute.LESSER_THAN; |
|
|
|
import static org.eclipse.leshan.core.attributes.Attribute.MAXIMUM_PERIOD; |
|
|
|
import static org.eclipse.leshan.core.attributes.Attribute.MINIMUM_PERIOD; |
|
|
|
import static org.eclipse.leshan.core.attributes.Attribute.STEP; |
|
|
|
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId; |
|
|
|
|
|
|
|
@Slf4j |
|
|
|
@Service |
|
|
|
@ -110,8 +122,31 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
@Override |
|
|
|
public void sendReadRequest(LwM2mClient client, TbLwM2MReadRequest request, DownlinkRequestCallback<ReadRequest, ReadResponse> callback) { |
|
|
|
validateVersionedId(client, request); |
|
|
|
ReadRequest downlink = new ReadRequest(getContentFormat(client, request), request.getObjectId()); |
|
|
|
sendRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
ReadRequest downlink = new ReadRequest(getRequestContentFormat(client, request, this.config.getModelProvider()), request.getObjectId()); |
|
|
|
sendSimpleRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void sendReadCompositeRequest(LwM2mClient client, TbLwM2MReadCompositeRequest request, DownlinkRequestCallback<ReadCompositeRequest, ReadCompositeResponse> callback) { |
|
|
|
validateVersionedIds(client, request); |
|
|
|
ContentFormat requestContentFormat = ContentFormat.SENML_JSON; |
|
|
|
ContentFormat responseContentFormat = ContentFormat.SENML_JSON; |
|
|
|
|
|
|
|
ReadCompositeRequest downlink = new ReadCompositeRequest(requestContentFormat, responseContentFormat, request.getObjectIds()); |
|
|
|
sendCompositeRequest(client, downlink, this.config.getTimeout(), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void sendWriteCompositeRequest(LwM2mClient client, Map<String, Object> nodes, DefaultLwM2MUplinkMsgHandler handler) { |
|
|
|
// ResourceModel resourceModelWrite = client.getResourceModel(request.getVersionedId(), this.config.getModelProvider());
|
|
|
|
TbLwM2MWriteResponseCompositeCallback callback = new TbLwM2MWriteResponseCompositeCallback(handler, logService, client, null); |
|
|
|
ContentFormat contentFormat = ContentFormat.SENML_JSON; |
|
|
|
try { |
|
|
|
WriteCompositeRequest downlink = new WriteCompositeRequest(contentFormat, nodes); |
|
|
|
sendWriteCompositeRequest(client, downlink, this.config.getTimeout(), callback); |
|
|
|
} catch (Exception e) { |
|
|
|
callback.onError(JacksonUtil.toString(nodes), e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -121,7 +156,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
Set<Observation> observations = context.getServer().getObservationService().getObservations(client.getRegistration()); |
|
|
|
if (observations.stream().noneMatch(observation -> observation.getPath().equals(resultIds))) { |
|
|
|
ObserveRequest downlink; |
|
|
|
ContentFormat contentFormat = getContentFormat(client, request); |
|
|
|
ContentFormat contentFormat = getRequestContentFormat(client, request, this.config.getModelProvider()); |
|
|
|
if (resultIds.isResource()) { |
|
|
|
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId()); |
|
|
|
} else if (resultIds.isObjectInstance()) { |
|
|
|
@ -130,7 +165,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId()); |
|
|
|
} |
|
|
|
log.info("[{}] Send observation: {}.", client.getEndpoint(), request.getVersionedId()); |
|
|
|
sendRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
} else { |
|
|
|
callback.onValidationError(resultIds.toString(), "Observation is already registered!"); |
|
|
|
} |
|
|
|
@ -158,13 +193,13 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
} else { |
|
|
|
downlink = new ExecuteRequest(request.getVersionedId()); |
|
|
|
} |
|
|
|
sendRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void sendDeleteRequest(LwM2mClient client, TbLwM2MDeleteRequest request, DownlinkRequestCallback<DeleteRequest, DeleteResponse> callback) { |
|
|
|
sendRequest(client, new DeleteRequest(request.getObjectId()), request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, new DeleteRequest(request.getObjectId()), request.getTimeout(), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -182,7 +217,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
@Override |
|
|
|
public void sendDiscoverRequest(LwM2mClient client, TbLwM2MDiscoverRequest request, DownlinkRequestCallback<DiscoverRequest, DiscoverResponse> callback) { |
|
|
|
validateVersionedId(client, request); |
|
|
|
sendRequest(client, new DiscoverRequest(request.getObjectId()), request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, new DiscoverRequest(request.getObjectId()), request.getTimeout(), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -202,7 +237,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
addAttribute(attributes, LESSER_THAN, params.getLt()); |
|
|
|
addAttribute(attributes, STEP, params.getSt()); |
|
|
|
AttributeSet attributeSet = new AttributeSet(attributes); |
|
|
|
sendRequest(client, new WriteAttributesRequest(request.getObjectId(), attributeSet), request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, new WriteAttributesRequest(request.getObjectId(), attributeSet), request.getTimeout(), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -214,7 +249,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
LwM2mPath path = new LwM2mPath(request.getObjectId()); |
|
|
|
WriteRequest downlink = this.getWriteRequestSingleResource(resourceModelWrite.type, contentFormat, |
|
|
|
path.getObjectId(), path.getObjectInstanceId(), path.getResourceId(), request.getValue()); |
|
|
|
sendRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
} catch (Exception e) { |
|
|
|
callback.onError(JacksonUtil.toString(request), e); |
|
|
|
} |
|
|
|
@ -236,7 +271,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
ContentFormat contentFormat = request.getObjectContentFormat() != null ? request.getObjectContentFormat() : convertResourceModelTypeToContentFormat(client, resourceModelWrite.type); |
|
|
|
WriteRequest downlink = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, resultIds.getObjectId(), |
|
|
|
resultIds.getObjectInstanceId(), resources); |
|
|
|
sendRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
} else if (resultIds.isObjectInstance()) { |
|
|
|
/* |
|
|
|
* params = "{\"id\":0,\"resources\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}" |
|
|
|
@ -245,9 +280,9 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
*/ |
|
|
|
Collection<LwM2mResource> resources = client.getNewResourcesForInstance(request.getVersionedId(), request.getValue(), this.config.getModelProvider(), this.converter); |
|
|
|
if (resources.size() > 0) { |
|
|
|
ContentFormat contentFormat = request.getObjectContentFormat() != null ? request.getObjectContentFormat() : client.getDefaultContentFormat(); |
|
|
|
ContentFormat contentFormat = request.getObjectContentFormat() != null ? request.getObjectContentFormat() : ContentFormat.DEFAULT; |
|
|
|
WriteRequest downlink = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resources); |
|
|
|
sendRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
sendSimpleRequest(client, downlink, request.getTimeout(), callback); |
|
|
|
} else { |
|
|
|
callback.onValidationError(JacksonUtil.toString(request), "No resources to update!"); |
|
|
|
} |
|
|
|
@ -256,10 +291,42 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private <R extends SimpleDownlinkRequest<T>, T extends LwM2mResponse> void sendRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback<R, T> callback) { |
|
|
|
|
|
|
|
private <R extends SimpleDownlinkRequest<T>, T extends LwM2mResponse> void sendSimpleRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback<R, T> callback) { |
|
|
|
sendRequest(client, request, timeoutInMs, callback, r -> request.getPath().toString()); |
|
|
|
} |
|
|
|
|
|
|
|
private <R extends CompositeDownlinkRequest<T>, T extends LwM2mResponse> void sendCompositeRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback<R, T> callback) { |
|
|
|
sendRequest(client, request, timeoutInMs, callback, r -> request.getPaths().toString()); |
|
|
|
} |
|
|
|
|
|
|
|
private <R extends DownlinkRequest<T>, T extends LwM2mResponse> void sendRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback<R, T> callback, Function<R, String> pathToStringFunction) { |
|
|
|
Registration registration = client.getRegistration(); |
|
|
|
try { |
|
|
|
logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), pathToStringFunction.apply(request))); |
|
|
|
context.getServer().send(registration, request, timeoutInMs, response -> { |
|
|
|
executor.submit(() -> { |
|
|
|
try { |
|
|
|
callback.onSuccess(request, response); |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}] failed to process successful response [{}] ", registration.getEndpoint(), response, e); |
|
|
|
} |
|
|
|
}); |
|
|
|
}, e -> { |
|
|
|
executor.submit(() -> { |
|
|
|
callback.onError(JacksonUtil.toString(request), e); |
|
|
|
}); |
|
|
|
}); |
|
|
|
} catch (Exception e) { |
|
|
|
callback.onError(JacksonUtil.toString(request), e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private <R extends SimpleDownlinkRequest<T>, T extends LwM2mResponse> void sendWriteCompositeRequest(LwM2mClient client, WriteCompositeRequest request, long timeoutInMs, |
|
|
|
DownlinkRequestCallback<WriteCompositeRequest, WriteCompositeResponse> callback) { |
|
|
|
Registration registration = client.getRegistration(); |
|
|
|
try { |
|
|
|
logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), request.getPath())); |
|
|
|
logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), request.getPaths())); |
|
|
|
context.getServer().send(registration, request, timeoutInMs, response -> { |
|
|
|
executor.submit(() -> { |
|
|
|
try { |
|
|
|
@ -278,6 +345,52 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
// private <R extends DownlinkRequest<T>, T extends LwM2mResponse> void sendReadRequestComposite(LwM2mClient client, ReadCompositeRequest request, long timeoutInMs,
|
|
|
|
// DownlinkRequestCallback<ReadCompositeRequest, ReadCompositeResponse> callback) {
|
|
|
|
// Registration registration = client.getRegistration();
|
|
|
|
// try {
|
|
|
|
// logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), request.getPaths()));
|
|
|
|
// context.getServer().send(registration, request, timeoutInMs, response -> {
|
|
|
|
// executor.submit(() -> {
|
|
|
|
// try {
|
|
|
|
// /**
|
|
|
|
// * [{"bn":"/3/0/","n":"0","vs":"Thingsboard Test Device"},
|
|
|
|
// * {"n":"1","vs":"Model 500"},
|
|
|
|
// * {"n":"2","vs":"TH-500-000-0001"},
|
|
|
|
// * {"n":"3","vs":"TestThingsboard@TestMore1024_2.04"},
|
|
|
|
// * {"n":"6","v":1},{"n":"7","v":56},
|
|
|
|
// * {"n":"8","v":42},{"n":"9","v":16},
|
|
|
|
// * {"n":"10","v":127619},{"n":"13","v":1624520988},
|
|
|
|
// * {"n":"14","vs":"+03"},{"n":"15","vs":"Europe/Kiev"},
|
|
|
|
// * {"n":"16","vs":"U"},{"n":"17","vs":"smart meters"},
|
|
|
|
// * {"n":"18","vs":"1.01"},{"n":"19","vs":"1.02"},
|
|
|
|
// * {"n":"20","v":3},{"n":"21","v":256000},
|
|
|
|
// * {"bn":"/5/0/","n":"1","vs":""},
|
|
|
|
// * {"n":"3","v":0},{"n":"5","v":0},
|
|
|
|
// * {"n":"6","vs":""},{"n":"7","vs":""},
|
|
|
|
// * {"n":"8/0","v":0},{"n":"8/1","v":1},
|
|
|
|
// * {"n":"9","v":2},
|
|
|
|
// * {"bn":"/1/0/","n":"0","v":123},
|
|
|
|
// * {"n":"1","v":300},
|
|
|
|
// * {"n":"6","vb":false},
|
|
|
|
// * {"n":"22","vs":"U"},
|
|
|
|
// * {"n":"7","vs":"U"}]
|
|
|
|
// */
|
|
|
|
// callback.onSuccess(request, response);
|
|
|
|
// } catch (Exception e) {
|
|
|
|
// log.error("[{}] failed to process successful response [{}] ", registration.getEndpoint(), response, e);
|
|
|
|
// }
|
|
|
|
// });
|
|
|
|
// }, e -> {
|
|
|
|
// executor.submit(() -> {
|
|
|
|
// callback.onError(JacksonUtil.toString(request), e);
|
|
|
|
// });
|
|
|
|
// });
|
|
|
|
// } catch (Exception e) {
|
|
|
|
// callback.onError(JacksonUtil.toString(request), e);
|
|
|
|
// }
|
|
|
|
// }
|
|
|
|
|
|
|
|
private WriteRequest getWriteRequestSingleResource(ResourceModel.Type type, ContentFormat contentFormat, int objectId, int instanceId, int resourceId, Object value) { |
|
|
|
switch (type) { |
|
|
|
case STRING: // String
|
|
|
|
@ -308,14 +421,23 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
} |
|
|
|
|
|
|
|
private void validateVersionedId(LwM2mClient client, HasVersionedId request) { |
|
|
|
if (!client.isValidObjectVersion(request.getVersionedId())) { |
|
|
|
throw new IllegalArgumentException("Specified resource id is not configured in the device profile!"); |
|
|
|
} |
|
|
|
client.isValidObjectVersion(request.getVersionedId()); |
|
|
|
if (request.getObjectId() == null) { |
|
|
|
throw new IllegalArgumentException("Specified object id is null!"); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void validateVersionedIds(LwM2mClient client, HasVersionedIds request) { |
|
|
|
for (String versionedId : request.getVersionedIds()) { |
|
|
|
client.isValidObjectVersion(versionedId); |
|
|
|
} |
|
|
|
for (String objectId : request.getObjectIds()) { |
|
|
|
if (objectId == null) { |
|
|
|
throw new IllegalArgumentException("Specified object id is null!"); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private static <T> void addAttribute(List<Attribute> attributes, String attributeName, T value) { |
|
|
|
addAttribute(attributes, attributeName, value, null, null); |
|
|
|
} |
|
|
|
@ -347,7 +469,23 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im |
|
|
|
throw new CodecException("Invalid ResourceModel_Type for %s ContentFormat.", type); |
|
|
|
} |
|
|
|
|
|
|
|
private static ContentFormat getContentFormat(LwM2mClient client, HasContentFormat request) { |
|
|
|
return request.getContentFormat() != null ? request.getContentFormat() : client.getDefaultContentFormat(); |
|
|
|
private static ContentFormat getRequestContentFormat(LwM2mClient client, HasContentFormat request, LwM2mModelProvider modelProvider) { |
|
|
|
if (request.getRequestContentFormat() != null) { |
|
|
|
return request.getRequestContentFormat(); |
|
|
|
} else { |
|
|
|
String versionedId = null; |
|
|
|
if (request instanceof TbLwM2MReadRequest) { |
|
|
|
versionedId = ((TbLwM2MReadRequest) request).getVersionedId(); |
|
|
|
} else if (request instanceof TbLwM2MObserveRequest) { |
|
|
|
versionedId = ((TbLwM2MObserveRequest) request).getVersionedId(); |
|
|
|
} |
|
|
|
String id = fromVersionedIdToObjectId(versionedId); |
|
|
|
if (id != null && new LwM2mPath(id).isResource() && !client.isResourceMultiInstances(versionedId, modelProvider)) { |
|
|
|
return client.getDefaultContentFormat(); |
|
|
|
} |
|
|
|
else { |
|
|
|
return ContentFormat.DEFAULT; |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|