Browse Source

Improvements to Tenant rules and plugins startup sequence

pull/266/head
Andrew Shvayka 9 years ago
parent
commit
c4cd601dcb
  1. 3
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 2
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  3. 7
      application/src/main/java/org/thingsboard/server/actors/shared/plugin/TenantPluginManager.java
  4. 10
      application/src/main/java/org/thingsboard/server/actors/shared/rule/RuleManager.java
  5. 7
      application/src/main/java/org/thingsboard/server/actors/shared/rule/TenantRuleManager.java
  6. 11
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  7. 14
      application/src/main/java/org/thingsboard/server/controller/AlarmController.java
  8. 2
      application/src/main/java/org/thingsboard/server/controller/AssetController.java
  9. 2
      application/src/main/java/org/thingsboard/server/controller/DeviceController.java
  10. 2
      application/src/main/resources/thingsboard.yml
  11. 74
      application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java
  12. 5
      common/data/src/main/java/org/thingsboard/server/common/data/asset/AssetSearchQuery.java
  13. 2
      dao/src/main/java/org/thingsboard/server/common/data/device/DeviceSearchQuery.java
  14. 1
      dao/src/main/java/org/thingsboard/server/dao/asset/AssetService.java
  15. 1
      dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
  16. 1
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
  17. 1
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  18. 66
      transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

3
application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java

@ -136,6 +136,9 @@ public class ActorSystemContext {
@Value("${actors.statistics.persist_frequency}") @Value("${actors.statistics.persist_frequency}")
@Getter private long statisticsPersistFrequency; @Getter private long statisticsPersistFrequency;
@Value("${actors.tenant.create_components_on_init}")
@Getter private boolean tenantComponentsInitEnabled;
@Getter @Setter private ActorSystem actorSystem; @Getter @Setter private ActorSystem actorSystem;
@Getter @Setter private ActorRef appActor; @Getter @Setter private ActorRef appActor;

2
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -174,7 +174,7 @@ public class AppActor extends ContextAwareActor {
TenantId tenantId = toDeviceActorMsg.getTenantId(); TenantId tenantId = toDeviceActorMsg.getTenantId();
ActorRef tenantActor = getOrCreateTenantActor(tenantId); ActorRef tenantActor = getOrCreateTenantActor(tenantId);
if (toDeviceActorMsg.getPayload().getMsgType().requiresRulesProcessing()) { if (toDeviceActorMsg.getPayload().getMsgType().requiresRulesProcessing()) {
tenantActor.tell(new RuleChainDeviceMsg(toDeviceActorMsg, ruleManager.getRuleChain()), context().self()); tenantActor.tell(new RuleChainDeviceMsg(toDeviceActorMsg, ruleManager.getRuleChain(this.context())), context().self());
} else { } else {
tenantActor.tell(toDeviceActorMsg, context().self()); tenantActor.tell(toDeviceActorMsg, context().self());
} }

7
application/src/main/java/org/thingsboard/server/actors/shared/plugin/TenantPluginManager.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.actors.shared.plugin; package org.thingsboard.server.actors.shared.plugin;
import akka.actor.ActorContext;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -30,6 +31,12 @@ public class TenantPluginManager extends PluginManager {
this.tenantId = tenantId; this.tenantId = tenantId;
} }
public void init(ActorContext context) {
if (systemContext.isTenantComponentsInitEnabled()) {
super.init(context);
}
}
@Override @Override
FetchFunction<PluginMetaData> getFetchPluginsFunction() { FetchFunction<PluginMetaData> getFetchPluginsFunction() {
return link -> pluginService.findTenantPlugins(tenantId, link); return link -> pluginService.findTenantPlugins(tenantId, link);

10
application/src/main/java/org/thingsboard/server/actors/shared/rule/RuleManager.java

@ -25,7 +25,6 @@ import org.thingsboard.server.actors.rule.RuleActorChain;
import org.thingsboard.server.actors.rule.RuleActorMetaData; import org.thingsboard.server.actors.rule.RuleActorMetaData;
import org.thingsboard.server.actors.rule.SimpleRuleActorChain; import org.thingsboard.server.actors.rule.SimpleRuleActorChain;
import org.thingsboard.server.actors.service.ContextAwareActor; import org.thingsboard.server.actors.service.ContextAwareActor;
import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.common.data.id.RuleId; import org.thingsboard.server.common.data.id.RuleId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.page.PageDataIterable;
@ -72,6 +71,9 @@ public abstract class RuleManager {
} }
public Optional<ActorRef> update(ActorContext context, RuleId ruleId, ComponentLifecycleEvent event) { public Optional<ActorRef> update(ActorContext context, RuleId ruleId, ComponentLifecycleEvent event) {
if (ruleMap == null) {
init(context);
}
RuleMetaData rule; RuleMetaData rule;
if (event != ComponentLifecycleEvent.DELETED) { if (event != ComponentLifecycleEvent.DELETED) {
rule = systemContext.getRuleService().findRuleById(ruleId); rule = systemContext.getRuleService().findRuleById(ruleId);
@ -111,11 +113,13 @@ public abstract class RuleManager {
.withDispatcher(getDispatcherName()), rId.toString())); .withDispatcher(getDispatcherName()), rId.toString()));
} }
public RuleActorChain getRuleChain() { public RuleActorChain getRuleChain(ActorContext context) {
if (ruleMap == null) {
init(context);
}
return ruleChain; return ruleChain;
} }
private void refreshRuleChain() { private void refreshRuleChain() {
Set<RuleActorMetaData> activeRuleSet = new HashSet<>(); Set<RuleActorMetaData> activeRuleSet = new HashSet<>();
for (Map.Entry<RuleMetaData, RuleActorMetaData> rule : ruleMap.entrySet()) { for (Map.Entry<RuleMetaData, RuleActorMetaData> rule : ruleMap.entrySet()) {

7
application/src/main/java/org/thingsboard/server/actors/shared/rule/TenantRuleManager.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.actors.shared.rule; package org.thingsboard.server.actors.shared.rule;
import akka.actor.ActorContext;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -27,6 +28,12 @@ public class TenantRuleManager extends RuleManager {
super(systemContext, tenantId); super(systemContext, tenantId);
} }
public void init(ActorContext context) {
if (systemContext.isTenantComponentsInitEnabled()) {
super.init(context);
}
}
@Override @Override
FetchFunction<RuleMetaData> getFetchRulesFunction() { FetchFunction<RuleMetaData> getFetchRulesFunction() {
return link -> ruleService.findTenantRules(tenantId, link); return link -> ruleService.findTenantRules(tenantId, link);

11
application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java

@ -151,18 +151,13 @@ public class TenantActor extends ContextAwareActor {
private void process(RuleChainDeviceMsg msg) { private void process(RuleChainDeviceMsg msg) {
ToDeviceActorMsg toDeviceActorMsg = msg.getToDeviceActorMsg(); ToDeviceActorMsg toDeviceActorMsg = msg.getToDeviceActorMsg();
ActorRef deviceActor = getOrCreateDeviceActor(toDeviceActorMsg.getDeviceId()); ActorRef deviceActor = getOrCreateDeviceActor(toDeviceActorMsg.getDeviceId());
RuleActorChain chain = new ComplexRuleActorChain(msg.getRuleChain(), ruleManager.getRuleChain()); RuleActorChain chain = new ComplexRuleActorChain(msg.getRuleChain(), ruleManager.getRuleChain(this.context()));
deviceActor.tell(new RuleChainDeviceMsg(toDeviceActorMsg, chain), context().self()); deviceActor.tell(new RuleChainDeviceMsg(toDeviceActorMsg, chain), context().self());
} }
private ActorRef getOrCreateDeviceActor(DeviceId deviceId) { private ActorRef getOrCreateDeviceActor(DeviceId deviceId) {
ActorRef deviceActor = deviceActors.get(deviceId); return deviceActors.computeIfAbsent(deviceId, k -> context().actorOf(Props.create(new DeviceActor.ActorCreator(systemContext, tenantId, deviceId))
if (deviceActor == null) { .withDispatcher(DefaultActorService.CORE_DISPATCHER_NAME), deviceId.toString()));
deviceActor = context().actorOf(Props.create(new DeviceActor.ActorCreator(systemContext, tenantId, deviceId))
.withDispatcher(DefaultActorService.CORE_DISPATCHER_NAME), deviceId.toString());
deviceActors.put(deviceId, deviceActor);
}
return deviceActor;
} }
public static class ActorCreator extends ContextBasedCreator<TenantActor> { public static class ActorCreator extends ContextBasedCreator<TenantActor> {

14
application/src/main/java/org/thingsboard/server/controller/AlarmController.java

@ -15,30 +15,16 @@
*/ */
package org.thingsboard.server.controller; package org.thingsboard.server.controller;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.http.HttpStatus; import org.springframework.http.HttpStatus;
import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.*; import org.springframework.web.bind.annotation.*;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Event;
import org.thingsboard.server.common.data.alarm.*; import org.thingsboard.server.common.data.alarm.*;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.*; import org.thingsboard.server.common.data.id.*;
import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.page.TimePageData;
import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.asset.AssetSearchQuery;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.exception.ThingsboardErrorCode; import org.thingsboard.server.exception.ThingsboardErrorCode;
import org.thingsboard.server.exception.ThingsboardException; import org.thingsboard.server.exception.ThingsboardException;
import org.thingsboard.server.service.security.model.SecurityUser;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Collectors;
@RestController @RestController
@RequestMapping("/api") @RequestMapping("/api")

2
application/src/main/java/org/thingsboard/server/controller/AssetController.java

@ -27,7 +27,7 @@ import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.asset.AssetSearchQuery; import org.thingsboard.server.common.data.asset.AssetSearchQuery;
import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.exception.ThingsboardException; import org.thingsboard.server.exception.ThingsboardException;

2
application/src/main/java/org/thingsboard/server/controller/DeviceController.java

@ -28,7 +28,7 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.dao.device.DeviceSearchQuery; import org.thingsboard.server.common.data.device.DeviceSearchQuery;
import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.exception.ThingsboardException; import org.thingsboard.server.exception.ThingsboardException;

2
application/src/main/resources/thingsboard.yml

@ -159,6 +159,8 @@ cassandra:
# Actor system parameters # Actor system parameters
actors: actors:
tenant:
create_components_on_init: true
session: session:
sync: sync:
# Default timeout for processing request using synchronous session (HTTP, CoAP) in milliseconds # Default timeout for processing request using synchronous session (HTTP, CoAP) in milliseconds

74
application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java

@ -97,28 +97,28 @@ public abstract class AbstractControllerTest {
protected static final String SYS_ADMIN_EMAIL = "sysadmin@thingsboard.org"; protected static final String SYS_ADMIN_EMAIL = "sysadmin@thingsboard.org";
private static final String SYS_ADMIN_PASSWORD = "sysadmin"; private static final String SYS_ADMIN_PASSWORD = "sysadmin";
protected static final String TENANT_ADMIN_EMAIL = "testtenant@thingsboard.org"; protected static final String TENANT_ADMIN_EMAIL = "testtenant@thingsboard.org";
private static final String TENANT_ADMIN_PASSWORD = "tenant"; private static final String TENANT_ADMIN_PASSWORD = "tenant";
protected static final String CUSTOMER_USER_EMAIL = "testcustomer@thingsboard.org"; protected static final String CUSTOMER_USER_EMAIL = "testcustomer@thingsboard.org";
private static final String CUSTOMER_USER_PASSWORD = "customer"; private static final String CUSTOMER_USER_PASSWORD = "customer";
protected MediaType contentType = new MediaType(MediaType.APPLICATION_JSON.getType(), protected MediaType contentType = new MediaType(MediaType.APPLICATION_JSON.getType(),
MediaType.APPLICATION_JSON.getSubtype(), MediaType.APPLICATION_JSON.getSubtype(),
Charset.forName("utf8")); Charset.forName("utf8"));
protected MockMvc mockMvc; protected MockMvc mockMvc;
protected String token; protected String token;
protected String refreshToken; protected String refreshToken;
protected String username; protected String username;
private TenantId tenantId; private TenantId tenantId;
@SuppressWarnings("rawtypes") @SuppressWarnings("rawtypes")
private HttpMessageConverter mappingJackson2HttpMessageConverter; private HttpMessageConverter mappingJackson2HttpMessageConverter;
@Autowired @Autowired
private WebApplicationContext webApplicationContext; private WebApplicationContext webApplicationContext;
@ -132,7 +132,7 @@ public abstract class AbstractControllerTest {
log.info("Finished test: {}", description.getMethodName()); log.info("Finished test: {}", description.getMethodName());
} }
}; };
@Autowired @Autowired
void setConverters(HttpMessageConverter<?>[] converters) { void setConverters(HttpMessageConverter<?>[] converters) {
@ -144,7 +144,7 @@ public abstract class AbstractControllerTest {
Assert.assertNotNull("the JSON message converter must not be null", Assert.assertNotNull("the JSON message converter must not be null",
this.mappingJackson2HttpMessageConverter); this.mappingJackson2HttpMessageConverter);
} }
@Before @Before
public void setup() throws Exception { public void setup() throws Exception {
log.info("Executing setup"); log.info("Executing setup");
@ -188,7 +188,7 @@ public abstract class AbstractControllerTest {
public void teardown() throws Exception { public void teardown() throws Exception {
log.info("Executing teardown"); log.info("Executing teardown");
loginSysAdmin(); loginSysAdmin();
doDelete("/api/tenant/"+tenantId.getId().toString()) doDelete("/api/tenant/" + tenantId.getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
log.info("Executed teardown"); log.info("Executed teardown");
} }
@ -196,7 +196,7 @@ public abstract class AbstractControllerTest {
protected void loginSysAdmin() throws Exception { protected void loginSysAdmin() throws Exception {
login(SYS_ADMIN_EMAIL, SYS_ADMIN_PASSWORD); login(SYS_ADMIN_EMAIL, SYS_ADMIN_PASSWORD);
} }
protected void loginTenantAdmin() throws Exception { protected void loginTenantAdmin() throws Exception {
login(TENANT_ADMIN_EMAIL, TENANT_ADMIN_PASSWORD); login(TENANT_ADMIN_EMAIL, TENANT_ADMIN_PASSWORD);
} }
@ -204,13 +204,13 @@ public abstract class AbstractControllerTest {
protected void loginCustomerUser() throws Exception { protected void loginCustomerUser() throws Exception {
login(CUSTOMER_USER_EMAIL, CUSTOMER_USER_PASSWORD); login(CUSTOMER_USER_EMAIL, CUSTOMER_USER_PASSWORD);
} }
protected User createUserAndLogin(User user, String password) throws Exception { protected User createUserAndLogin(User user, String password) throws Exception {
User savedUser = doPost("/api/user", user, User.class); User savedUser = doPost("/api/user", user, User.class);
logout(); logout();
doGet("/api/noauth/activate?activateToken={activateToken}", TestMailService.currentActivateToken) doGet("/api/noauth/activate?activateToken={activateToken}", TestMailService.currentActivateToken)
.andExpect(status().isSeeOther()) .andExpect(status().isSeeOther())
.andExpect(header().string(HttpHeaders.LOCATION, "/login/createPassword?activateToken=" + TestMailService.currentActivateToken)); .andExpect(header().string(HttpHeaders.LOCATION, "/login/createPassword?activateToken=" + TestMailService.currentActivateToken));
JsonNode tokenInfo = readResponse(doPost("/api/noauth/activate", "activateToken", TestMailService.currentActivateToken, "password", password).andExpect(status().isOk()), JsonNode.class); JsonNode tokenInfo = readResponse(doPost("/api/noauth/activate", "activateToken", TestMailService.currentActivateToken, "password", password).andExpect(status().isOk()), JsonNode.class);
validateAndSetJwtToken(tokenInfo, user.getEmail()); validateAndSetJwtToken(tokenInfo, user.getEmail());
return savedUser; return savedUser;
@ -247,14 +247,14 @@ public abstract class AbstractControllerTest {
Assert.assertNotNull(token); Assert.assertNotNull(token);
Assert.assertFalse(token.isEmpty()); Assert.assertFalse(token.isEmpty());
int i = token.lastIndexOf('.'); int i = token.lastIndexOf('.');
Assert.assertTrue(i>0); Assert.assertTrue(i > 0);
String withoutSignature = token.substring(0, i+1); String withoutSignature = token.substring(0, i + 1);
Jwt<Header,Claims> jwsClaims = Jwts.parser().parseClaimsJwt(withoutSignature); Jwt<Header, Claims> jwsClaims = Jwts.parser().parseClaimsJwt(withoutSignature);
Claims claims = jwsClaims.getBody(); Claims claims = jwsClaims.getBody();
String subject = claims.getSubject(); String subject = claims.getSubject();
Assert.assertEquals(username, subject); Assert.assertEquals(username, subject);
} }
protected void logout() throws Exception { protected void logout() throws Exception {
this.token = null; this.token = null;
this.refreshToken = null; this.refreshToken = null;
@ -266,24 +266,24 @@ public abstract class AbstractControllerTest {
request.header(ThingsboardSecurityConfiguration.JWT_TOKEN_HEADER_PARAM, "Bearer " + this.token); request.header(ThingsboardSecurityConfiguration.JWT_TOKEN_HEADER_PARAM, "Bearer " + this.token);
} }
} }
protected ResultActions doGet(String urlTemplate, Object... urlVariables) throws Exception { protected ResultActions doGet(String urlTemplate, Object... urlVariables) throws Exception {
MockHttpServletRequestBuilder getRequest = get(urlTemplate, urlVariables); MockHttpServletRequestBuilder getRequest = get(urlTemplate, urlVariables);
setJwtToken(getRequest); setJwtToken(getRequest);
return mockMvc.perform(getRequest); return mockMvc.perform(getRequest);
} }
protected <T> T doGet(String urlTemplate, Class<T> responseClass, Object... urlVariables) throws Exception { protected <T> T doGet(String urlTemplate, Class<T> responseClass, Object... urlVariables) throws Exception {
return readResponse(doGet(urlTemplate, urlVariables).andExpect(status().isOk()), responseClass); return readResponse(doGet(urlTemplate, urlVariables).andExpect(status().isOk()), responseClass);
} }
protected <T> T doGetTyped(String urlTemplate, TypeReference<T> responseType, Object... urlVariables) throws Exception { protected <T> T doGetTyped(String urlTemplate, TypeReference<T> responseType, Object... urlVariables) throws Exception {
return readResponse(doGet(urlTemplate, urlVariables).andExpect(status().isOk()), responseType); return readResponse(doGet(urlTemplate, urlVariables).andExpect(status().isOk()), responseType);
} }
protected <T> T doGetTypedWithPageLink(String urlTemplate, TypeReference<T> responseType, protected <T> T doGetTypedWithPageLink(String urlTemplate, TypeReference<T> responseType,
TextPageLink pageLink, TextPageLink pageLink,
Object... urlVariables) throws Exception { Object... urlVariables) throws Exception {
List<Object> pageLinkVariables = new ArrayList<>(); List<Object> pageLinkVariables = new ArrayList<>();
urlTemplate += "limit={limit}"; urlTemplate += "limit={limit}";
pageLinkVariables.add(pageLink.getLimit()); pageLinkVariables.add(pageLink.getLimit());
@ -299,18 +299,18 @@ public abstract class AbstractControllerTest {
urlTemplate += "&textOffset={textOffset}"; urlTemplate += "&textOffset={textOffset}";
pageLinkVariables.add(pageLink.getTextOffset()); pageLinkVariables.add(pageLink.getTextOffset());
} }
Object[] vars = new Object[urlVariables.length + pageLinkVariables.size()]; Object[] vars = new Object[urlVariables.length + pageLinkVariables.size()];
System.arraycopy(urlVariables, 0, vars, 0, urlVariables.length); System.arraycopy(urlVariables, 0, vars, 0, urlVariables.length);
System.arraycopy(pageLinkVariables.toArray(), 0, vars, urlVariables.length, pageLinkVariables.size()); System.arraycopy(pageLinkVariables.toArray(), 0, vars, urlVariables.length, pageLinkVariables.size());
return readResponse(doGet(urlTemplate, vars).andExpect(status().isOk()), responseType); return readResponse(doGet(urlTemplate, vars).andExpect(status().isOk()), responseType);
} }
protected <T> T doPost(String urlTemplate, Class<T> responseClass, String... params) throws Exception { protected <T> T doPost(String urlTemplate, Class<T> responseClass, String... params) throws Exception {
return readResponse(doPost(urlTemplate, params).andExpect(status().isOk()), responseClass); return readResponse(doPost(urlTemplate, params).andExpect(status().isOk()), responseClass);
} }
protected <T> T doPost(String urlTemplate, T content, Class<T> responseClass, String... params) throws Exception { protected <T> T doPost(String urlTemplate, T content, Class<T> responseClass, String... params) throws Exception {
return readResponse(doPost(urlTemplate, content, params).andExpect(status().isOk()), responseClass); return readResponse(doPost(urlTemplate, content, params).andExpect(status().isOk()), responseClass);
} }
@ -318,15 +318,15 @@ public abstract class AbstractControllerTest {
protected <T> T doDelete(String urlTemplate, Class<T> responseClass, String... params) throws Exception { protected <T> T doDelete(String urlTemplate, Class<T> responseClass, String... params) throws Exception {
return readResponse(doDelete(urlTemplate, params).andExpect(status().isOk()), responseClass); return readResponse(doDelete(urlTemplate, params).andExpect(status().isOk()), responseClass);
} }
protected ResultActions doPost(String urlTemplate, String... params) throws Exception { protected ResultActions doPost(String urlTemplate, String... params) throws Exception {
MockHttpServletRequestBuilder postRequest = post(urlTemplate); MockHttpServletRequestBuilder postRequest = post(urlTemplate);
setJwtToken(postRequest); setJwtToken(postRequest);
populateParams(postRequest, params); populateParams(postRequest, params);
return mockMvc.perform(postRequest); return mockMvc.perform(postRequest);
} }
protected <T> ResultActions doPost(String urlTemplate, T content, String... params) throws Exception { protected <T> ResultActions doPost(String urlTemplate, T content, String... params) throws Exception {
MockHttpServletRequestBuilder postRequest = post(urlTemplate); MockHttpServletRequestBuilder postRequest = post(urlTemplate);
setJwtToken(postRequest); setJwtToken(postRequest);
String json = json(content); String json = json(content);
@ -334,25 +334,25 @@ public abstract class AbstractControllerTest {
populateParams(postRequest, params); populateParams(postRequest, params);
return mockMvc.perform(postRequest); return mockMvc.perform(postRequest);
} }
protected ResultActions doDelete(String urlTemplate, String... params) throws Exception { protected ResultActions doDelete(String urlTemplate, String... params) throws Exception {
MockHttpServletRequestBuilder deleteRequest = delete(urlTemplate); MockHttpServletRequestBuilder deleteRequest = delete(urlTemplate);
setJwtToken(deleteRequest); setJwtToken(deleteRequest);
populateParams(deleteRequest, params); populateParams(deleteRequest, params);
return mockMvc.perform(deleteRequest); return mockMvc.perform(deleteRequest);
} }
protected void populateParams(MockHttpServletRequestBuilder request, String... params) { protected void populateParams(MockHttpServletRequestBuilder request, String... params) {
if (params != null && params.length > 0) { if (params != null && params.length > 0) {
Assert.assertEquals(params.length % 2, 0); Assert.assertEquals(params.length % 2, 0);
MultiValueMap<String, String> paramsMap = new LinkedMultiValueMap<String, String>(); MultiValueMap<String, String> paramsMap = new LinkedMultiValueMap<String, String>();
for (int i=0;i<params.length;i+=2) { for (int i = 0; i < params.length; i += 2) {
paramsMap.add(params[i], params[i+1]); paramsMap.add(params[i], params[i + 1]);
} }
request.params(paramsMap); request.params(paramsMap);
} }
} }
@SuppressWarnings("unchecked") @SuppressWarnings("unchecked")
protected String json(Object o) throws IOException { protected String json(Object o) throws IOException {
MockHttpOutputMessage mockHttpOutputMessage = new MockHttpOutputMessage(); MockHttpOutputMessage mockHttpOutputMessage = new MockHttpOutputMessage();
@ -360,14 +360,14 @@ public abstract class AbstractControllerTest {
o, MediaType.APPLICATION_JSON, mockHttpOutputMessage); o, MediaType.APPLICATION_JSON, mockHttpOutputMessage);
return mockHttpOutputMessage.getBodyAsString(); return mockHttpOutputMessage.getBodyAsString();
} }
@SuppressWarnings("unchecked") @SuppressWarnings("unchecked")
protected <T> T readResponse(ResultActions result, Class<T> responseClass) throws Exception { protected <T> T readResponse(ResultActions result, Class<T> responseClass) throws Exception {
byte[] content = result.andReturn().getResponse().getContentAsByteArray(); byte[] content = result.andReturn().getResponse().getContentAsByteArray();
MockHttpInputMessage mockHttpInputMessage = new MockHttpInputMessage(content); MockHttpInputMessage mockHttpInputMessage = new MockHttpInputMessage(content);
return (T) this.mappingJackson2HttpMessageConverter.read(responseClass, mockHttpInputMessage); return (T) this.mappingJackson2HttpMessageConverter.read(responseClass, mockHttpInputMessage);
} }
protected <T> T readResponse(ResultActions result, TypeReference<T> type) throws Exception { protected <T> T readResponse(ResultActions result, TypeReference<T> type) throws Exception {
byte[] content = result.andReturn().getResponse().getContentAsByteArray(); byte[] content = result.andReturn().getResponse().getContentAsByteArray();
ObjectMapper mapper = new ObjectMapper(); ObjectMapper mapper = new ObjectMapper();

5
dao/src/main/java/org/thingsboard/server/dao/asset/AssetSearchQuery.java → common/data/src/main/java/org/thingsboard/server/common/data/asset/AssetSearchQuery.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and * See the License for the specific language governing permissions and
* limitations under the License. * limitations under the License.
*/ */
package org.thingsboard.server.dao.asset; package org.thingsboard.server.common.data.asset;
import lombok.Data; import lombok.Data;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
@ -22,7 +22,6 @@ import org.thingsboard.server.common.data.relation.EntityRelationsQuery;
import org.thingsboard.server.common.data.relation.EntityTypeFilter; import org.thingsboard.server.common.data.relation.EntityTypeFilter;
import org.thingsboard.server.common.data.relation.RelationsSearchParameters; import org.thingsboard.server.common.data.relation.RelationsSearchParameters;
import javax.annotation.Nullable;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
@ -33,9 +32,7 @@ import java.util.List;
public class AssetSearchQuery { public class AssetSearchQuery {
private RelationsSearchParameters parameters; private RelationsSearchParameters parameters;
@Nullable
private String relationType; private String relationType;
@Nullable
private List<String> assetTypes; private List<String> assetTypes;
public EntityRelationsQuery toEntitySearchQuery() { public EntityRelationsQuery toEntitySearchQuery() {

2
dao/src/main/java/org/thingsboard/server/dao/device/DeviceSearchQuery.java → dao/src/main/java/org/thingsboard/server/common/data/device/DeviceSearchQuery.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and * See the License for the specific language governing permissions and
* limitations under the License. * limitations under the License.
*/ */
package org.thingsboard.server.dao.device; package org.thingsboard.server.common.data.device;
import lombok.Data; import lombok.Data;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;

1
dao/src/main/java/org/thingsboard/server/dao/asset/AssetService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.asset;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetSearchQuery;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;

1
dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java

@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetSearchQuery;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;

1
dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.device;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.device.DeviceSearchQuery;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;

1
dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java

@ -25,6 +25,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
import org.thingsboard.server.common.data.*; import org.thingsboard.server.common.data.*;
import org.thingsboard.server.common.data.device.DeviceSearchQuery;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;

66
transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -16,6 +16,7 @@
package org.thingsboard.server.transport.mqtt; package org.thingsboard.server.transport.mqtt;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.handler.codec.mqtt.*; import io.netty.handler.codec.mqtt.*;
@ -45,6 +46,8 @@ import org.thingsboard.server.transport.mqtt.util.SslUtil;
import javax.net.ssl.SSLPeerUnverifiedException; import javax.net.ssl.SSLPeerUnverifiedException;
import javax.security.cert.X509Certificate; import javax.security.cert.X509Certificate;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
@ -71,6 +74,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private final RelationService relationService; private final RelationService relationService;
private final SslHandler sslHandler; private final SslHandler sslHandler;
private volatile boolean connected; private volatile boolean connected;
private volatile InetSocketAddress address;
private volatile GatewaySessionCtx gatewaySessionCtx; private volatile GatewaySessionCtx gatewaySessionCtx;
public MqttTransportHandler(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, RelationService relationService, public MqttTransportHandler(SessionMsgProcessor processor, DeviceService deviceService, DeviceAuthService authService, RelationService relationService,
@ -94,30 +98,36 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
private void processMqttMsg(ChannelHandlerContext ctx, MqttMessage msg) { private void processMqttMsg(ChannelHandlerContext ctx, MqttMessage msg) {
deviceSessionCtx.setChannel(ctx); address = (InetSocketAddress) ctx.channel().remoteAddress();
switch (msg.fixedHeader().messageType()) { if (msg.fixedHeader() == null) {
case CONNECT: log.info("[{}:{}] Invalid message received", address.getHostName(), address.getPort());
processConnect(ctx, (MqttConnectMessage) msg); processDisconnect(ctx);
break; } else {
case PUBLISH: deviceSessionCtx.setChannel(ctx);
processPublish(ctx, (MqttPublishMessage) msg); switch (msg.fixedHeader().messageType()) {
break; case CONNECT:
case SUBSCRIBE: processConnect(ctx, (MqttConnectMessage) msg);
processSubscribe(ctx, (MqttSubscribeMessage) msg); break;
break; case PUBLISH:
case UNSUBSCRIBE: processPublish(ctx, (MqttPublishMessage) msg);
processUnsubscribe(ctx, (MqttUnsubscribeMessage) msg); break;
break; case SUBSCRIBE:
case PINGREQ: processSubscribe(ctx, (MqttSubscribeMessage) msg);
if (checkConnected(ctx)) { break;
ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); case UNSUBSCRIBE:
} processUnsubscribe(ctx, (MqttUnsubscribeMessage) msg);
break; break;
case DISCONNECT: case PINGREQ:
if (checkConnected(ctx)) { if (checkConnected(ctx)) {
processDisconnect(ctx); ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0)));
} }
break; break;
case DISCONNECT:
if (checkConnected(ctx)) {
processDisconnect(ctx);
}
break;
}
} }
} }
@ -313,9 +323,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private void processDisconnect(ChannelHandlerContext ctx) { private void processDisconnect(ChannelHandlerContext ctx) {
ctx.close(); ctx.close();
processor.process(SessionCloseMsg.onDisconnect(deviceSessionCtx.getSessionId())); if (connected) {
if (gatewaySessionCtx != null) { processor.process(SessionCloseMsg.onDisconnect(deviceSessionCtx.getSessionId()));
gatewaySessionCtx.onGatewayDisconnect(); if (gatewaySessionCtx != null) {
gatewaySessionCtx.onGatewayDisconnect();
}
} }
} }

Loading…
Cancel
Save