Browse Source

Actor message processors NPE fix.

pull/1500/head
Igor Kulikov 8 years ago
parent
commit
adb08a1d7b
  1. 22
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 68
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
  3. 13
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
  4. 2
      application/src/test/java/org/thingsboard/server/mqtt/telemetry/AbstractMqttTelemetryIntegrationTest.java

22
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -118,17 +118,23 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
this.rpcSubscriptions = new HashMap<>(); this.rpcSubscriptions = new HashMap<>();
this.toDeviceRpcPendingMap = new HashMap<>(); this.toDeviceRpcPendingMap = new HashMap<>();
this.toServerRpcPendingMap = new HashMap<>(); this.toServerRpcPendingMap = new HashMap<>();
initAttributes(); if (initAttributes()) {
restoreSessions(); restoreSessions();
}
} }
private void initAttributes() { private boolean initAttributes() {
Device device = systemContext.getDeviceService().findDeviceById(tenantId, deviceId); Device device = systemContext.getDeviceService().findDeviceById(tenantId, deviceId);
this.deviceName = device.getName(); if (device != null) {
this.deviceType = device.getType(); this.deviceName = device.getName();
this.defaultMetaData = new TbMsgMetaData(); this.deviceType = device.getType();
this.defaultMetaData.putValue("deviceName", deviceName); this.defaultMetaData = new TbMsgMetaData();
this.defaultMetaData.putValue("deviceType", deviceType); this.defaultMetaData.putValue("deviceName", deviceName);
this.defaultMetaData.putValue("deviceType", deviceType);
return true;
} else {
return false;
}
} }
void processRpcRequest(ActorContext context, ToDeviceRpcRequestActorMsg msg) { void processRpcRequest(ActorContext context, ToDeviceRpcRequestActorMsg msg) {

68
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java

@ -91,17 +91,19 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
public void start(ActorContext context) { public void start(ActorContext context) {
if (!started) { if (!started) {
RuleChain ruleChain = service.findRuleChainById(tenantId, entityId); RuleChain ruleChain = service.findRuleChainById(tenantId, entityId);
ruleChainName = ruleChain.getName(); if (ruleChain != null) {
List<RuleNode> ruleNodeList = service.getRuleChainNodes(tenantId, entityId); ruleChainName = ruleChain.getName();
log.trace("[{}][{}] Starting rule chain with {} nodes", tenantId, entityId, ruleNodeList.size()); List<RuleNode> ruleNodeList = service.getRuleChainNodes(tenantId, entityId);
// Creating and starting the actors; log.trace("[{}][{}] Starting rule chain with {} nodes", tenantId, entityId, ruleNodeList.size());
for (RuleNode ruleNode : ruleNodeList) { // Creating and starting the actors;
log.trace("[{}][{}] Creating rule node [{}]: {}", entityId, ruleNode.getId(), ruleNode.getName(), ruleNode); for (RuleNode ruleNode : ruleNodeList) {
ActorRef ruleNodeActor = createRuleNodeActor(context, ruleNode); log.trace("[{}][{}] Creating rule node [{}]: {}", entityId, ruleNode.getId(), ruleNode.getName(), ruleNode);
nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode)); ActorRef ruleNodeActor = createRuleNodeActor(context, ruleNode);
nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode));
}
initRoutes(ruleChain, ruleNodeList);
started = true;
} }
initRoutes(ruleChain, ruleNodeList);
started = true;
} else { } else {
onUpdate(context); onUpdate(context);
} }
@ -110,31 +112,33 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
@Override @Override
public void onUpdate(ActorContext context) { public void onUpdate(ActorContext context) {
RuleChain ruleChain = service.findRuleChainById(tenantId, entityId); RuleChain ruleChain = service.findRuleChainById(tenantId, entityId);
ruleChainName = ruleChain.getName(); if (ruleChain != null) {
List<RuleNode> ruleNodeList = service.getRuleChainNodes(tenantId, entityId); ruleChainName = ruleChain.getName();
log.trace("[{}][{}] Updating rule chain with {} nodes", tenantId, entityId, ruleNodeList.size()); List<RuleNode> ruleNodeList = service.getRuleChainNodes(tenantId, entityId);
for (RuleNode ruleNode : ruleNodeList) { log.trace("[{}][{}] Updating rule chain with {} nodes", tenantId, entityId, ruleNodeList.size());
RuleNodeCtx existing = nodeActors.get(ruleNode.getId()); for (RuleNode ruleNode : ruleNodeList) {
if (existing == null) { RuleNodeCtx existing = nodeActors.get(ruleNode.getId());
log.trace("[{}][{}] Creating rule node [{}]: {}", entityId, ruleNode.getId(), ruleNode.getName(), ruleNode); if (existing == null) {
ActorRef ruleNodeActor = createRuleNodeActor(context, ruleNode); log.trace("[{}][{}] Creating rule node [{}]: {}", entityId, ruleNode.getId(), ruleNode.getName(), ruleNode);
nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode)); ActorRef ruleNodeActor = createRuleNodeActor(context, ruleNode);
} else { nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode));
log.trace("[{}][{}] Updating rule node [{}]: {}", entityId, ruleNode.getId(), ruleNode.getName(), ruleNode); } else {
existing.setSelf(ruleNode); log.trace("[{}][{}] Updating rule node [{}]: {}", entityId, ruleNode.getId(), ruleNode.getName(), ruleNode);
existing.getSelfActor().tell(new ComponentLifecycleMsg(tenantId, existing.getSelf().getId(), ComponentLifecycleEvent.UPDATED), self); existing.setSelf(ruleNode);
existing.getSelfActor().tell(new ComponentLifecycleMsg(tenantId, existing.getSelf().getId(), ComponentLifecycleEvent.UPDATED), self);
}
} }
}
Set<RuleNodeId> existingNodes = ruleNodeList.stream().map(RuleNode::getId).collect(Collectors.toSet()); Set<RuleNodeId> existingNodes = ruleNodeList.stream().map(RuleNode::getId).collect(Collectors.toSet());
List<RuleNodeId> removedRules = nodeActors.keySet().stream().filter(node -> !existingNodes.contains(node)).collect(Collectors.toList()); List<RuleNodeId> removedRules = nodeActors.keySet().stream().filter(node -> !existingNodes.contains(node)).collect(Collectors.toList());
removedRules.forEach(ruleNodeId -> { removedRules.forEach(ruleNodeId -> {
log.trace("[{}][{}] Removing rule node [{}]", tenantId, entityId, ruleNodeId); log.trace("[{}][{}] Removing rule node [{}]", tenantId, entityId, ruleNodeId);
RuleNodeCtx removed = nodeActors.remove(ruleNodeId); RuleNodeCtx removed = nodeActors.remove(ruleNodeId);
removed.getSelfActor().tell(new ComponentLifecycleMsg(tenantId, removed.getSelf().getId(), ComponentLifecycleEvent.DELETED), self); removed.getSelfActor().tell(new ComponentLifecycleMsg(tenantId, removed.getSelf().getId(), ComponentLifecycleEvent.DELETED), self);
}); });
initRoutes(ruleChain, ruleNodeList); initRoutes(ruleChain, ruleNodeList);
}
} }
@Override @Override

13
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java

@ -55,7 +55,9 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
@Override @Override
public void start(ActorContext context) throws Exception { public void start(ActorContext context) throws Exception {
tbNode = initComponent(ruleNode); tbNode = initComponent(ruleNode);
state = ComponentLifecycleState.ACTIVE; if (tbNode != null) {
state = ComponentLifecycleState.ACTIVE;
}
} }
@Override @Override
@ -118,9 +120,12 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
} }
private TbNode initComponent(RuleNode ruleNode) throws Exception { private TbNode initComponent(RuleNode ruleNode) throws Exception {
Class<?> componentClazz = Class.forName(ruleNode.getType()); TbNode tbNode = null;
TbNode tbNode = (TbNode) (componentClazz.newInstance()); if (ruleNode != null) {
tbNode.init(defaultCtx, new TbNodeConfiguration(ruleNode.getConfiguration())); Class<?> componentClazz = Class.forName(ruleNode.getType());
tbNode = (TbNode) (componentClazz.newInstance());
tbNode.init(defaultCtx, new TbNodeConfiguration(ruleNode.getConfiguration()));
}
return tbNode; return tbNode;
} }

2
application/src/test/java/org/thingsboard/server/mqtt/telemetry/AbstractMqttTelemetryIntegrationTest.java

@ -111,7 +111,7 @@ public abstract class AbstractMqttTelemetryIntegrationTest extends AbstractContr
client.subscribe("v1/devices/me/attributes", MqttQoS.AT_MOST_ONCE.value()); client.subscribe("v1/devices/me/attributes", MqttQoS.AT_MOST_ONCE.value());
String payload = "{\"key\":\"value\"}"; String payload = "{\"key\":\"value\"}";
String result = doPostAsync("/api/plugins/telemetry/" + savedDevice.getId() + "/SHARED_SCOPE", payload, String.class, status().isOk()); String result = doPostAsync("/api/plugins/telemetry/" + savedDevice.getId() + "/SHARED_SCOPE", payload, String.class, status().isOk());
latch.await(3, TimeUnit.SECONDS); latch.await(10, TimeUnit.SECONDS);
assertEquals(payload, callback.getPayload()); assertEquals(payload, callback.getPayload());
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS());
} }

Loading…
Cancel
Save