Browse Source
# Conflicts: # dao/src/test/java/org/thingsboard/server/dao/attributes/CachedAttributesServiceTest.javapull/4832/head
329 changed files with 13185 additions and 6610 deletions
@ -0,0 +1,77 @@ |
|||
/** |
|||
* 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.service.rpc; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.RpcId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.common.data.rpc.Rpc; |
|||
import org.thingsboard.server.common.data.rpc.RpcStatus; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.dao.rpc.RpcService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.queue.TbClusterService; |
|||
|
|||
@TbCoreComponent |
|||
@Service |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class TbRpcService { |
|||
private final RpcService rpcService; |
|||
private final TbClusterService tbClusterService; |
|||
|
|||
public Rpc save(TenantId tenantId, Rpc rpc) { |
|||
Rpc saved = rpcService.save(rpc); |
|||
pushRpcMsgToRuleEngine(tenantId, saved); |
|||
return saved; |
|||
} |
|||
|
|||
public void save(TenantId tenantId, RpcId rpcId, RpcStatus newStatus, JsonNode response) { |
|||
Rpc foundRpc = rpcService.findById(tenantId, rpcId); |
|||
if (foundRpc != null) { |
|||
foundRpc.setStatus(newStatus); |
|||
if (response != null) { |
|||
foundRpc.setResponse(response); |
|||
} |
|||
Rpc saved = rpcService.save(foundRpc); |
|||
pushRpcMsgToRuleEngine(tenantId, saved); |
|||
} else { |
|||
log.warn("[{}] Failed to update RPC status because RPC was already deleted", rpcId); |
|||
} |
|||
} |
|||
|
|||
private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) { |
|||
TbMsg msg = TbMsg.newMsg("RPC_" + rpc.getStatus().name(), rpc.getDeviceId(), TbMsgMetaData.EMPTY, JacksonUtil.toString(rpc)); |
|||
tbClusterService.pushMsgToRuleEngine(tenantId, rpc.getId(), msg, null); |
|||
} |
|||
|
|||
public Rpc findRpcById(TenantId tenantId, RpcId rpcId) { |
|||
return rpcService.findById(tenantId, rpcId); |
|||
} |
|||
|
|||
public PageData<Rpc> findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId deviceId, RpcStatus rpcStatus, PageLink pageLink) { |
|||
return rpcService.findAllByDeviceIdAndStatus(tenantId, deviceId, rpcStatus, pageLink); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,83 @@ |
|||
/** |
|||
* 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.service.ttl.rpc; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.scheduling.annotation.Scheduled; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; |
|||
import org.thingsboard.server.common.msg.queue.ServiceType; |
|||
import org.thingsboard.server.dao.rpc.RpcDao; |
|||
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; |
|||
import org.thingsboard.server.dao.tenant.TenantDao; |
|||
import org.thingsboard.server.queue.discovery.PartitionService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
import java.util.Date; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@TbCoreComponent |
|||
@Service |
|||
@Slf4j |
|||
@RequiredArgsConstructor |
|||
public class RpcCleanUpService { |
|||
@Value("${sql.ttl.rpc.enabled}") |
|||
private boolean ttlTaskExecutionEnabled; |
|||
|
|||
private final TenantDao tenantDao; |
|||
private final PartitionService partitionService; |
|||
private final TbTenantProfileCache tenantProfileCache; |
|||
private final RpcDao rpcDao; |
|||
|
|||
@Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.rpc.checking_interval})}", fixedDelayString = "${sql.ttl.rpc.checking_interval}") |
|||
public void cleanUp() { |
|||
if (ttlTaskExecutionEnabled) { |
|||
PageLink tenantsBatchRequest = new PageLink(10_000, 0); |
|||
PageData<TenantId> tenantsIds; |
|||
do { |
|||
tenantsIds = tenantDao.findTenantsIds(tenantsBatchRequest); |
|||
for (TenantId tenantId : tenantsIds.getData()) { |
|||
if (!partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) { |
|||
continue; |
|||
} |
|||
|
|||
Optional<DefaultTenantProfileConfiguration> tenantProfileConfiguration = tenantProfileCache.get(tenantId).getProfileConfiguration(); |
|||
if (tenantProfileConfiguration.isEmpty() || tenantProfileConfiguration.get().getRpcTtlDays() == 0) { |
|||
continue; |
|||
} |
|||
|
|||
long ttl = TimeUnit.DAYS.toMillis(tenantProfileConfiguration.get().getRpcTtlDays()); |
|||
long expirationTime = System.currentTimeMillis() - ttl; |
|||
|
|||
long totalRemoved = rpcDao.deleteOutdatedRpcByTenantId(tenantId, expirationTime); |
|||
|
|||
if (totalRemoved > 0) { |
|||
log.info("Removed {} outdated rpc(s) for tenant {} older than {}", totalRemoved, tenantId, new Date(expirationTime)); |
|||
} |
|||
} |
|||
|
|||
tenantsBatchRequest = tenantsBatchRequest.nextPageLink(); |
|||
} while (tenantsIds.hasNext()); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -1,44 +0,0 @@ |
|||
/** |
|||
* 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.service.ttl.timeseries; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.dao.model.ModelConstants; |
|||
import org.thingsboard.server.dao.util.PsqlDao; |
|||
import org.thingsboard.server.dao.util.SqlTsDao; |
|||
|
|||
import java.sql.Connection; |
|||
import java.sql.SQLException; |
|||
|
|||
@SqlTsDao |
|||
@PsqlDao |
|||
@Service |
|||
@Slf4j |
|||
public class PsqlTimeseriesCleanUpService extends AbstractTimeseriesCleanUpService { |
|||
|
|||
@Value("${sql.postgres.ts_key_value_partitioning}") |
|||
private String partitionType; |
|||
|
|||
@Override |
|||
protected void doCleanUp(Connection connection) throws SQLException { |
|||
long totalPartitionsRemoved = executeQuery(connection, "call drop_partitions_by_max_ttl('" + partitionType + "'," + systemTtl + ", 0);"); |
|||
log.info("Total partitions removed by TTL: [{}]", totalPartitionsRemoved); |
|||
long totalEntitiesTelemetryRemoved = executeQuery(connection, "call cleanup_timeseries_by_ttl('" + ModelConstants.NULL_UUID + "'," + systemTtl + ", 0);"); |
|||
log.info("Total telemetry removed stats by TTL for entities: [{}]", totalEntitiesTelemetryRemoved); |
|||
} |
|||
} |
|||
@ -1,36 +0,0 @@ |
|||
/** |
|||
* 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.service.ttl.timeseries; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.dao.model.ModelConstants; |
|||
import org.thingsboard.server.dao.util.TimescaleDBTsDao; |
|||
|
|||
import java.sql.Connection; |
|||
import java.sql.SQLException; |
|||
|
|||
@TimescaleDBTsDao |
|||
@Service |
|||
@Slf4j |
|||
public class TimescaleTimeseriesCleanUpService extends AbstractTimeseriesCleanUpService { |
|||
|
|||
@Override |
|||
protected void doCleanUp(Connection connection) throws SQLException { |
|||
long totalEntitiesTelemetryRemoved = executeQuery(connection, "call cleanup_timeseries_by_ttl('" + ModelConstants.NULL_UUID + "'," + systemTtl + ", 0);"); |
|||
log.info("Total telemetry removed stats by TTL for entities: [{}]", totalEntitiesTelemetryRemoved); |
|||
} |
|||
} |
|||
@ -0,0 +1,43 @@ |
|||
/** |
|||
* 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; |
|||
|
|||
import org.eclipse.leshan.client.object.Security; |
|||
import org.eclipse.leshan.core.util.Hex; |
|||
import org.junit.Test; |
|||
import org.thingsboard.server.common.data.device.credentials.lwm2m.PSKClientCredentials; |
|||
|
|||
import java.nio.charset.StandardCharsets; |
|||
|
|||
import static org.eclipse.leshan.client.object.Security.psk; |
|||
|
|||
public class PskLwm2mIntegrationTest extends AbstractLwM2MIntegrationTest { |
|||
|
|||
@Test |
|||
public void testConnectWithPSKAndObserveTelemetry() throws Exception { |
|||
String pskIdentity = "SOME_PSK_ID"; |
|||
String pskKey = "73656372657450534b"; |
|||
PSKClientCredentials clientCredentials = new PSKClientCredentials(); |
|||
clientCredentials.setEndpoint(ENDPOINT); |
|||
clientCredentials.setKey(pskKey); |
|||
clientCredentials.setIdentity(pskIdentity); |
|||
Security security = psk(SECURE_URI, |
|||
123, |
|||
pskIdentity.getBytes(StandardCharsets.UTF_8), |
|||
Hex.decodeHex(pskKey.toCharArray())); |
|||
super.basicTestConnectionObserveTelemetry(security, clientCredentials, SECURE_COAP_CONFIG, ENDPOINT); |
|||
} |
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* 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; |
|||
|
|||
import org.eclipse.leshan.client.object.Security; |
|||
import org.eclipse.leshan.core.util.Hex; |
|||
import org.junit.Test; |
|||
import org.thingsboard.server.common.data.device.credentials.lwm2m.RPKClientCredentials; |
|||
|
|||
import static org.eclipse.leshan.client.object.Security.rpk; |
|||
|
|||
public class RpkLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest { |
|||
@Test |
|||
public void testConnectWithRPKAndObserveTelemetry() throws Exception { |
|||
RPKClientCredentials rpkClientCredentials = new RPKClientCredentials(); |
|||
rpkClientCredentials.setEndpoint(ENDPOINT); |
|||
rpkClientCredentials.setKey(Hex.encodeHexString(clientPublicKey.getEncoded())); |
|||
Security security = rpk(SECURE_URI, |
|||
123, |
|||
clientPublicKey.getEncoded(), |
|||
clientPrivateKey.getEncoded(), |
|||
serverX509Cert.getPublicKey().getEncoded()); |
|||
super.basicTestConnectionObserveTelemetry(security, rpkClientCredentials, SECURE_COAP_CONFIG, ENDPOINT); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* 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.dao.rpc; |
|||
|
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.RpcId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.common.data.rpc.Rpc; |
|||
import org.thingsboard.server.common.data.rpc.RpcStatus; |
|||
|
|||
public interface RpcService { |
|||
Rpc save(Rpc rpc); |
|||
|
|||
void deleteRpc(TenantId tenantId, RpcId id); |
|||
|
|||
void deleteAllRpcByTenantId(TenantId tenantId); |
|||
|
|||
Rpc findById(TenantId tenantId, RpcId id); |
|||
|
|||
ListenableFuture<Rpc> findRpcByIdAsync(TenantId tenantId, RpcId id); |
|||
|
|||
PageData<Rpc> findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId deviceId, RpcStatus rpcStatus, PageLink pageLink); |
|||
} |
|||
@ -0,0 +1,20 @@ |
|||
/** |
|||
* 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.common.data.device.data; |
|||
|
|||
public enum PowerMode { |
|||
PSM, DRX, E_DRX |
|||
} |
|||
@ -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.common.data.device.data.lwm2m; |
|||
|
|||
import lombok.Data; |
|||
|
|||
import java.util.Map; |
|||
|
|||
@Data |
|||
public class BootstrapConfiguration { |
|||
|
|||
//TODO: define the objects;
|
|||
private Map<String, Object> servers; |
|||
private Map<String, Object> lwm2mServer; |
|||
private Map<String, Object> bootstrapServer; |
|||
|
|||
} |
|||
@ -0,0 +1,33 @@ |
|||
/** |
|||
* 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.common.data.device.data.lwm2m; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonInclude; |
|||
import lombok.Data; |
|||
|
|||
@Data |
|||
@JsonInclude(JsonInclude.Include.NON_NULL) |
|||
public class ObjectAttributes { |
|||
|
|||
private Long dim; |
|||
private String ver; |
|||
private Long pmin; |
|||
private Long pmax; |
|||
private Double gt; |
|||
private Double lt; |
|||
private Double st; |
|||
|
|||
} |
|||
@ -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.common.data.device.data.lwm2m; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.device.data.PowerMode; |
|||
|
|||
@Data |
|||
public class OtherConfiguration { |
|||
|
|||
private Integer fwUpdateStrategy; |
|||
private Integer swUpdateStrategy; |
|||
private Integer clientOnlyObserveAfterConnect; |
|||
private PowerMode powerMode; |
|||
private String fwUpdateResource; |
|||
private String swUpdateResource; |
|||
private boolean compositeOperationsSupport; |
|||
|
|||
} |
|||
@ -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.common.data.device.data.lwm2m; |
|||
|
|||
import lombok.Data; |
|||
|
|||
import java.util.Map; |
|||
import java.util.Set; |
|||
|
|||
@Data |
|||
public class TelemetryMappingConfiguration { |
|||
|
|||
private Map<String, String> keyName; |
|||
private Set<String> observe; |
|||
private Set<String> attribute; |
|||
private Set<String> telemetry; |
|||
private Map<String, ObjectAttributes> attributeLwm2m; |
|||
|
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* 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.common.data.id; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonCreator; |
|||
import com.fasterxml.jackson.annotation.JsonIgnore; |
|||
import com.fasterxml.jackson.annotation.JsonProperty; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
public final class RpcId extends UUIDBased implements EntityId { |
|||
|
|||
private static final long serialVersionUID = 1L; |
|||
|
|||
@JsonCreator |
|||
public RpcId(@JsonProperty("id") UUID id) { |
|||
super(id); |
|||
} |
|||
|
|||
@JsonIgnore |
|||
@Override |
|||
public EntityType getEntityType() { |
|||
return EntityType.RPC; |
|||
} |
|||
} |
|||
@ -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.common.data.rpc; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.Data; |
|||
import lombok.EqualsAndHashCode; |
|||
import org.thingsboard.server.common.data.BaseData; |
|||
import org.thingsboard.server.common.data.HasTenantId; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.RpcId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
|
|||
@Data |
|||
@EqualsAndHashCode(callSuper = true) |
|||
public class Rpc extends BaseData<RpcId> implements HasTenantId { |
|||
private TenantId tenantId; |
|||
private DeviceId deviceId; |
|||
private long expirationTime; |
|||
private JsonNode request; |
|||
private JsonNode response; |
|||
private RpcStatus status; |
|||
|
|||
public Rpc() { |
|||
super(); |
|||
} |
|||
|
|||
public Rpc(RpcId id) { |
|||
super(id); |
|||
} |
|||
|
|||
public Rpc(Rpc rpc) { |
|||
super(rpc); |
|||
this.tenantId = rpc.getTenantId(); |
|||
this.deviceId = rpc.getDeviceId(); |
|||
this.expirationTime = rpc.getExpirationTime(); |
|||
this.request = rpc.getRequest(); |
|||
this.response = rpc.getResponse(); |
|||
this.status = rpc.getStatus(); |
|||
} |
|||
} |
|||
@ -0,0 +1,20 @@ |
|||
/** |
|||
* 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.common.data.rpc; |
|||
|
|||
public enum RpcStatus { |
|||
QUEUED, DELIVERED, SUCCESSFUL, TIMEOUT, FAILED |
|||
} |
|||
@ -0,0 +1,211 @@ |
|||
/** |
|||
* 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.queue.common; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.junit.After; |
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.junit.runner.RunWith; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.MockitoJUnitRunner; |
|||
import org.thingsboard.server.queue.TbQueueAdmin; |
|||
import org.thingsboard.server.queue.TbQueueConsumer; |
|||
import org.thingsboard.server.queue.TbQueueMsg; |
|||
import org.thingsboard.server.queue.TbQueueProducer; |
|||
|
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.CountDownLatch; |
|||
import java.util.concurrent.ExecutorService; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.atomic.AtomicLong; |
|||
|
|||
import static org.hamcrest.Matchers.equalTo; |
|||
import static org.hamcrest.Matchers.greaterThanOrEqualTo; |
|||
import static org.hamcrest.Matchers.is; |
|||
import static org.hamcrest.Matchers.lessThan; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.anyLong; |
|||
import static org.mockito.BDDMockito.willAnswer; |
|||
import static org.mockito.BDDMockito.willDoNothing; |
|||
import static org.mockito.BDDMockito.willReturn; |
|||
import static org.mockito.Mockito.RETURNS_DEEP_STUBS; |
|||
import static org.mockito.Mockito.atLeastOnce; |
|||
import static org.mockito.Mockito.mock; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.spy; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
|
|||
import static org.hamcrest.MatcherAssert.assertThat; |
|||
import static org.mockito.hamcrest.MockitoHamcrest.longThat; |
|||
|
|||
@Slf4j |
|||
@RunWith(MockitoJUnitRunner.class) |
|||
public class DefaultTbQueueRequestTemplateTest { |
|||
|
|||
@Mock |
|||
TbQueueAdmin queueAdmin; |
|||
@Mock |
|||
TbQueueProducer<TbQueueMsg> requestTemplate; |
|||
@Mock |
|||
TbQueueConsumer<TbQueueMsg> responseTemplate; |
|||
@Mock |
|||
ExecutorService executorMock; |
|||
|
|||
ExecutorService executor; |
|||
String topic = "js-responses-tb-node-0"; |
|||
long maxRequestTimeout = 10; |
|||
long maxPendingRequests = 32; |
|||
long pollInterval = 5; |
|||
|
|||
DefaultTbQueueRequestTemplate inst; |
|||
|
|||
@Before |
|||
public void setUp() throws Exception { |
|||
willReturn(topic).given(responseTemplate).getTopic(); |
|||
inst = spy(new DefaultTbQueueRequestTemplate( |
|||
queueAdmin, requestTemplate, responseTemplate, |
|||
maxRequestTimeout, maxPendingRequests, pollInterval, executorMock)); |
|||
|
|||
} |
|||
|
|||
@After |
|||
public void tearDown() throws Exception { |
|||
if (executor != null) { |
|||
executor.shutdownNow(); |
|||
} |
|||
} |
|||
|
|||
@Test |
|||
public void givenInstance_whenVerifyInitialParameters_thenOK() { |
|||
assertThat(inst.maxPendingRequests, equalTo(maxPendingRequests)); |
|||
assertThat(inst.maxRequestTimeoutNs, equalTo(TimeUnit.MILLISECONDS.toNanos(maxRequestTimeout))); |
|||
assertThat(inst.pollInterval, equalTo(pollInterval)); |
|||
assertThat(inst.executor, is(executorMock)); |
|||
assertThat(inst.stopped, is(false)); |
|||
assertThat(inst.internalExecutor, is(false)); |
|||
} |
|||
|
|||
@Test |
|||
public void givenExternalExecutor_whenInitStop_thenOK() { |
|||
inst.init(); |
|||
assertThat(inst.nextCleanupNs, equalTo(0L)); |
|||
verify(queueAdmin, times(1)).createTopicIfNotExists(topic); |
|||
verify(requestTemplate, times(1)).init(); |
|||
verify(responseTemplate, times(1)).subscribe(); |
|||
verify(executorMock, times(1)).submit(any(Runnable.class)); |
|||
|
|||
inst.stop(); |
|||
assertThat(inst.stopped, is(true)); |
|||
verify(responseTemplate, times(1)).unsubscribe(); |
|||
verify(requestTemplate, times(1)).stop(); |
|||
verify(executorMock, never()).shutdownNow(); |
|||
} |
|||
|
|||
@Test |
|||
public void givenMainLoop_whenLoopFewTimes_thenVerifyInvocationCount() throws InterruptedException { |
|||
executor = inst.createExecutor(); |
|||
CountDownLatch latch = new CountDownLatch(5); |
|||
willDoNothing().given(inst).sleep(anyLong()); |
|||
willAnswer(invocation -> { |
|||
if (latch.getCount() == 1) { |
|||
inst.stop(); //stop the loop in natural way
|
|||
} |
|||
if (latch.getCount() == 3 || latch.getCount() == 4) { |
|||
latch.countDown(); |
|||
throw new RuntimeException("test catch block"); |
|||
} |
|||
latch.countDown(); |
|||
return null; |
|||
}).given(inst).fetchAndProcessResponses(); |
|||
|
|||
executor.submit(inst::mainLoop); |
|||
latch.await(10, TimeUnit.SECONDS); |
|||
|
|||
verify(inst, times(5)).fetchAndProcessResponses(); |
|||
verify(inst, times(2)).sleep(longThat(lessThan(TimeUnit.MILLISECONDS.toNanos(inst.pollInterval)))); |
|||
} |
|||
|
|||
@Test |
|||
public void givenMessages_whenSend_thenOK() { |
|||
willDoNothing().given(inst).sendToRequestTemplate(any(), any(), any(), any()); |
|||
inst.init(); |
|||
final int msgCount = 10; |
|||
for (int i = 0; i < msgCount; i++) { |
|||
inst.send(getRequestMsgMock()); |
|||
} |
|||
assertThat(inst.pendingRequests.mappingCount(), equalTo((long) msgCount)); |
|||
verify(inst, times(msgCount)).sendToRequestTemplate(any(), any(), any(), any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenMessagesOverMaxPendingRequests_whenSend_thenImmediateFailedFutureForTheOfRequests() { |
|||
willDoNothing().given(inst).sendToRequestTemplate(any(), any(), any(), any()); |
|||
inst.init(); |
|||
int msgOverflowCount = 10; |
|||
for (int i = 0; i < inst.maxPendingRequests; i++) { |
|||
assertThat(inst.send(getRequestMsgMock()).isDone(), is(false)); //SettableFuture future - pending only
|
|||
} |
|||
for (int i = 0; i < msgOverflowCount; i++) { |
|||
assertThat("max pending requests overflow", inst.send(getRequestMsgMock()).isDone(), is(true)); //overflow, immediate failed future
|
|||
} |
|||
assertThat(inst.pendingRequests.mappingCount(), equalTo(inst.maxPendingRequests)); |
|||
verify(inst, times((int) inst.maxPendingRequests)).sendToRequestTemplate(any(), any(), any(), any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenNothing_whenSendAndFetchAndProcessResponsesWithTimeout_thenFail() { |
|||
//given
|
|||
AtomicLong currentTime = new AtomicLong(); |
|||
willAnswer(x -> { |
|||
log.info("currentTime={}", currentTime.get()); |
|||
return currentTime.get(); |
|||
}).given(inst).getCurrentClockNs(); |
|||
inst.init(); |
|||
inst.setupNextCleanup(); |
|||
willReturn(Collections.emptyList()).given(inst).doPoll(); |
|||
|
|||
//when
|
|||
long stepNs = TimeUnit.MILLISECONDS.toNanos(1); |
|||
for (long i = 0; i <= inst.maxRequestTimeoutNs * 2; i = i + stepNs) { |
|||
currentTime.addAndGet(stepNs); |
|||
assertThat(inst.send(getRequestMsgMock()).isDone(), is(false)); //SettableFuture future - pending only
|
|||
if (i % (inst.maxRequestTimeoutNs * 3 / 2) == 0) { |
|||
inst.fetchAndProcessResponses(); |
|||
} |
|||
} |
|||
|
|||
//then
|
|||
ArgumentCaptor<DefaultTbQueueRequestTemplate.ResponseMetaData> argumentCaptorResp = ArgumentCaptor.forClass(DefaultTbQueueRequestTemplate.ResponseMetaData.class); |
|||
ArgumentCaptor<UUID> argumentCaptorUUID = ArgumentCaptor.forClass(UUID.class); |
|||
ArgumentCaptor<Long> argumentCaptorLong = ArgumentCaptor.forClass(Long.class); |
|||
verify(inst, atLeastOnce()).setTimeoutException(argumentCaptorUUID.capture(), argumentCaptorResp.capture(), argumentCaptorLong.capture()); |
|||
|
|||
List<DefaultTbQueueRequestTemplate.ResponseMetaData> responseMetaDataList = argumentCaptorResp.getAllValues(); |
|||
List<Long> tickTsList = argumentCaptorLong.getAllValues(); |
|||
for (int i = 0; i < responseMetaDataList.size(); i++) { |
|||
assertThat("tickTs >= calculatedExpTime", tickTsList.get(i), greaterThanOrEqualTo(responseMetaDataList.get(i).getSubmitTime() + responseMetaDataList.get(i).getTimeout())); |
|||
} |
|||
} |
|||
|
|||
TbQueueMsg getRequestMsgMock() { |
|||
return mock(TbQueueMsg.class, RETURNS_DEEP_STUBS); |
|||
} |
|||
} |
|||
@ -0,0 +1,150 @@ |
|||
/** |
|||
* 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.coap; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.eclipse.californium.core.coap.CoAP; |
|||
import org.eclipse.californium.core.coap.Request; |
|||
import org.eclipse.californium.core.coap.Response; |
|||
import org.eclipse.californium.core.network.Exchange; |
|||
import org.eclipse.californium.core.observe.ObserveRelation; |
|||
import org.eclipse.californium.core.server.resources.CoapExchange; |
|||
import org.eclipse.californium.core.server.resources.Resource; |
|||
import org.eclipse.californium.core.server.resources.ResourceObserver; |
|||
import org.thingsboard.common.util.ThingsBoardExecutors; |
|||
import org.thingsboard.server.common.data.DeviceTransportType; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.ota.OtaPackageType; |
|||
import org.thingsboard.server.common.data.security.DeviceTokenCredentials; |
|||
import org.thingsboard.server.common.transport.TransportServiceCallback; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
|
|||
import java.util.List; |
|||
import java.util.Optional; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.ExecutorService; |
|||
|
|||
@Slf4j |
|||
public class OtaPackageTransportResource extends AbstractCoapTransportResource { |
|||
private static final int ACCESS_TOKEN_POSITION = 2; |
|||
|
|||
private final OtaPackageType otaPackageType; |
|||
|
|||
public OtaPackageTransportResource(CoapTransportContext ctx, OtaPackageType otaPackageType) { |
|||
super(ctx, otaPackageType.getKeyPrefix()); |
|||
this.otaPackageType = otaPackageType; |
|||
|
|||
this.setObservable(true); |
|||
} |
|||
|
|||
@Override |
|||
protected void processHandleGet(CoapExchange exchange) { |
|||
log.trace("Processing {}", exchange.advanced().getRequest()); |
|||
exchange.accept(); |
|||
Exchange advanced = exchange.advanced(); |
|||
Request request = advanced.getRequest(); |
|||
processAccessTokenRequest(exchange, request); |
|||
} |
|||
|
|||
@Override |
|||
protected void processHandlePost(CoapExchange exchange) { |
|||
exchange.respond(CoAP.ResponseCode.METHOD_NOT_ALLOWED); |
|||
} |
|||
|
|||
private void processAccessTokenRequest(CoapExchange exchange, Request request) { |
|||
Optional<DeviceTokenCredentials> credentials = decodeCredentials(request); |
|||
if (credentials.isEmpty()) { |
|||
exchange.respond(CoAP.ResponseCode.UNAUTHORIZED); |
|||
return; |
|||
} |
|||
transportService.process(DeviceTransportType.COAP, TransportProtos.ValidateDeviceTokenRequestMsg.newBuilder().setToken(credentials.get().getCredentialsId()).build(), |
|||
new CoapDeviceAuthCallback(transportContext, exchange, (sessionInfo, deviceProfile) -> { |
|||
getOtaPackageCallback(sessionInfo, exchange, otaPackageType); |
|||
})); |
|||
} |
|||
|
|||
private void getOtaPackageCallback(TransportProtos.SessionInfoProto sessionInfo, CoapExchange exchange, OtaPackageType firmwareType) { |
|||
TransportProtos.GetOtaPackageRequestMsg requestMsg = TransportProtos.GetOtaPackageRequestMsg.newBuilder() |
|||
.setTenantIdMSB(sessionInfo.getTenantIdMSB()) |
|||
.setTenantIdLSB(sessionInfo.getTenantIdLSB()) |
|||
.setDeviceIdMSB(sessionInfo.getDeviceIdMSB()) |
|||
.setDeviceIdLSB(sessionInfo.getDeviceIdLSB()) |
|||
.setType(firmwareType.name()).build(); |
|||
transportContext.getTransportService().process(sessionInfo, requestMsg, new OtaPackageCallback(exchange)); |
|||
} |
|||
|
|||
private Optional<DeviceTokenCredentials> decodeCredentials(Request request) { |
|||
List<String> uriPath = request.getOptions().getUriPath(); |
|||
if (uriPath.size() == ACCESS_TOKEN_POSITION) { |
|||
return Optional.of(new DeviceTokenCredentials(uriPath.get(ACCESS_TOKEN_POSITION - 1))); |
|||
} else { |
|||
return Optional.empty(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Resource getChild(String name) { |
|||
return this; |
|||
} |
|||
|
|||
private class OtaPackageCallback implements TransportServiceCallback<TransportProtos.GetOtaPackageResponseMsg> { |
|||
private final CoapExchange exchange; |
|||
|
|||
OtaPackageCallback(CoapExchange exchange) { |
|||
this.exchange = exchange; |
|||
} |
|||
|
|||
@Override |
|||
public void onSuccess(TransportProtos.GetOtaPackageResponseMsg msg) { |
|||
String title = exchange.getQueryParameter("title"); |
|||
String version = exchange.getQueryParameter("version"); |
|||
if (msg.getResponseStatus().equals(TransportProtos.ResponseStatus.SUCCESS)) { |
|||
String firmwareId = new UUID(msg.getOtaPackageIdMSB(), msg.getOtaPackageIdLSB()).toString(); |
|||
if ((title == null || msg.getTitle().equals(title)) && (version == null || msg.getVersion().equals(version))) { |
|||
String strChunkSize = exchange.getQueryParameter("size"); |
|||
String strChunk = exchange.getQueryParameter("chunk"); |
|||
int chunkSize = StringUtils.isEmpty(strChunkSize) ? 0 : Integer.parseInt(strChunkSize); |
|||
int chunk = StringUtils.isEmpty(strChunk) ? 0 : Integer.parseInt(strChunk); |
|||
respondOtaPackage(exchange, transportContext.getOtaPackageDataCache().get(firmwareId, chunkSize, chunk)); |
|||
} else { |
|||
exchange.respond(CoAP.ResponseCode.BAD_REQUEST); |
|||
} |
|||
} else { |
|||
exchange.respond(CoAP.ResponseCode.NOT_FOUND); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onError(Throwable e) { |
|||
log.warn("Failed to process request", e); |
|||
exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); |
|||
} |
|||
} |
|||
|
|||
private void respondOtaPackage(CoapExchange exchange, byte[] data) { |
|||
Response response = new Response(CoAP.ResponseCode.CONTENT); |
|||
if (data != null && data.length > 0) { |
|||
response.setPayload(data); |
|||
if (exchange.getRequestOptions().getBlock2() != null) { |
|||
int chunkSize = exchange.getRequestOptions().getBlock2().getSzx(); |
|||
boolean lastFlag = data.length <= chunkSize; |
|||
response.getOptions().setBlock2(chunkSize, lastFlag, 0); |
|||
} |
|||
transportContext.getExecutor().submit(() -> exchange.respond(response)); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,66 @@ |
|||
/** |
|||
* 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.config; |
|||
|
|||
import lombok.Getter; |
|||
import org.eclipse.leshan.core.LwM2m.Version; |
|||
import org.eclipse.leshan.core.request.ContentFormat; |
|||
|
|||
public enum LwM2mVersion { |
|||
VERSION_1_0(0, Version.V1_0, ContentFormat.TLV), |
|||
VERSION_1_1(1, Version.V1_1, ContentFormat.TEXT); |
|||
|
|||
@Getter |
|||
private final int code; |
|||
@Getter |
|||
private final Version version; |
|||
@Getter |
|||
private final ContentFormat contentFormat; |
|||
|
|||
LwM2mVersion(int code, Version version, ContentFormat contentFormat) { |
|||
this.code = code; |
|||
this.version = version; |
|||
this.contentFormat = contentFormat; |
|||
} |
|||
|
|||
public static LwM2mVersion fromVersion(Version version) { |
|||
for (LwM2mVersion to : LwM2mVersion.values()) { |
|||
if (to.version.equals(version)) { |
|||
return to; |
|||
} |
|||
} |
|||
throw new IllegalArgumentException(String.format("Unsupported typeLwM2mVersion type : %s", version)); |
|||
} |
|||
|
|||
public static LwM2mVersion fromVersionStr(String versionStr) { |
|||
for (LwM2mVersion to : LwM2mVersion.values()) { |
|||
if (to.version.toString().equals(versionStr)) { |
|||
return to; |
|||
} |
|||
} |
|||
throw new IllegalArgumentException(String.format("Unsupported contentFormatLwM2mVersion version : %s", versionStr)); |
|||
} |
|||
|
|||
public static LwM2mVersion fromCode(int code) { |
|||
for (LwM2mVersion to : LwM2mVersion.values()) { |
|||
if (to.code == code) { |
|||
return to; |
|||
} |
|||
} |
|||
throw new IllegalArgumentException(String.format("Unsupported codeLwM2mVersion code : %d", code)); |
|||
} |
|||
} |
|||
|
|||
File diff suppressed because it is too large
@ -0,0 +1,94 @@ |
|||
/** |
|||
* 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; |
|||
|
|||
import lombok.Getter; |
|||
|
|||
/** |
|||
* Define the behavior of a write request. |
|||
*/ |
|||
public enum LwM2mOperationType { |
|||
|
|||
READ(0, "Read", true), |
|||
READ_COMPOSITE(1, "ReadComposite", false, true), |
|||
DISCOVER(2, "Discover", true), |
|||
DISCOVER_ALL(3, "DiscoverAll", false), |
|||
OBSERVE_READ_ALL(4, "ObserveReadAll", false), |
|||
|
|||
OBSERVE(5, "Observe", true), |
|||
OBSERVE_COMPOSITE(6, "ObserveComposite", false, true), |
|||
OBSERVE_CANCEL(7, "ObserveCancel", true), |
|||
OBSERVE_COMPOSITE_CANCEL(8, "ObserveCompositeCancel", false, true), |
|||
OBSERVE_CANCEL_ALL(9, "ObserveCancelAll", false), |
|||
EXECUTE(10, "Execute", true), |
|||
/** |
|||
* Replaces the Object Instance or the Resource(s) with the new value provided in the “Write” operation. (see |
|||
* section 5.3.3 of the LW M2M spec). |
|||
* if all resources are to be replaced |
|||
*/ |
|||
WRITE_REPLACE(11, "WriteReplace", true), |
|||
|
|||
/** |
|||
* Adds or updates Resources provided in the new value and leaves other existing Resources unchanged. (see section |
|||
* 5.3.3 of the LW M2M spec). |
|||
* if this is a partial update request |
|||
*/ |
|||
WRITE_UPDATE(12, "WriteUpdate", true), |
|||
WRITE_COMPOSITE(14, "WriteComposite", false, true), |
|||
WRITE_ATTRIBUTES(15, "WriteAttributes", true), |
|||
DELETE(16, "Delete", true), |
|||
|
|||
// only for RPC
|
|||
FW_UPDATE(17, "FirmwareUpdate", false); |
|||
|
|||
// FW_READ_INFO(18, "FirmwareReadInfo"),
|
|||
// SW_READ_INFO(19, "SoftwareReadInfo"),
|
|||
// SW_UPDATE(20, "SoftwareUpdate"),
|
|||
// SW_UNINSTALL(21, "SoftwareUninstall");
|
|||
|
|||
@Getter |
|||
private final int code; |
|||
@Getter |
|||
private final String type; |
|||
@Getter |
|||
private final boolean hasObjectId; |
|||
|
|||
@Getter |
|||
private final boolean composite; |
|||
|
|||
LwM2mOperationType(int code, String type, boolean hasObjectId) { |
|||
this(code, type, hasObjectId, false); |
|||
} |
|||
|
|||
LwM2mOperationType(int code, String type, boolean hasObjectId, boolean composite) { |
|||
this.code = code; |
|||
this.type = type; |
|||
this.hasObjectId = hasObjectId; |
|||
this.composite = composite; |
|||
if(hasObjectId && composite){ |
|||
throw new IllegalArgumentException("Can't set both Composite and hasObjectId for the same operation!"); |
|||
} |
|||
} |
|||
|
|||
public static LwM2mOperationType fromType(String type) { |
|||
for (LwM2mOperationType to : LwM2mOperationType.values()) { |
|||
if (to.type.equals(type)) { |
|||
return to; |
|||
} |
|||
} |
|||
return null; |
|||
} |
|||
} |
|||
@ -1,613 +0,0 @@ |
|||
/** |
|||
* 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; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.eclipse.californium.core.coap.CoAP; |
|||
import org.eclipse.californium.core.coap.Response; |
|||
import org.eclipse.leshan.core.Link; |
|||
import org.eclipse.leshan.core.model.ResourceModel; |
|||
import org.eclipse.leshan.core.node.LwM2mNode; |
|||
import org.eclipse.leshan.core.node.LwM2mObject; |
|||
import org.eclipse.leshan.core.node.LwM2mObjectInstance; |
|||
import org.eclipse.leshan.core.node.LwM2mPath; |
|||
import org.eclipse.leshan.core.node.LwM2mResource; |
|||
import org.eclipse.leshan.core.node.LwM2mSingleResource; |
|||
import org.eclipse.leshan.core.node.ObjectLink; |
|||
import org.eclipse.leshan.core.observation.Observation; |
|||
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.ExecuteRequest; |
|||
import org.eclipse.leshan.core.request.ObserveRequest; |
|||
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.WriteRequest; |
|||
import org.eclipse.leshan.core.request.exception.ClientSleepingException; |
|||
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.ReadResponse; |
|||
import org.eclipse.leshan.core.response.ResponseCallback; |
|||
import org.eclipse.leshan.core.response.WriteAttributesResponse; |
|||
import org.eclipse.leshan.core.response.WriteResponse; |
|||
import org.eclipse.leshan.core.util.Hex; |
|||
import org.eclipse.leshan.core.util.NamedThreadFactory; |
|||
import org.eclipse.leshan.server.registration.Registration; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; |
|||
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientRpcRequest; |
|||
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import java.util.Arrays; |
|||
import java.util.Collection; |
|||
import java.util.Date; |
|||
import java.util.Set; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.ExecutorService; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.stream.Collectors; |
|||
|
|||
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.CONTENT; |
|||
import static org.eclipse.leshan.core.ResponseCode.BAD_REQUEST; |
|||
import static org.eclipse.leshan.core.ResponseCode.NOT_FOUND; |
|||
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.transport.lwm2m.server.LwM2mTransportServerHelper.getContentFormatByResourceModelType; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.DEFAULT_TIMEOUT; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_PACKAGE_5_ID; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_UPDATE_ID; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_ERROR; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_INFO; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_VALUE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.DISCOVER; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.EXECUTE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_READ_ALL; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_ATTRIBUTES; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_REPLACE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_UPDATE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.RESPONSE_REQUEST_CHANNEL; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SW_INSTALL_ID; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SW_PACKAGE_ID; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromIdVerToObjectId; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromObjectIdToIdVer; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.createWriteAttributeRequest; |
|||
|
|||
@Slf4j |
|||
@Service |
|||
@TbLwM2mTransportComponent |
|||
@RequiredArgsConstructor |
|||
public class LwM2mTransportRequest { |
|||
private ExecutorService responseRequestExecutor; |
|||
|
|||
public LwM2mValueConverterImpl converter; |
|||
|
|||
private final LwM2mTransportContext context; |
|||
private final LwM2MTransportServerConfig config; |
|||
private final LwM2mClientContext lwM2mClientContext; |
|||
private final DefaultLwM2MTransportMsgHandler handler; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
this.converter = LwM2mValueConverterImpl.getInstance(); |
|||
responseRequestExecutor = Executors.newFixedThreadPool(this.config.getResponsePoolSize(), |
|||
new NamedThreadFactory(String.format("LwM2M %s channel response after request", RESPONSE_REQUEST_CHANNEL))); |
|||
} |
|||
|
|||
public void sendAllRequest(LwM2mClient lwM2MClient, String targetIdVer, LwM2mTypeOper typeOper, Object params, long timeoutInMs, LwM2mClientRpcRequest lwm2mClientRpcRequest) { |
|||
sendAllRequest(lwM2MClient, targetIdVer, typeOper, lwM2MClient.getDefaultContentFormat(), params, timeoutInMs, lwm2mClientRpcRequest); |
|||
} |
|||
|
|||
|
|||
public void sendAllRequest(LwM2mClient lwM2MClient, String targetIdVer, LwM2mTypeOper typeOper, |
|||
ContentFormat contentFormat, Object params, long timeoutInMs, LwM2mClientRpcRequest lwm2mClientRpcRequest) { |
|||
Registration registration = lwM2MClient.getRegistration(); |
|||
try { |
|||
String target = convertPathFromIdVerToObjectId(targetIdVer); |
|||
if(contentFormat == null){ |
|||
contentFormat = ContentFormat.DEFAULT; |
|||
} |
|||
LwM2mPath resultIds = target != null ? new LwM2mPath(target) : null; |
|||
if (!OBSERVE_CANCEL.name().equals(typeOper.name()) && resultIds != null && registration != null && resultIds.getObjectId() >= 0 && lwM2MClient != null) { |
|||
if (lwM2MClient.isValidObjectVersion(targetIdVer)) { |
|||
timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT; |
|||
SimpleDownlinkRequest request = createRequest(registration, lwM2MClient, typeOper, contentFormat, target, |
|||
targetIdVer, resultIds, params, lwm2mClientRpcRequest); |
|||
if (request != null) { |
|||
try { |
|||
this.sendRequest(registration, lwM2MClient, request, timeoutInMs, lwm2mClientRpcRequest); |
|||
} catch (ClientSleepingException e) { |
|||
SimpleDownlinkRequest finalRequest = request; |
|||
long finalTimeoutInMs = timeoutInMs; |
|||
LwM2mClientRpcRequest finalRpcRequest = lwm2mClientRpcRequest; |
|||
lwM2MClient.getQueuedRequests().add(() -> sendRequest(registration, lwM2MClient, finalRequest, finalTimeoutInMs, finalRpcRequest)); |
|||
} catch (Exception e) { |
|||
log.error("[{}] [{}] [{}] Failed to send downlink.", registration.getEndpoint(), targetIdVer, typeOper.name(), e); |
|||
} |
|||
} else if (WRITE_UPDATE.name().equals(typeOper.name())) { |
|||
if (lwm2mClientRpcRequest != null) { |
|||
String errorMsg = String.format("Path %s params is not valid", targetIdVer); |
|||
handler.sentRpcResponse(lwm2mClientRpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR); |
|||
} |
|||
} else if (WRITE_REPLACE.name().equals(typeOper.name()) || EXECUTE.name().equals(typeOper.name())) { |
|||
if (lwm2mClientRpcRequest != null) { |
|||
String errorMsg = String.format("Path %s object model is absent", targetIdVer); |
|||
handler.sentRpcResponse(lwm2mClientRpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR); |
|||
} |
|||
} else if (!OBSERVE_CANCEL.name().equals(typeOper.name())) { |
|||
log.error("[{}], [{}] - [{}] error SendRequest", registration.getEndpoint(), typeOper.name(), targetIdVer); |
|||
if (lwm2mClientRpcRequest != null) { |
|||
ResourceModel resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.config.getModelProvider()); |
|||
String errorMsg = resourceModel == null ? String.format("Path %s not found in object version", targetIdVer) : "SendRequest - null"; |
|||
handler.sentRpcResponse(lwm2mClientRpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR); |
|||
} |
|||
} |
|||
} else if (lwm2mClientRpcRequest != null) { |
|||
String errorMsg = String.format("Path %s not found in object version", targetIdVer); |
|||
handler.sentRpcResponse(lwm2mClientRpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR); |
|||
} |
|||
} else { |
|||
switch (typeOper) { |
|||
case OBSERVE_READ_ALL: |
|||
case DISCOVER_ALL: |
|||
Set<String> paths; |
|||
if (OBSERVE_READ_ALL.name().equals(typeOper.name())) { |
|||
Set<Observation> observations = context.getServer().getObservationService().getObservations(registration); |
|||
paths = observations.stream().map(observation -> observation.getPath().toString()).collect(Collectors.toUnmodifiableSet()); |
|||
} else { |
|||
assert registration != null; |
|||
Link[] objectLinks = registration.getSortedObjectLinks(); |
|||
paths = Arrays.stream(objectLinks).map(Link::toString).collect(Collectors.toUnmodifiableSet()); |
|||
} |
|||
String msg = String.format("%s: type operation %s paths - %s", LOG_LW2M_INFO, |
|||
typeOper.name(), paths); |
|||
this.handler.sendLogsToThingsboard(lwM2MClient, msg); |
|||
if (lwm2mClientRpcRequest != null) { |
|||
String valueMsg = String.format("Paths - %s", paths); |
|||
handler.sentRpcResponse(lwm2mClientRpcRequest, CONTENT.name(), valueMsg, LOG_LW2M_VALUE); |
|||
} |
|||
break; |
|||
case OBSERVE_CANCEL: |
|||
case OBSERVE_CANCEL_ALL: |
|||
int observeCancelCnt = 0; |
|||
String observeCancelMsg = null; |
|||
if (OBSERVE_CANCEL.name().equals(typeOper)) { |
|||
observeCancelCnt = context.getServer().getObservationService().cancelObservations(registration, target); |
|||
observeCancelMsg = String.format("%s: type operation %s paths: %s count: %d", LOG_LW2M_INFO, |
|||
OBSERVE_CANCEL.name(), target, observeCancelCnt); |
|||
} else { |
|||
observeCancelCnt = context.getServer().getObservationService().cancelObservations(registration); |
|||
observeCancelMsg = String.format("%s: type operation %s paths: All count: %d", LOG_LW2M_INFO, |
|||
OBSERVE_CANCEL.name(), observeCancelCnt); |
|||
} |
|||
this.afterObserveCancel(lwM2MClient, observeCancelCnt, observeCancelMsg, lwm2mClientRpcRequest); |
|||
break; |
|||
// lwm2mClientRpcRequest != null
|
|||
case FW_UPDATE: |
|||
handler.getInfoFirmwareUpdate(lwM2MClient, lwm2mClientRpcRequest); |
|||
break; |
|||
} |
|||
} |
|||
} catch (Exception e) { |
|||
String msg = String.format("%s: type operation %s %s", LOG_LW2M_ERROR, |
|||
typeOper.name(), e.getMessage()); |
|||
handler.sendLogsToThingsboard(lwM2MClient, msg); |
|||
if (lwm2mClientRpcRequest != null) { |
|||
String errorMsg = String.format("Path %s type operation %s %s", targetIdVer, typeOper.name(), e.getMessage()); |
|||
handler.sentRpcResponse(lwm2mClientRpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private SimpleDownlinkRequest createRequest(Registration registration, LwM2mClient lwM2MClient, LwM2mTypeOper typeOper, |
|||
ContentFormat contentFormat, String target, String targetIdVer, |
|||
LwM2mPath resultIds, Object params, LwM2mClientRpcRequest rpcRequest) { |
|||
SimpleDownlinkRequest request = null; |
|||
switch (typeOper) { |
|||
case READ: |
|||
request = new ReadRequest(contentFormat, target); |
|||
break; |
|||
case DISCOVER: |
|||
request = new DiscoverRequest(target); |
|||
break; |
|||
case OBSERVE: |
|||
String msg = String.format("%s: Send Observation %s.", LOG_LW2M_INFO, targetIdVer); |
|||
log.warn(msg); |
|||
if (resultIds.isResource()) { |
|||
Set<Observation> observations = context.getServer().getObservationService().getObservations(registration); |
|||
Set<Observation> paths = observations.stream().filter(observation -> observation.getPath().equals(resultIds)).collect(Collectors.toSet()); |
|||
if (paths.size() == 0) { |
|||
request = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId()); |
|||
} else { |
|||
request = new ReadRequest(contentFormat, target); |
|||
} |
|||
} else if (resultIds.isObjectInstance()) { |
|||
request = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId()); |
|||
} else if (resultIds.getObjectId() >= 0) { |
|||
request = new ObserveRequest(contentFormat, resultIds.getObjectId()); |
|||
} |
|||
break; |
|||
case EXECUTE: |
|||
ResourceModel resourceModelExecute = lwM2MClient.getResourceModel(targetIdVer, this.config.getModelProvider()); |
|||
if (resourceModelExecute != null) { |
|||
if (params != null && !resourceModelExecute.multiple) { |
|||
request = new ExecuteRequest(target, (String) this.converter.convertValue(params, resourceModelExecute.type, ResourceModel.Type.STRING, resultIds)); |
|||
} else { |
|||
request = new ExecuteRequest(target); |
|||
} |
|||
} |
|||
break; |
|||
case WRITE_REPLACE: |
|||
/** |
|||
* Request to write a <b>String Single-Instance Resource</b> using the TLV content format. |
|||
* Type from resourceModel -> STRING, INTEGER, FLOAT, BOOLEAN, OPAQUE, TIME, OBJLNK |
|||
* contentFormat -> TLV, TLV, TLV, TLV, OPAQUE, TLV, LINK |
|||
* JSON, TEXT; |
|||
**/ |
|||
ResourceModel resourceModelWrite = lwM2MClient.getResourceModel(targetIdVer, this.config.getModelProvider()); |
|||
if (resourceModelWrite != null) { |
|||
contentFormat = getContentFormatByResourceModelType(resourceModelWrite, contentFormat); |
|||
request = this.getWriteRequestSingleResource(contentFormat, resultIds.getObjectId(), |
|||
resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resourceModelWrite.type, |
|||
lwM2MClient, rpcRequest); |
|||
} |
|||
break; |
|||
case WRITE_UPDATE: |
|||
if (resultIds.isResource()) { |
|||
/** |
|||
* send request: path = '/3/0' node == wM2mObjectInstance |
|||
* with params == "\"resources\": {15: resource:{id:15. value:'+01'...}} |
|||
**/ |
|||
Collection<LwM2mResource> resources = lwM2MClient.getNewResourceForInstance( |
|||
targetIdVer, params, |
|||
this.config.getModelProvider(), |
|||
this.converter); |
|||
contentFormat = getContentFormatByResourceModelType(lwM2MClient.getResourceModel(targetIdVer, this.config.getModelProvider()), |
|||
contentFormat); |
|||
request = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, resultIds.getObjectId(), |
|||
resultIds.getObjectInstanceId(), resources); |
|||
} |
|||
/** |
|||
* params = "{\"id\":0,\"resources\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}" |
|||
* int rscId = resultIds.getObjectInstanceId(); |
|||
* contentFormat – Format of the payload (TLV or JSON). |
|||
*/ |
|||
else if (resultIds.isObjectInstance()) { |
|||
if (((ConcurrentHashMap) params).size() > 0) { |
|||
Collection<LwM2mResource> resources = lwM2MClient.getNewResourcesForInstance( |
|||
targetIdVer, params, |
|||
this.config.getModelProvider(), |
|||
this.converter); |
|||
if (resources.size() > 0) { |
|||
contentFormat = contentFormat.equals(ContentFormat.JSON) ? contentFormat : ContentFormat.TLV; |
|||
request = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, resultIds.getObjectId(), |
|||
resultIds.getObjectInstanceId(), resources); |
|||
} |
|||
} |
|||
} else if (resultIds.getObjectId() >= 0) { |
|||
request = new ObserveRequest(resultIds.getObjectId()); |
|||
} |
|||
break; |
|||
case WRITE_ATTRIBUTES: |
|||
request = createWriteAttributeRequest(target, params, this.handler); |
|||
break; |
|||
case DELETE: |
|||
request = new DeleteRequest(target); |
|||
break; |
|||
} |
|||
return request; |
|||
} |
|||
|
|||
/** |
|||
* @param registration - |
|||
* @param request - |
|||
* @param timeoutInMs - |
|||
*/ |
|||
|
|||
@SuppressWarnings({"error sendRequest"}) |
|||
private void sendRequest(Registration registration, LwM2mClient lwM2MClient, SimpleDownlinkRequest request, |
|||
long timeoutInMs, LwM2mClientRpcRequest rpcRequest) { |
|||
context.getServer().send(registration, request, timeoutInMs, (ResponseCallback<?>) response -> { |
|||
|
|||
if (!lwM2MClient.isInit()) { |
|||
lwM2MClient.initReadValue(this.handler, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration)); |
|||
} |
|||
if (CoAP.ResponseCode.isSuccess(((Response) response.getCoapResponse()).getCode())) { |
|||
this.handleResponse(lwM2MClient, request.getPath().toString(), response, request, rpcRequest); |
|||
} else { |
|||
String msg = String.format("%s: SendRequest %s: CoapCode - %s Lwm2m code - %d name - %s Resource path - %s", LOG_LW2M_ERROR, request.getClass().getName().toString(), |
|||
((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString()); |
|||
handler.sendLogsToThingsboard(lwM2MClient, msg); |
|||
log.error("[{}] [{}], [{}] - [{}] [{}] error SendRequest", request.getClass().getName().toString(), registration.getEndpoint(), |
|||
((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString()); |
|||
if (!lwM2MClient.isInit()) { |
|||
lwM2MClient.initReadValue(this.handler, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration)); |
|||
} |
|||
/** Not Found */ |
|||
if (rpcRequest != null) { |
|||
handler.sentRpcResponse(rpcRequest, response.getCode().getName(), response.getErrorMessage(), LOG_LW2M_ERROR); |
|||
} |
|||
/** Not Found |
|||
set setClient_fw_info... = empty |
|||
**/ |
|||
if (lwM2MClient.getFwUpdate() != null && lwM2MClient.getFwUpdate().isInfoFwSwUpdate()) { |
|||
lwM2MClient.getFwUpdate().initReadValue(handler, this, request.getPath().toString()); |
|||
} |
|||
if (lwM2MClient.getSwUpdate() != null && lwM2MClient.getSwUpdate().isInfoFwSwUpdate()) { |
|||
lwM2MClient.getSwUpdate().initReadValue(handler, this, request.getPath().toString()); |
|||
} |
|||
if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) { |
|||
this.afterWriteFwSWUpdateError(registration, request, response.getErrorMessage()); |
|||
} |
|||
if (request.getPath().toString().equals(FW_UPDATE_ID) || request.getPath().toString().equals(SW_INSTALL_ID)) { |
|||
this.afterExecuteFwSwUpdateError(registration, request, response.getErrorMessage()); |
|||
} |
|||
} |
|||
}, e -> { |
|||
/** version == null |
|||
set setClient_fw_info... = empty |
|||
**/ |
|||
if (lwM2MClient.getFwUpdate() != null && lwM2MClient.getFwUpdate().isInfoFwSwUpdate()) { |
|||
lwM2MClient.getFwUpdate().initReadValue(handler, this, request.getPath().toString()); |
|||
} |
|||
if (lwM2MClient.getSwUpdate() != null && lwM2MClient.getSwUpdate().isInfoFwSwUpdate()) { |
|||
lwM2MClient.getSwUpdate().initReadValue(handler, this, request.getPath().toString()); |
|||
} |
|||
if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) { |
|||
this.afterWriteFwSWUpdateError(registration, request, e.getMessage()); |
|||
} |
|||
if (request.getPath().toString().equals(FW_UPDATE_ID) || request.getPath().toString().equals(SW_INSTALL_ID)) { |
|||
this.afterExecuteFwSwUpdateError(registration, request, e.getMessage()); |
|||
} |
|||
if (!lwM2MClient.isInit()) { |
|||
lwM2MClient.initReadValue(this.handler, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration)); |
|||
} |
|||
String msg = String.format("%s: SendRequest %s: Resource path - %s msg error - %s", |
|||
LOG_LW2M_ERROR, request.getClass().getName().toString(), request.getPath().toString(), e.getMessage()); |
|||
handler.sendLogsToThingsboard(lwM2MClient, msg); |
|||
log.error("[{}] [{}] - [{}] error SendRequest", request.getClass().getName().toString(), request.getPath().toString(), e.toString()); |
|||
if (rpcRequest != null) { |
|||
handler.sentRpcResponse(rpcRequest, CoAP.CodeClass.ERROR_RESPONSE.name(), e.getMessage(), LOG_LW2M_ERROR); |
|||
} |
|||
}); |
|||
} |
|||
|
|||
private WriteRequest getWriteRequestSingleResource(ContentFormat contentFormat, Integer objectId, Integer instanceId, |
|||
Integer resourceId, Object value, ResourceModel.Type type, |
|||
LwM2mClient client, LwM2mClientRpcRequest rpcRequest) { |
|||
try { |
|||
if (type != null) { |
|||
switch (type) { |
|||
case STRING: // String
|
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, value.toString()) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, value.toString()); |
|||
case INTEGER: // Long
|
|||
final long valueInt = Integer.toUnsignedLong(Integer.parseInt(value.toString())); |
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, valueInt) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, valueInt); |
|||
case OBJLNK: // ObjectLink
|
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())); |
|||
case BOOLEAN: // Boolean
|
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())); |
|||
case FLOAT: // Double
|
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Double.parseDouble(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Double.parseDouble(value.toString())); |
|||
case TIME: // Date
|
|||
Date date = new Date(Long.decode(value.toString())); |
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, date) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, date); |
|||
case OPAQUE: // byte[] value, base64
|
|||
byte[] valueRequest = value instanceof byte[] ? (byte[]) value : Hex.decodeHex(value.toString().toCharArray()); |
|||
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, valueRequest) : |
|||
new WriteRequest(contentFormat, objectId, instanceId, resourceId, valueRequest); |
|||
default: |
|||
} |
|||
} |
|||
if (rpcRequest != null) { |
|||
String patn = "/" + objectId + "/" + instanceId + "/" + resourceId; |
|||
String errorMsg = String.format("Bad ResourceModel Operations (E): Resource path - %s ResourceModel type - %s", patn, type); |
|||
rpcRequest.setErrorMsg(errorMsg); |
|||
} |
|||
return null; |
|||
} catch (NumberFormatException e) { |
|||
String patn = "/" + objectId + "/" + instanceId + "/" + resourceId; |
|||
String msg = String.format(LOG_LW2M_ERROR + ": NumberFormatException: Resource path - %s type - %s value - %s msg error - %s SendRequest to Client", |
|||
patn, type, value, e.toString()); |
|||
handler.sendLogsToThingsboard(client, msg); |
|||
log.error("Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]", patn, type, value, e.toString()); |
|||
if (rpcRequest != null) { |
|||
String errorMsg = String.format("NumberFormatException: Resource path - %s type - %s value - %s", patn, type, value); |
|||
handler.sentRpcResponse(rpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR); |
|||
} |
|||
return null; |
|||
} |
|||
} |
|||
|
|||
private void handleResponse(LwM2mClient lwM2mClient, final String path, LwM2mResponse response, |
|||
SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) { |
|||
responseRequestExecutor.submit(() -> { |
|||
try { |
|||
this.sendResponse(lwM2mClient, path, response, request, rpcRequest); |
|||
} catch (Exception e) { |
|||
log.error("[{}] endpoint [{}] path [{}] Exception Unable to after send response.", lwM2mClient.getRegistration().getEndpoint(), path, e); |
|||
} |
|||
}); |
|||
} |
|||
|
|||
/** |
|||
* processing a response from a client |
|||
* |
|||
* @param path - |
|||
* @param response - |
|||
*/ |
|||
private void sendResponse(LwM2mClient lwM2mClient, String path, LwM2mResponse response, |
|||
SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) { |
|||
Registration registration = lwM2mClient.getRegistration(); |
|||
String pathIdVer = convertPathFromObjectIdToIdVer(path, registration); |
|||
String msgLog = ""; |
|||
if (response instanceof ReadResponse) { |
|||
handler.onUpdateValueAfterReadResponse(registration, pathIdVer, (ReadResponse) response, rpcRequest); |
|||
} else if (response instanceof DeleteResponse) { |
|||
log.warn("11) [{}] Path [{}] DeleteResponse", pathIdVer, response); |
|||
if (rpcRequest != null) { |
|||
rpcRequest.setInfoMsg(null); |
|||
handler.sentRpcResponse(rpcRequest, response.getCode().getName(), null, null); |
|||
} |
|||
} else if (response instanceof DiscoverResponse) { |
|||
String discoverValue = Link.serialize(((DiscoverResponse) response).getObjectLinks()); |
|||
msgLog = String.format("%s: type operation: %s path: %s value: %s", |
|||
LOG_LW2M_INFO, DISCOVER.name(), request.getPath().toString(), discoverValue); |
|||
handler.sendLogsToThingsboard(lwM2mClient, msgLog); |
|||
log.warn("DiscoverResponse: [{}]", (DiscoverResponse) response); |
|||
if (rpcRequest != null) { |
|||
handler.sentRpcResponse(rpcRequest, response.getCode().getName(), discoverValue, LOG_LW2M_VALUE); |
|||
} |
|||
} else if (response instanceof ExecuteResponse) { |
|||
msgLog = String.format("%s: type operation: %s path: %s", |
|||
LOG_LW2M_INFO, EXECUTE.name(), request.getPath().toString()); |
|||
log.warn("9) [{}] ", msgLog); |
|||
handler.sendLogsToThingsboard(lwM2mClient, msgLog); |
|||
if (rpcRequest != null) { |
|||
msgLog = String.format("Start %s path: %S. Preparation finished: %s", EXECUTE.name(), path, rpcRequest.getInfoMsg()); |
|||
rpcRequest.setInfoMsg(msgLog); |
|||
handler.sentRpcResponse(rpcRequest, response.getCode().getName(), path, LOG_LW2M_INFO); |
|||
} |
|||
|
|||
} else if (response instanceof WriteAttributesResponse) { |
|||
msgLog = String.format("%s: type operation: %s path: %s value: %s", |
|||
LOG_LW2M_INFO, WRITE_ATTRIBUTES.name(), request.getPath().toString(), ((WriteAttributesRequest) request).getAttributes().toString()); |
|||
handler.sendLogsToThingsboard(lwM2mClient, msgLog); |
|||
log.warn("12) [{}] Path [{}] WriteAttributesResponse", pathIdVer, response); |
|||
if (rpcRequest != null) { |
|||
handler.sentRpcResponse(rpcRequest, response.getCode().getName(), response.toString(), LOG_LW2M_VALUE); |
|||
} |
|||
} else if (response instanceof WriteResponse) { |
|||
msgLog = String.format("Type operation: Write path: %s", pathIdVer); |
|||
log.warn("10) [{}] response: [{}]", msgLog, response); |
|||
this.infoWriteResponse(lwM2mClient, response, request, rpcRequest); |
|||
handler.onWriteResponseOk(registration, pathIdVer, (WriteRequest) request); |
|||
} |
|||
} |
|||
|
|||
private void infoWriteResponse(LwM2mClient lwM2mClient, LwM2mResponse response, SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) { |
|||
try { |
|||
Registration registration = lwM2mClient.getRegistration(); |
|||
LwM2mNode node = ((WriteRequest) request).getNode(); |
|||
String msg = null; |
|||
Object value; |
|||
if (node instanceof LwM2mObject) { |
|||
msg = String.format("%s: Update finished successfully: Lwm2m code - %d Source path: %s value: %s", |
|||
LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), ((LwM2mObject) node).toString()); |
|||
} else if (node instanceof LwM2mObjectInstance) { |
|||
msg = String.format("%s: Update finished successfully: Lwm2m code - %d Source path: %s value: %s", |
|||
LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), ((LwM2mObjectInstance) node).prettyPrint()); |
|||
} else if (node instanceof LwM2mSingleResource) { |
|||
LwM2mSingleResource singleResource = (LwM2mSingleResource) node; |
|||
if (singleResource.getType() == ResourceModel.Type.STRING || singleResource.getType() == ResourceModel.Type.OPAQUE) { |
|||
int valueLength; |
|||
if (singleResource.getType() == ResourceModel.Type.STRING) { |
|||
valueLength = ((String) singleResource.getValue()).length(); |
|||
value = ((String) singleResource.getValue()) |
|||
.substring(Math.min(valueLength, config.getLogMaxLength())).trim(); |
|||
|
|||
} else { |
|||
valueLength = ((byte[]) singleResource.getValue()).length; |
|||
value = new String(Arrays.copyOf(((byte[]) singleResource.getValue()), |
|||
Math.min(valueLength, config.getLogMaxLength()))).trim(); |
|||
} |
|||
value = valueLength > config.getLogMaxLength() ? value + "..." : value; |
|||
msg = String.format("%s: Update finished successfully: Lwm2m code - %d Resource path: %s length: %s value: %s", |
|||
LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), valueLength, value); |
|||
} else { |
|||
value = this.converter.convertValue(singleResource.getValue(), |
|||
singleResource.getType(), ResourceModel.Type.STRING, request.getPath()); |
|||
msg = String.format("%s: Update finished successfully. Lwm2m code: %d Resource path: %s value: %s", |
|||
LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), value); |
|||
} |
|||
} |
|||
if (msg != null) { |
|||
handler.sendLogsToThingsboard(lwM2mClient, msg); |
|||
if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) { |
|||
this.afterWriteSuccessFwSwUpdate(registration, request); |
|||
if (rpcRequest != null) { |
|||
rpcRequest.setInfoMsg(msg); |
|||
} |
|||
} |
|||
else if (rpcRequest != null) { |
|||
handler.sentRpcResponse(rpcRequest, response.getCode().getName(), msg, LOG_LW2M_INFO); |
|||
} |
|||
} |
|||
} catch (Exception e) { |
|||
log.trace("Fail convert value from request to string. ", e); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* After finish operation FwSwUpdate Write (success): |
|||
* fw_state/sw_state = DOWNLOADED |
|||
* send operation Execute |
|||
*/ |
|||
private void afterWriteSuccessFwSwUpdate(Registration registration, SimpleDownlinkRequest request) { |
|||
LwM2mClient lwM2MClient = this.lwM2mClientContext.getClientByRegistrationId(registration.getId()); |
|||
if (request.getPath().toString().equals(FW_PACKAGE_5_ID) && lwM2MClient.getFwUpdate() != null) { |
|||
lwM2MClient.getFwUpdate().setStateUpdate(DOWNLOADED.name()); |
|||
lwM2MClient.getFwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_INFO, null); |
|||
} |
|||
if (request.getPath().toString().equals(SW_PACKAGE_ID) && lwM2MClient.getSwUpdate() != null) { |
|||
lwM2MClient.getSwUpdate().setStateUpdate(DOWNLOADED.name()); |
|||
lwM2MClient.getSwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_INFO, null); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* After finish operation FwSwUpdate Write (error): fw_state = FAILED |
|||
*/ |
|||
private void afterWriteFwSWUpdateError(Registration registration, SimpleDownlinkRequest request, String msgError) { |
|||
LwM2mClient lwM2MClient = this.lwM2mClientContext.getClientByRegistrationId(registration.getId()); |
|||
if (request.getPath().toString().equals(FW_PACKAGE_5_ID) && lwM2MClient.getFwUpdate() != null) { |
|||
lwM2MClient.getFwUpdate().setStateUpdate(FAILED.name()); |
|||
lwM2MClient.getFwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_ERROR, msgError); |
|||
} |
|||
if (request.getPath().toString().equals(SW_PACKAGE_ID) && lwM2MClient.getSwUpdate() != null) { |
|||
lwM2MClient.getSwUpdate().setStateUpdate(FAILED.name()); |
|||
lwM2MClient.getSwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_ERROR, msgError); |
|||
} |
|||
} |
|||
|
|||
private void afterExecuteFwSwUpdateError(Registration registration, SimpleDownlinkRequest request, String msgError) { |
|||
LwM2mClient lwM2MClient = this.lwM2mClientContext.getClientByRegistrationId(registration.getId()); |
|||
if (request.getPath().toString().equals(FW_UPDATE_ID) && lwM2MClient.getFwUpdate() != null) { |
|||
lwM2MClient.getFwUpdate().sendLogs(this.handler, EXECUTE.name(), LOG_LW2M_ERROR, msgError); |
|||
} |
|||
if (request.getPath().toString().equals(SW_INSTALL_ID) && lwM2MClient.getSwUpdate() != null) { |
|||
lwM2MClient.getSwUpdate().sendLogs(this.handler, EXECUTE.name(), LOG_LW2M_ERROR, msgError); |
|||
} |
|||
} |
|||
|
|||
private void afterObserveCancel(LwM2mClient lwM2mClient, int observeCancelCnt, String observeCancelMsg, LwM2mClientRpcRequest rpcRequest) { |
|||
handler.sendLogsToThingsboard(lwM2mClient, observeCancelMsg); |
|||
log.warn("[{}]", observeCancelMsg); |
|||
if (rpcRequest != null) { |
|||
rpcRequest.setInfoMsg(String.format("Count: %d", observeCancelCnt)); |
|||
handler.sentRpcResponse(rpcRequest, CONTENT.name(), null, LOG_LW2M_INFO); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,221 @@ |
|||
/** |
|||
* 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.attributes; |
|||
|
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import com.google.common.util.concurrent.SettableFuture; |
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.eclipse.leshan.core.model.ResourceModel; |
|||
import org.eclipse.leshan.core.node.LwM2mPath; |
|||
import org.eclipse.leshan.core.node.LwM2mResource; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.transport.TransportService; |
|||
import org.thingsboard.server.common.transport.TransportServiceCallback; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg; |
|||
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; |
|||
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.LwM2mDownlinkMsgHandler; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteReplaceRequest; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteResponseCallback; |
|||
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; |
|||
import org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService; |
|||
import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MOtaUpdateService; |
|||
import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; |
|||
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Collection; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.atomic.AtomicInteger; |
|||
|
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_ERROR; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId; |
|||
|
|||
@Slf4j |
|||
@Service |
|||
@TbLwM2mTransportComponent |
|||
@RequiredArgsConstructor |
|||
public class DefaultLwM2MAttributesService implements LwM2MAttributesService { |
|||
|
|||
//TODO: add timeout logic
|
|||
private final AtomicInteger reqIdSeq = new AtomicInteger(); |
|||
private final Map<Integer, SettableFuture<List<TransportProtos.TsKvProto>>> futures; |
|||
|
|||
private final TransportService transportService; |
|||
private final LwM2mTransportServerHelper helper; |
|||
private final LwM2mClientContext clientContext; |
|||
private final LwM2MTransportServerConfig config; |
|||
private final LwM2mUplinkMsgHandler uplinkHandler; |
|||
private final LwM2mDownlinkMsgHandler downlinkHandler; |
|||
private final LwM2MTelemetryLogService logService; |
|||
private final LwM2MOtaUpdateService otaUpdateService; |
|||
|
|||
@Override |
|||
public ListenableFuture<List<TransportProtos.TsKvProto>> getSharedAttributes(LwM2mClient client, Collection<String> keys) { |
|||
SettableFuture<List<TransportProtos.TsKvProto>> future = SettableFuture.create(); |
|||
int requestId = reqIdSeq.incrementAndGet(); |
|||
futures.put(requestId, future); |
|||
transportService.process(client.getSession(), TransportProtos.GetAttributeRequestMsg.newBuilder().setRequestId(requestId). |
|||
addAllSharedAttributeNames(keys).build(), new TransportServiceCallback<Void>() { |
|||
@Override |
|||
public void onSuccess(Void msg) { |
|||
|
|||
} |
|||
|
|||
@Override |
|||
public void onError(Throwable e) { |
|||
SettableFuture<List<TransportProtos.TsKvProto>> callback = futures.remove(requestId); |
|||
if (callback != null) { |
|||
callback.setException(e); |
|||
} |
|||
} |
|||
}); |
|||
return future; |
|||
} |
|||
|
|||
@Override |
|||
public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse, TransportProtos.SessionInfoProto sessionInfo) { |
|||
var callback = futures.remove(getAttributesResponse.getRequestId()); |
|||
if (callback != null) { |
|||
callback.set(getAttributesResponse.getSharedAttributeListList()); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* Update - send request in change value resources in Client |
|||
* 1. FirmwareUpdate: |
|||
* - If msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0 |
|||
* 2. Shared Other AttributeUpdate |
|||
* -- Path to resources from profile equal keyName or from ModelObject equal name |
|||
* -- Only for resources: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase) |
|||
* 3. Delete - nothing |
|||
* |
|||
* @param msg - |
|||
*/ |
|||
@Override |
|||
public void onAttributesUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo) { |
|||
LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo); |
|||
if (msg.getSharedUpdatedCount() > 0 && lwM2MClient != null) { |
|||
String newFirmwareTitle = null; |
|||
String newFirmwareVersion = null; |
|||
String newFirmwareUrl = null; |
|||
String newSoftwareTitle = null; |
|||
String newSoftwareVersion = null; |
|||
List<TransportProtos.TsKvProto> otherAttributes = new ArrayList<>(); |
|||
for (TransportProtos.TsKvProto tsKvProto : msg.getSharedUpdatedList()) { |
|||
String attrName = tsKvProto.getKv().getKey(); |
|||
if (DefaultLwM2MOtaUpdateService.FIRMWARE_TITLE.equals(attrName)) { |
|||
newFirmwareTitle = getStrValue(tsKvProto); |
|||
} else if (DefaultLwM2MOtaUpdateService.FIRMWARE_VERSION.equals(attrName)) { |
|||
newFirmwareVersion = getStrValue(tsKvProto); |
|||
} else if (DefaultLwM2MOtaUpdateService.FIRMWARE_URL.equals(attrName)) { |
|||
newFirmwareUrl = getStrValue(tsKvProto); |
|||
} else if (DefaultLwM2MOtaUpdateService.SOFTWARE_TITLE.equals(attrName)) { |
|||
newSoftwareTitle = getStrValue(tsKvProto); |
|||
} else if (DefaultLwM2MOtaUpdateService.SOFTWARE_VERSION.equals(attrName)) { |
|||
newSoftwareVersion = getStrValue(tsKvProto); |
|||
} else { |
|||
otherAttributes.add(tsKvProto); |
|||
} |
|||
} |
|||
if (newFirmwareTitle != null || newFirmwareVersion != null) { |
|||
otaUpdateService.onTargetFirmwareUpdate(lwM2MClient, newFirmwareTitle, newFirmwareVersion, Optional.ofNullable(newFirmwareUrl)); |
|||
} |
|||
if (newSoftwareTitle != null || newSoftwareVersion != null) { |
|||
otaUpdateService.onTargetSoftwareUpdate(lwM2MClient, newSoftwareTitle, newSoftwareVersion); |
|||
} |
|||
if (!otherAttributes.isEmpty()) { |
|||
onAttributesUpdate(lwM2MClient, otherAttributes); |
|||
} |
|||
} else if (lwM2MClient == null) { |
|||
log.error("OnAttributeUpdate, lwM2MClient is null"); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* #1.1 If two names have equal path => last time attribute |
|||
* #2.1 if there is a difference in values between the current resource values and the shared attribute values |
|||
* => send to client Request Update of value (new value from shared attribute) |
|||
* and LwM2MClient.delayedRequests.add(path) |
|||
* #2.1 if there is not a difference in values between the current resource values and the shared attribute values |
|||
* |
|||
*/ |
|||
@Override |
|||
public void onAttributesUpdate(LwM2mClient lwM2MClient, List<TransportProtos.TsKvProto> tsKvProtos) { |
|||
log.trace("[{}] onAttributesUpdate [{}]", lwM2MClient.getEndpoint(), tsKvProtos); |
|||
tsKvProtos.forEach(tsKvProto -> { |
|||
String pathIdVer = clientContext.getObjectIdByKeyNameFromProfile(lwM2MClient, tsKvProto.getKv().getKey()); |
|||
if (pathIdVer != null) { |
|||
// #1.1
|
|||
if (lwM2MClient.getSharedAttributes().containsKey(pathIdVer)) { |
|||
if (tsKvProto.getTs() > lwM2MClient.getSharedAttributes().get(pathIdVer).getTs()) { |
|||
lwM2MClient.getSharedAttributes().put(pathIdVer, tsKvProto); |
|||
} |
|||
} else { |
|||
lwM2MClient.getSharedAttributes().put(pathIdVer, tsKvProto); |
|||
} |
|||
} |
|||
}); |
|||
clientContext.update(lwM2MClient); |
|||
// #2.1
|
|||
lwM2MClient.getSharedAttributes().forEach((pathIdVer, tsKvProto) -> { |
|||
this.pushUpdateToClientIfNeeded(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), |
|||
getValueFromKvProto(tsKvProto.getKv()), pathIdVer); |
|||
}); |
|||
} |
|||
|
|||
private void pushUpdateToClientIfNeeded(LwM2mClient lwM2MClient, Object valueOld, Object newValue, String versionedId) { |
|||
if (newValue != null && (valueOld == null || !newValue.toString().equals(valueOld.toString()))) { |
|||
TbLwM2MWriteReplaceRequest request = TbLwM2MWriteReplaceRequest.builder().versionedId(versionedId).value(newValue).timeout(this.config.getTimeout()).build(); |
|||
downlinkHandler.sendWriteReplaceRequest(lwM2MClient, request, new TbLwM2MWriteResponseCallback(uplinkHandler, logService, lwM2MClient, versionedId)); |
|||
} else { |
|||
log.error("Failed update resource [{}] [{}]", versionedId, newValue); |
|||
String logMsg = String.format("%s: Failed update resource versionedId - %s value - %s. Value is not changed or bad", |
|||
LOG_LWM2M_ERROR, versionedId, newValue); |
|||
logService.log(lwM2MClient, logMsg); |
|||
log.info("Failed update resource [{}] [{}]", versionedId, newValue); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* @param pathIdVer - path resource |
|||
* @return - value of Resource into format KvProto or null |
|||
*/ |
|||
private Object getResourceValueFormatKv(LwM2mClient lwM2MClient, String pathIdVer) { |
|||
LwM2mResource resourceValue = LwM2mTransportUtil.getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer); |
|||
if (resourceValue != null) { |
|||
ResourceModel.Type currentType = resourceValue.getType(); |
|||
ResourceModel.Type expectedType = helper.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer); |
|||
return LwM2mValueConverterImpl.getInstance().convertValue(resourceValue.getValue(), currentType, expectedType, |
|||
new LwM2mPath(fromVersionedIdToObjectId(pathIdVer))); |
|||
} else { |
|||
return null; |
|||
} |
|||
} |
|||
|
|||
private String getStrValue(TransportProtos.TsKvProto tsKvProto) { |
|||
return tsKvProto.getKv().getStringV(); |
|||
} |
|||
} |
|||
@ -0,0 +1,34 @@ |
|||
/** |
|||
* 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.attributes; |
|||
|
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
|
|||
import java.util.Collection; |
|||
import java.util.List; |
|||
|
|||
public interface LwM2MAttributesService { |
|||
|
|||
ListenableFuture<List<TransportProtos.TsKvProto>> getSharedAttributes(LwM2mClient client, Collection<String> keys); |
|||
|
|||
void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg getAttributesResponse, TransportProtos.SessionInfoProto sessionInfo); |
|||
|
|||
void onAttributesUpdate(TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotification, TransportProtos.SessionInfoProto sessionInfo); |
|||
|
|||
void onAttributesUpdate(LwM2mClient lwM2MClient, List<TransportProtos.TsKvProto> tsKvProtos); |
|||
} |
|||
@ -0,0 +1,22 @@ |
|||
/** |
|||
* 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; |
|||
|
|||
public class LwM2MAuthException extends RuntimeException { |
|||
|
|||
private static final long serialVersionUID = 4202690897971364044L; |
|||
|
|||
} |
|||
@ -1,113 +0,0 @@ |
|||
/** |
|||
* 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 com.google.gson.Gson; |
|||
import com.google.gson.JsonArray; |
|||
import com.google.gson.JsonObject; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2MClientStrategy; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2MFirmwareUpdateStrategy; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2MSoftwareUpdateStrategy; |
|||
|
|||
@Data |
|||
public class LwM2mClientProfile { |
|||
private final String clientStrategyStr = "clientStrategy"; |
|||
private final String fwUpdateStrategyStr = "fwUpdateStrategy"; |
|||
private final String swUpdateStrategyStr = "swUpdateStrategy"; |
|||
|
|||
private TenantId tenantId; |
|||
/** |
|||
* "clientLwM2mSettings": { |
|||
* "fwUpdateStrategy": "1", |
|||
* "swUpdateStrategy": "1", |
|||
* "clientStrategy": "1" |
|||
* } |
|||
**/ |
|||
private JsonObject postClientLwM2mSettings; |
|||
|
|||
/** |
|||
* {"keyName": { |
|||
* "/3_1.0/0/1": "modelNumber", |
|||
* "/3_1.0/0/0": "manufacturer", |
|||
* "/3_1.0/0/2": "serialNumber" |
|||
* } |
|||
**/ |
|||
private JsonObject postKeyNameProfile; |
|||
|
|||
/** |
|||
* [ "/3_1.0/0/0", "/3_1.0/0/1"] |
|||
*/ |
|||
private JsonArray postAttributeProfile; |
|||
|
|||
/** |
|||
* [ "/3_1.0/0/0", "/3_1.0/0/2"] |
|||
*/ |
|||
private JsonArray postTelemetryProfile; |
|||
|
|||
/** |
|||
* [ "/3_1.0/0", "/3_1.0/0/1, "/3_1.0/0/2"] |
|||
*/ |
|||
private JsonArray postObserveProfile; |
|||
|
|||
/** |
|||
* "attributeLwm2m": {"/3_1.0": {"ver": "currentTimeTest11"}, |
|||
* "/3_1.0/0": {"gt": 17}, |
|||
* "/3_1.0/0/9": {"pmax": 45}, "/3_1.2": {ver": "3_1.2"}} |
|||
*/ |
|||
private JsonObject postAttributeLwm2mProfile; |
|||
|
|||
public LwM2mClientProfile clone() { |
|||
LwM2mClientProfile lwM2mClientProfile = new LwM2mClientProfile(); |
|||
lwM2mClientProfile.postClientLwM2mSettings = this.deepCopy(this.postClientLwM2mSettings, JsonObject.class); |
|||
lwM2mClientProfile.postKeyNameProfile = this.deepCopy(this.postKeyNameProfile, JsonObject.class); |
|||
lwM2mClientProfile.postAttributeProfile = this.deepCopy(this.postAttributeProfile, JsonArray.class); |
|||
lwM2mClientProfile.postTelemetryProfile = this.deepCopy(this.postTelemetryProfile, JsonArray.class); |
|||
lwM2mClientProfile.postObserveProfile = this.deepCopy(this.postObserveProfile, JsonArray.class); |
|||
lwM2mClientProfile.postAttributeLwm2mProfile = this.deepCopy(this.postAttributeLwm2mProfile, JsonObject.class); |
|||
return lwM2mClientProfile; |
|||
} |
|||
|
|||
|
|||
private <T> T deepCopy(T elements, Class<T> type) { |
|||
try { |
|||
Gson gson = new Gson(); |
|||
return gson.fromJson(gson.toJson(elements), type); |
|||
} catch (Exception e) { |
|||
e.printStackTrace(); |
|||
return null; |
|||
} |
|||
} |
|||
|
|||
public int getClientStrategy() { |
|||
return this.postClientLwM2mSettings.getAsJsonObject().has(this.clientStrategyStr) ? |
|||
Integer.parseInt(this.postClientLwM2mSettings.getAsJsonObject().get(this.clientStrategyStr).getAsString()) : |
|||
LwM2MClientStrategy.CLIENT_STRATEGY_1.code; |
|||
} |
|||
|
|||
public int getFwUpdateStrategy() { |
|||
return this.postClientLwM2mSettings.getAsJsonObject().has(this.fwUpdateStrategyStr) ? |
|||
Integer.parseInt(this.postClientLwM2mSettings.getAsJsonObject().get(this.fwUpdateStrategyStr).getAsString()) : |
|||
LwM2MFirmwareUpdateStrategy.OBJ_5_BINARY.code; |
|||
} |
|||
|
|||
public int getSwUpdateStrategy() { |
|||
return this.postClientLwM2mSettings.getAsJsonObject().has(this.swUpdateStrategyStr) ? |
|||
Integer.parseInt(this.postClientLwM2mSettings.getAsJsonObject().get(this.swUpdateStrategyStr).getAsString()) : |
|||
LwM2MSoftwareUpdateStrategy.BINARY.code; |
|||
} |
|||
} |
|||
@ -1,281 +0,0 @@ |
|||
/** |
|||
* 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 com.google.gson.Gson; |
|||
import com.google.gson.JsonObject; |
|||
import com.google.gson.reflect.TypeToken; |
|||
import lombok.Data; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.eclipse.leshan.core.node.LwM2mPath; |
|||
import org.eclipse.leshan.server.registration.Registration; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
import org.thingsboard.server.transport.lwm2m.server.DefaultLwM2MTransportMsgHandler; |
|||
|
|||
import java.util.Map; |
|||
import java.util.Objects; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.TimeoutException; |
|||
|
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.ERROR_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FINISH_JSON_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FINISH_VALUE_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.INFO_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.KEY_NAME_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.DISCOVER_ALL; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.EXECUTE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.FW_UPDATE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_READ_ALL; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_ATTRIBUTES; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_REPLACE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_UPDATE; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.METHOD_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.PARAMS_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.RESULT_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SEPARATOR_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.START_JSON_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.TARGET_ID_VER_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.VALUE_KEY; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromIdVerToObjectId; |
|||
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validPathIdVer; |
|||
|
|||
@Slf4j |
|||
@Data |
|||
public class LwM2mClientRpcRequest { |
|||
|
|||
private Registration registration; |
|||
private TransportProtos.SessionInfoProto sessionInfo; |
|||
private String bodyParams; |
|||
private int requestId; |
|||
|
|||
private LwM2mTypeOper typeOper; |
|||
private String key; |
|||
private String targetIdVer; |
|||
private Object value; |
|||
private Map<String, Object> params; |
|||
|
|||
private String errorMsg; |
|||
private String valueMsg; |
|||
private String infoMsg; |
|||
private String responseCode; |
|||
|
|||
public LwM2mClientRpcRequest() { |
|||
} |
|||
|
|||
public LwM2mClientRpcRequest(LwM2mTypeOper lwM2mTypeOper, String bodyParams, int requestId, |
|||
TransportProtos.SessionInfoProto sessionInfo, Registration registration, DefaultLwM2MTransportMsgHandler handler) { |
|||
this.registration = registration; |
|||
this.sessionInfo = sessionInfo; |
|||
this.requestId = requestId; |
|||
if (lwM2mTypeOper != null) { |
|||
this.typeOper = lwM2mTypeOper; |
|||
} else { |
|||
this.errorMsg = METHOD_KEY + " - " + typeOper + " is not valid."; |
|||
} |
|||
if (this.errorMsg == null && !bodyParams.equals("null")) { |
|||
this.bodyParams = bodyParams; |
|||
this.init(handler); |
|||
} |
|||
} |
|||
|
|||
public TransportProtos.ToDeviceRpcResponseMsg getDeviceRpcResponseResultMsg() { |
|||
JsonObject payloadResp = new JsonObject(); |
|||
payloadResp.addProperty(RESULT_KEY, this.responseCode); |
|||
if (this.errorMsg != null) { |
|||
payloadResp.addProperty(ERROR_KEY, this.errorMsg); |
|||
} else if (this.valueMsg != null) { |
|||
payloadResp.addProperty(VALUE_KEY, this.valueMsg); |
|||
} else if (this.infoMsg != null) { |
|||
payloadResp.addProperty(INFO_KEY, this.infoMsg); |
|||
} |
|||
return TransportProtos.ToDeviceRpcResponseMsg.newBuilder() |
|||
.setPayload(payloadResp.getAsJsonObject().toString()) |
|||
.setRequestId(this.requestId) |
|||
.build(); |
|||
} |
|||
|
|||
private void init(DefaultLwM2MTransportMsgHandler handler) { |
|||
try { |
|||
// #1
|
|||
if (this.bodyParams.contains(KEY_NAME_KEY)) { |
|||
String targetIdVerStr = this.getValueKeyFromBody(KEY_NAME_KEY); |
|||
if (targetIdVerStr != null) { |
|||
String targetIdVer = handler.getPresentPathIntoProfile(sessionInfo, targetIdVerStr); |
|||
if (targetIdVer != null) { |
|||
this.targetIdVer = targetIdVer; |
|||
this.setInfoMsg(String.format("Changed by: key - %s, pathIdVer - %s", |
|||
targetIdVerStr, targetIdVer)); |
|||
} |
|||
} |
|||
} |
|||
if (this.getTargetIdVer() == null && this.bodyParams.contains(TARGET_ID_VER_KEY)) { |
|||
this.setValidTargetIdVerKey(); |
|||
} |
|||
if (this.bodyParams.contains(VALUE_KEY)) { |
|||
this.value = this.getValueKeyFromBody(VALUE_KEY); |
|||
} |
|||
try { |
|||
if (this.bodyParams.contains(PARAMS_KEY)) { |
|||
this.setValidParamsKey(handler); |
|||
} |
|||
} catch (Exception e) { |
|||
this.setErrorMsg(String.format("Params of request is bad Json format. %s", e.getMessage())); |
|||
} |
|||
|
|||
if (this.getTargetIdVer() == null |
|||
&& !(OBSERVE_READ_ALL == this.getTypeOper() |
|||
|| DISCOVER_ALL == this.getTypeOper() |
|||
|| OBSERVE_CANCEL == this.getTypeOper() |
|||
|| FW_UPDATE == this.getTypeOper())) { |
|||
this.setErrorMsg(TARGET_ID_VER_KEY + " and " + |
|||
KEY_NAME_KEY + " is null or bad format"); |
|||
} |
|||
/** |
|||
* EXECUTE && WRITE_REPLACE - only for Resource or ResourceInstance |
|||
*/ |
|||
else if (this.getTargetIdVer() != null |
|||
&& (EXECUTE == this.getTypeOper() |
|||
|| WRITE_REPLACE == this.getTypeOper()) |
|||
&& !(new LwM2mPath(Objects.requireNonNull(convertPathFromIdVerToObjectId(this.getTargetIdVer()))).isResource() |
|||
|| new LwM2mPath(Objects.requireNonNull(convertPathFromIdVerToObjectId(this.getTargetIdVer()))).isResourceInstance())) { |
|||
this.setErrorMsg("Invalid parameter " + TARGET_ID_VER_KEY |
|||
+ ". Only Resource or ResourceInstance can be this operation"); |
|||
} |
|||
} catch (Exception e) { |
|||
this.setErrorMsg(String.format("Bad format request. %s", e.getMessage())); |
|||
} |
|||
|
|||
} |
|||
|
|||
private void setValidTargetIdVerKey() { |
|||
String targetIdVerStr = this.getValueKeyFromBody(TARGET_ID_VER_KEY); |
|||
// targetIdVer without ver - ok
|
|||
try { |
|||
// targetIdVer with/without ver - ok
|
|||
this.targetIdVer = validPathIdVer(targetIdVerStr, this.registration); |
|||
if (this.targetIdVer != null) { |
|||
this.infoMsg = String.format("Changed by: pathIdVer - %s", this.targetIdVer); |
|||
} |
|||
} catch (Exception e) { |
|||
if (this.targetIdVer == null) { |
|||
this.errorMsg = TARGET_ID_VER_KEY + " - " + targetIdVerStr + " is not valid."; |
|||
} |
|||
} |
|||
} |
|||
|
|||
private void setValidParamsKey(DefaultLwM2MTransportMsgHandler handler) { |
|||
String paramsStr = this.getValueKeyFromBody(PARAMS_KEY); |
|||
if (paramsStr != null) { |
|||
String params2Json = |
|||
START_JSON_KEY |
|||
+ "\"" |
|||
+ paramsStr |
|||
.replaceAll(SEPARATOR_KEY, "\"" + SEPARATOR_KEY + "\"") |
|||
.replaceAll(FINISH_VALUE_KEY, "\"" + FINISH_VALUE_KEY + "\"") |
|||
+ "\"" |
|||
+ FINISH_JSON_KEY; |
|||
// jsonObject
|
|||
Map<String, Object> params = new Gson().fromJson(params2Json, new TypeToken<ConcurrentHashMap<String, Object>>() { |
|||
}.getType()); |
|||
if (WRITE_UPDATE == this.getTypeOper()) { |
|||
if (this.targetIdVer != null) { |
|||
Map<String, Object> paramsResourceId = this.convertParamsToResourceId((ConcurrentHashMap<String, Object>) params, handler); |
|||
if (paramsResourceId.size() > 0) { |
|||
this.setParams(paramsResourceId); |
|||
} |
|||
} |
|||
} else if (WRITE_ATTRIBUTES == this.getTypeOper()) { |
|||
this.setParams(params); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private String getValueKeyFromBody(String key) { |
|||
String valueKey = null; |
|||
int startInd = -1; |
|||
int finishInd = -1; |
|||
try { |
|||
switch (key) { |
|||
case KEY_NAME_KEY: |
|||
case TARGET_ID_VER_KEY: |
|||
case VALUE_KEY: |
|||
startInd = this.bodyParams.indexOf(SEPARATOR_KEY, this.bodyParams.indexOf(key)); |
|||
finishInd = this.bodyParams.indexOf(FINISH_VALUE_KEY, this.bodyParams.indexOf(key)); |
|||
if (startInd >= 0 && finishInd < 0) { |
|||
finishInd = this.bodyParams.indexOf(FINISH_JSON_KEY, this.bodyParams.indexOf(key)); |
|||
} |
|||
break; |
|||
case PARAMS_KEY: |
|||
startInd = this.bodyParams.indexOf(START_JSON_KEY, this.bodyParams.indexOf(key)); |
|||
finishInd = this.bodyParams.indexOf(FINISH_JSON_KEY, this.bodyParams.indexOf(key)); |
|||
} |
|||
if (startInd >= 0 && finishInd > 0) { |
|||
valueKey = this.bodyParams.substring(startInd + 1, finishInd); |
|||
} |
|||
} catch (Exception e) { |
|||
log.error("", new TimeoutException()); |
|||
} |
|||
/** |
|||
* ReplaceAll "\"" |
|||
*/ |
|||
if (StringUtils.trimToNull(valueKey) != null) { |
|||
char[] chars = valueKey.toCharArray(); |
|||
for (int i = 0; i < chars.length; i++) { |
|||
if (chars[i] == 92 || chars[i] == 34) chars[i] = 32; |
|||
} |
|||
return key.equals(PARAMS_KEY) ? String.valueOf(chars) : String.valueOf(chars).replaceAll(" ", ""); |
|||
} |
|||
return null; |
|||
} |
|||
|
|||
private ConcurrentHashMap<String, Object> convertParamsToResourceId(ConcurrentHashMap<String, Object> params, |
|||
DefaultLwM2MTransportMsgHandler serviceImpl) { |
|||
Map<String, Object> paramsIdVer = new ConcurrentHashMap<>(); |
|||
LwM2mPath targetId = new LwM2mPath(Objects.requireNonNull(convertPathFromIdVerToObjectId(this.targetIdVer))); |
|||
if (targetId.isObjectInstance()) { |
|||
params.forEach((k, v) -> { |
|||
try { |
|||
int id = Integer.parseInt(k); |
|||
paramsIdVer.put(String.valueOf(id), v); |
|||
} catch (NumberFormatException e) { |
|||
String targetIdVer = serviceImpl.getPresentPathIntoProfile(sessionInfo, k); |
|||
if (targetIdVer != null) { |
|||
LwM2mPath lwM2mPath = new LwM2mPath(Objects.requireNonNull(convertPathFromIdVerToObjectId(targetIdVer))); |
|||
paramsIdVer.put(String.valueOf(lwM2mPath.getResourceId()), v); |
|||
} |
|||
/** WRITE_UPDATE*/ |
|||
else { |
|||
String rezId = this.getRezIdByResourceNameAndObjectInstanceId(k, serviceImpl); |
|||
if (rezId != null) { |
|||
paramsIdVer.put(rezId, v); |
|||
} |
|||
} |
|||
} |
|||
}); |
|||
} |
|||
return (ConcurrentHashMap<String, Object>) paramsIdVer; |
|||
} |
|||
|
|||
private String getRezIdByResourceNameAndObjectInstanceId(String resourceName, DefaultLwM2MTransportMsgHandler handler) { |
|||
LwM2mClient lwM2mClient = handler.clientContext.getClientBySessionInfo(this.sessionInfo); |
|||
return lwM2mClient != null ? |
|||
lwM2mClient.getRezIdByResourceNameAndObjectInstanceId(resourceName, this.targetIdVer, handler.config.getModelProvider()) : |
|||
null; |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue