Browse Source

Merge branch 'master' of github.com:thingsboard/thingsboard into improvement/ruleNodes

pull/1479/head
ShvaykaD 8 years ago
parent
commit
b33b682cc6
  1. 5
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 6
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/executors/AbstractListeningExecutor.java
  4. 48
      application/src/main/java/org/thingsboard/server/service/executors/SharedEventLoopGroupService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java
  6. 6
      application/src/main/resources/thingsboard.yml
  7. 4
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  8. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java
  9. 2
      ui/package-lock.json
  10. 12
      ui/src/app/widget/lib/map-widget2.js

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

@ -67,6 +67,7 @@ import org.thingsboard.server.service.encoding.DataDecodingEncodingService;
import org.thingsboard.server.service.executors.ClusterRpcCallbackExecutorService;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.executors.ExternalCallExecutorService;
import org.thingsboard.server.service.executors.SharedEventLoopGroupService;
import org.thingsboard.server.service.mail.MailExecutorService;
import org.thingsboard.server.service.rpc.DeviceRpcService;
import org.thingsboard.server.service.script.JsExecutorService;
@ -206,6 +207,10 @@ public class ActorSystemContext {
@Getter
private ExternalCallExecutorService externalCallExecutorService;
@Autowired
@Getter
private SharedEventLoopGroupService sharedEventLoopGroupService;
@Autowired
@Getter
private MailService mailService;

6
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -17,6 +17,7 @@ package org.thingsboard.server.actors.ruleChain;
import akka.actor.ActorRef;
import com.datastax.driver.core.utils.UUIDs;
import io.netty.channel.EventLoopGroup;
import org.springframework.util.StringUtils;
import org.thingsboard.rule.engine.api.ListeningExecutor;
import org.thingsboard.rule.engine.api.MailService;
@ -238,6 +239,11 @@ class DefaultTbContext implements TbContext {
return mainCtx.getRuleChainTransactionService();
}
@Override
public EventLoopGroup getSharedEventLoop() {
return mainCtx.getSharedEventLoopGroupService().getSharedEventLoopGroup();
}
@Override
public MailService getMailService() {
if (mainCtx.isAllowSystemMailService()) {

2
application/src/main/java/org/thingsboard/server/service/executors/AbstractListeningExecutor.java

@ -34,7 +34,7 @@ public abstract class AbstractListeningExecutor implements ListeningExecutor {
@PostConstruct
public void init() {
this.service = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(getThreadPollSize()));
this.service = MoreExecutors.listeningDecorator(Executors.newWorkStealingPool(getThreadPollSize()));
}
@PreDestroy

48
application/src/main/java/org/thingsboard/server/service/executors/SharedEventLoopGroupService.java

@ -0,0 +1,48 @@
/**
* Copyright © 2016-2019 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.executors;
import com.google.common.util.concurrent.MoreExecutors;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import lombok.Getter;
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
@Component
public class SharedEventLoopGroupService {
@Getter
private EventLoopGroup sharedEventLoopGroup;
@PostConstruct
public void init() {
this.sharedEventLoopGroup = new NioEventLoopGroup();
}
@PreDestroy
public void destroy() {
if (this.sharedEventLoopGroup != null) {
this.sharedEventLoopGroup.shutdownGracefully(0, 5, TimeUnit.SECONDS);
}
}
}

2
application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java

@ -42,7 +42,7 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer
public void init() {
if (useJsSandbox()) {
sandbox = NashornSandboxes.create();
monitorExecutorService = Executors.newFixedThreadPool(getMonitorThreadPoolSize());
monitorExecutorService = Executors.newWorkStealingPool(getMonitorThreadPoolSize());
sandbox.setExecutor(monitorExecutorService);
sandbox.setMaxCPUTime(getMaxCpuTime());
sandbox.allowNoBraces(false);

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

@ -200,13 +200,13 @@ actors:
# Specify thread pool size for database request callbacks executor service
db_callback_thread_pool_size: "${ACTORS_RULE_DB_CALLBACK_THREAD_POOL_SIZE:1}"
# Specify thread pool size for javascript executor service
js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:10}"
js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:50}"
# Specify thread pool size for mail sender executor service
mail_thread_pool_size: "${ACTORS_RULE_MAIL_THREAD_POOL_SIZE:10}"
mail_thread_pool_size: "${ACTORS_RULE_MAIL_THREAD_POOL_SIZE:50}"
# Whether to allow usage of system mail service for rules
allow_system_mail_service: "${ACTORS_RULE_ALLOW_SYSTEM_MAIL_SERVICE:true}"
# Specify thread pool size for external call service
external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:10}"
external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:50}"
chain:
# Errors for particular actor are persisted once per specified amount of milliseconds
error_persist_frequency: "${ACTORS_RULE_CHAIN_ERROR_FREQUENCY:3000}"

4
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.rule.engine.api;
import io.netty.channel.EventLoopGroup;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId;
@ -105,4 +106,7 @@ public interface TbContext {
String getNodeId();
RuleChainTransactionService getRuleChainTransactionService();
EventLoopGroup getSharedEventLoop();
}

11
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java

@ -59,15 +59,13 @@ public class TbMqttNode implements TbNode {
private TbMqttNodeConfiguration config;
private EventLoopGroup eventLoopGroup;
private MqttClient mqttClient;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
try {
this.config = TbNodeUtils.convert(configuration, TbMqttNodeConfiguration.class);
this.eventLoopGroup = new NioEventLoopGroup();
this.mqttClient = initClient();
this.mqttClient = initClient(ctx);
} catch (Exception e) {
throw new TbNodeException(e);
}
@ -99,12 +97,9 @@ public class TbMqttNode implements TbNode {
if (this.mqttClient != null) {
this.mqttClient.disconnect();
}
if (this.eventLoopGroup != null) {
this.eventLoopGroup.shutdownGracefully(0, 5, TimeUnit.SECONDS);
}
}
private MqttClient initClient() throws Exception {
private MqttClient initClient(TbContext ctx) throws Exception {
Optional<SslContext> sslContextOpt = initSslContext();
MqttClientConfig config = sslContextOpt.isPresent() ? new MqttClientConfig(sslContextOpt.get()) : new MqttClientConfig();
if (!StringUtils.isEmpty(this.config.getClientId())) {
@ -113,7 +108,7 @@ public class TbMqttNode implements TbNode {
config.setCleanSession(this.config.isCleanSession());
this.config.getCredentials().configure(config);
MqttClient client = MqttClient.create(config, null);
client.setEventLoop(this.eventLoopGroup);
client.setEventLoop(ctx.getSharedEventLoop());
Future<MqttConnectResult> connectFuture = client.connect(this.config.getHost(), this.config.getPort());
MqttConnectResult result;
try {

2
ui/package-lock.json

@ -1,6 +1,6 @@
{
"name": "thingsboard",
"version": "2.3.0",
"version": "2.3.1",
"lockfileVersion": 1,
"requires": true,
"dependencies": {

12
ui/src/app/widget/lib/map-widget2.js

@ -305,12 +305,9 @@ export default class TbMapWidgetV2 {
}
function updateLocationPolygonColor(location, color) {
if (!location.settings.calculatedPolygonColor || location.settings.calculatedPolygonColor !== color) {
if (location.polygon && color) {
location.settings.calculatedPolygonColor = color;
if (location.polygon) {
tbMap.map.updatePolygonColor(location.polygon, location.settings, color);
}
tbMap.map.updatePolygonColor(location.polygon, location.settings, color);
}
}
@ -338,10 +335,8 @@ export default class TbMapWidgetV2 {
function updateLocationStyle(location, dataMap) {
updateLocationLabel(location, dataMap);
var color = calculateLocationColor(location, dataMap);
var polygonColor = calculateLocationPolygonColor(location, dataMap);
var image = calculateLocationMarkerImage(location, dataMap);
updateLocationColor(location, color, image);
if (location.settings.usePolygonColorFunction) updateLocationPolygonColor(location, polygonColor);
updateLocationMarkerIcon(location, image);
}
@ -441,6 +436,7 @@ export default class TbMapWidgetV2 {
if (location.marker) {
updateLocationStyle(location, dataMap);
}
}
}
return locationChanged;
@ -456,11 +452,13 @@ export default class TbMapWidgetV2 {
locationPolygonClick(event, location);
}, [location.dsIndex]);
tbMap.polygons.push(location.polygon);
if (location.settings.usePolygonColorFunction) updateLocationPolygonColor(location, calculateLocationPolygonColor(location, dataMap));
} else if (polygonLatLngs.length > 0) {
let prevPolygonArr = tbMap.map.getPolygonLatLngs(location.polygon);
if (!prevPolygonArr || !arraysEqual(prevPolygonArr, polygonLatLngs)) {
tbMap.map.setPolygonLatLngs(location.polygon, polygonLatLngs);
}
if (location.settings.usePolygonColorFunction) updateLocationPolygonColor(location, calculateLocationPolygonColor(location, dataMap));
}
}
}

Loading…
Cancel
Save