|
|
|
@ -18,14 +18,16 @@ package org.thingsboard.server.dao.device; |
|
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
|
import com.fasterxml.jackson.databind.node.ArrayNode; |
|
|
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
|
|
|
import lombok.RequiredArgsConstructor; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.bouncycastle.cert.X509CertificateHolder; |
|
|
|
import org.bouncycastle.openssl.PEMParser; |
|
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
|
import org.springframework.core.io.ByteArrayResource; |
|
|
|
import org.springframework.core.io.Resource; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.server.common.data.AdminSettings; |
|
|
|
import org.thingsboard.server.common.data.Device; |
|
|
|
import org.thingsboard.server.common.data.DeviceProfile; |
|
|
|
import org.thingsboard.server.common.data.DeviceTransportType; |
|
|
|
@ -33,11 +35,12 @@ import org.thingsboard.server.common.data.ResourceUtils; |
|
|
|
import org.thingsboard.server.common.data.StringUtils; |
|
|
|
import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.security.DeviceCredentials; |
|
|
|
import org.thingsboard.server.common.data.security.DeviceCredentialsType; |
|
|
|
import org.thingsboard.server.dao.settings.AdminSettingsService; |
|
|
|
import org.thingsboard.server.dao.util.DeviceConnectivityUtil; |
|
|
|
|
|
|
|
import javax.annotation.PostConstruct; |
|
|
|
import java.io.InputStream; |
|
|
|
import java.io.InputStreamReader; |
|
|
|
import java.net.URI; |
|
|
|
@ -64,6 +67,7 @@ import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.WINDOWS; |
|
|
|
|
|
|
|
@Service("DeviceConnectivityDaoService") |
|
|
|
@Slf4j |
|
|
|
@RequiredArgsConstructor |
|
|
|
public class DeviceConnectivityServiceImpl implements DeviceConnectivityService { |
|
|
|
|
|
|
|
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; |
|
|
|
@ -74,26 +78,12 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
|
|
|
|
private final Map<String, Resource> certs = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private DeviceCredentialsService deviceCredentialsService; |
|
|
|
private final DeviceCredentialsService deviceCredentialsService; |
|
|
|
private final DeviceProfileService deviceProfileService; |
|
|
|
private final AdminSettingsService adminSettingsService; |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private DeviceProfileService deviceProfileService; |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private DeviceConnectivityConfiguration deviceConnectivityConfiguration; |
|
|
|
|
|
|
|
@PostConstruct |
|
|
|
private void init() { |
|
|
|
DeviceConnectivityInfo mqtts = deviceConnectivityConfiguration.getConnectivity(MQTTS); |
|
|
|
if (mqtts != null && mqtts.isEnabled()) { |
|
|
|
String certFilePath = mqtts.getPemCertFile(); |
|
|
|
if (StringUtils.isBlank(certFilePath) || !ResourceUtils.resourceExists(this, certFilePath)) { |
|
|
|
String error = StringUtils.isBlank(certFilePath) ? "path is empty" : "file is not exists"; |
|
|
|
log.error("MQTTS is enabled but cert {}!", error); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
@Value("${device.connectivity.mqtts.pem_cert_file:}") |
|
|
|
private String mqttsPemCertFile; |
|
|
|
|
|
|
|
@Override |
|
|
|
public JsonNode findDevicePublishTelemetryCommands(String baseUrl, Device device) throws URISyntaxException { |
|
|
|
@ -149,11 +139,11 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
DeviceCredentials creds = deviceCredentialsService.findDeviceCredentialsByDeviceId(device.getTenantId(), deviceId); |
|
|
|
|
|
|
|
ObjectNode commands = JacksonUtil.newObjectNode(); |
|
|
|
if (deviceConnectivityConfiguration.isEnabled(MQTT)) { |
|
|
|
if (isEnabled(MQTT)) { |
|
|
|
Optional.ofNullable(getGatewayDockerCommands(baseUrl, creds, MQTT)) |
|
|
|
.ifPresent(v -> commands.set(MQTT, v)); |
|
|
|
} |
|
|
|
if (deviceConnectivityConfiguration.isEnabled(MQTTS)) { |
|
|
|
if (isEnabled(MQTTS)) { |
|
|
|
Optional.ofNullable(getGatewayDockerCommands(baseUrl, creds, MQTTS)) |
|
|
|
.ifPresent(v -> commands.set(MQTTS, v)); |
|
|
|
} |
|
|
|
@ -163,15 +153,15 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
@Override |
|
|
|
public Resource getPemCertFile(String protocol) { |
|
|
|
return certs.computeIfAbsent(protocol, key -> { |
|
|
|
DeviceConnectivityInfo connectivity = deviceConnectivityConfiguration.getConnectivity(protocol); |
|
|
|
if (connectivity == null) { |
|
|
|
DeviceConnectivityInfo connectivity = getConnectivity(protocol); |
|
|
|
if (!MQTTS.equals(protocol) || connectivity == null) { |
|
|
|
log.warn("Unknown connectivity protocol: {}", protocol); |
|
|
|
return null; |
|
|
|
} |
|
|
|
String certFilePath = connectivity.getPemCertFile(); |
|
|
|
if (StringUtils.isNotBlank(certFilePath) && ResourceUtils.resourceExists(this, certFilePath)) { |
|
|
|
|
|
|
|
if (StringUtils.isNotBlank(mqttsPemCertFile) && ResourceUtils.resourceExists(this, mqttsPemCertFile)) { |
|
|
|
try { |
|
|
|
return getCert(certFilePath); |
|
|
|
return getCert(mqttsPemCertFile); |
|
|
|
} catch (Exception e) { |
|
|
|
String msg = String.format("Failed to read %s server certificate!", protocol); |
|
|
|
log.warn(msg); |
|
|
|
@ -183,6 +173,20 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private DeviceConnectivityInfo getConnectivity(String protocol) { |
|
|
|
AdminSettings connectivitySettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "connectivity"); |
|
|
|
JsonNode connectivity; |
|
|
|
if (connectivitySettings != null && (connectivity = connectivitySettings.getJsonValue()) != null) { |
|
|
|
return JacksonUtil.convertValue(connectivity.get(protocol), DeviceConnectivityInfo.class); |
|
|
|
} |
|
|
|
return null; |
|
|
|
} |
|
|
|
|
|
|
|
public boolean isEnabled(String protocol) { |
|
|
|
var info = getConnectivity(protocol); |
|
|
|
return info != null && info.isEnabled(); |
|
|
|
} |
|
|
|
|
|
|
|
private Resource getCert(String path) throws Exception { |
|
|
|
StringBuilder pemContentBuilder = new StringBuilder(); |
|
|
|
|
|
|
|
@ -221,7 +225,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
} |
|
|
|
|
|
|
|
private String getHttpPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException { |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity(protocol); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(protocol); |
|
|
|
if (properties == null || !properties.isEnabled() || |
|
|
|
deviceCredentials.getCredentialsType() != DeviceCredentialsType.ACCESS_TOKEN) { |
|
|
|
return null; |
|
|
|
@ -247,7 +251,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
|
|
|
|
ObjectNode dockerMqttCommands = JacksonUtil.newObjectNode(); |
|
|
|
|
|
|
|
if (deviceConnectivityConfiguration.isEnabled(MQTT)) { |
|
|
|
if (isEnabled(MQTT)) { |
|
|
|
Optional.ofNullable(getMqttPublishCommand(baseUrl, topic, deviceCredentials)). |
|
|
|
ifPresent(v -> mqttCommands.put(MQTT, v)); |
|
|
|
|
|
|
|
@ -255,7 +259,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
.ifPresent(v -> dockerMqttCommands.put(MQTT, v)); |
|
|
|
} |
|
|
|
|
|
|
|
if (deviceConnectivityConfiguration.isEnabled(MQTTS)) { |
|
|
|
if (isEnabled(MQTTS)) { |
|
|
|
List<String> mqttsPublishCommand = getMqttsPublishCommand(baseUrl, topic, deviceCredentials); |
|
|
|
if (mqttsPublishCommand != null) { |
|
|
|
ArrayNode arrayNode = mqttCommands.putArray(MQTTS); |
|
|
|
@ -273,14 +277,14 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
} |
|
|
|
|
|
|
|
private String getMqttPublishCommand(String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException { |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity(MQTT); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(MQTT); |
|
|
|
String mqttHost = getHost(baseUrl, properties); |
|
|
|
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); |
|
|
|
return DeviceConnectivityUtil.getMqttPublishCommand(MQTT, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials); |
|
|
|
} |
|
|
|
|
|
|
|
private List<String> getMqttsPublishCommand(String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException { |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity(MQTTS); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(MQTTS); |
|
|
|
String mqttHost = getHost(baseUrl, properties); |
|
|
|
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); |
|
|
|
String pubCommand = DeviceConnectivityUtil.getMqttPublishCommand(MQTTS, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials); |
|
|
|
@ -296,7 +300,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
|
|
|
|
private JsonNode getGatewayDockerCommands(String baseUrl, DeviceCredentials deviceCredentials, String mqttType) throws URISyntaxException { |
|
|
|
ObjectNode dockerLaunchCommands = JacksonUtil.newObjectNode(); |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity().get(mqttType); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(mqttType); |
|
|
|
String mqttHost = getHost(baseUrl, properties); |
|
|
|
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); |
|
|
|
Optional.ofNullable(DeviceConnectivityUtil.getGatewayLaunchCommand(LINUX, mqttHost, mqttPort, deviceCredentials)) |
|
|
|
@ -307,7 +311,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
} |
|
|
|
|
|
|
|
private String getDockerMqttPublishCommand(String protocol, String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException { |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity(protocol); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(protocol); |
|
|
|
String mqttHost = getHost(baseUrl, properties); |
|
|
|
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); |
|
|
|
return DeviceConnectivityUtil.getDockerMqttPublishCommand(protocol, baseUrl, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials); |
|
|
|
@ -323,7 +327,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
|
|
|
|
ObjectNode dockerCoapCommands = JacksonUtil.newObjectNode(); |
|
|
|
|
|
|
|
if (deviceConnectivityConfiguration.isEnabled(COAP)) { |
|
|
|
if (isEnabled(COAP)) { |
|
|
|
Optional.ofNullable(getCoapPublishCommand(COAP, baseUrl, deviceCredentials)) |
|
|
|
.ifPresent(v -> coapCommands.put(COAP, v)); |
|
|
|
|
|
|
|
@ -331,7 +335,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
.ifPresent(v -> dockerCoapCommands.put(COAP, v)); |
|
|
|
} |
|
|
|
|
|
|
|
if (deviceConnectivityConfiguration.isEnabled(COAPS)) { |
|
|
|
if (isEnabled(COAPS)) { |
|
|
|
Optional.ofNullable(getCoapPublishCommand(COAPS, baseUrl, deviceCredentials)) |
|
|
|
.ifPresent(v -> coapCommands.put(COAPS, v)); |
|
|
|
|
|
|
|
@ -347,14 +351,14 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService |
|
|
|
} |
|
|
|
|
|
|
|
private String getCoapPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException { |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity(protocol); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(protocol); |
|
|
|
String hostName = getHost(baseUrl, properties); |
|
|
|
String port = properties.getPort().isEmpty() ? "" : ":" + properties.getPort(); |
|
|
|
return DeviceConnectivityUtil.getCoapPublishCommand(protocol, hostName, port, deviceCredentials); |
|
|
|
} |
|
|
|
|
|
|
|
private String getDockerCoapPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException { |
|
|
|
DeviceConnectivityInfo properties = deviceConnectivityConfiguration.getConnectivity(protocol); |
|
|
|
DeviceConnectivityInfo properties = getConnectivity(protocol); |
|
|
|
String host = getHost(baseUrl, properties); |
|
|
|
String port = properties.getPort().isEmpty() ? "" : ":" + properties.getPort(); |
|
|
|
return DeviceConnectivityUtil.getDockerCoapPublishCommand(protocol, host, port, deviceCredentials); |
|
|
|
|