|
|
@ -16,6 +16,7 @@ |
|
|
package org.thingsboard.rule.engine.action; |
|
|
package org.thingsboard.rule.engine.action; |
|
|
|
|
|
|
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
|
|
|
import com.fasterxml.jackson.databind.ObjectMapper; |
|
|
import com.google.common.base.Function; |
|
|
import com.google.common.base.Function; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
@ -31,6 +32,8 @@ import org.thingsboard.server.common.data.id.TenantId; |
|
|
import org.thingsboard.server.common.data.plugin.ComponentType; |
|
|
import org.thingsboard.server.common.data.plugin.ComponentType; |
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
|
|
|
|
|
|
|
|
|
import java.io.IOException; |
|
|
|
|
|
|
|
|
@Slf4j |
|
|
@Slf4j |
|
|
@RuleNode( |
|
|
@RuleNode( |
|
|
type = ComponentType.ACTION, |
|
|
type = ComponentType.ACTION, |
|
|
@ -49,6 +52,8 @@ import org.thingsboard.server.common.msg.TbMsg; |
|
|
) |
|
|
) |
|
|
public class TbCreateAlarmNode extends TbAbstractAlarmNode<TbCreateAlarmNodeConfiguration> { |
|
|
public class TbCreateAlarmNode extends TbAbstractAlarmNode<TbCreateAlarmNodeConfiguration> { |
|
|
|
|
|
|
|
|
|
|
|
private static ObjectMapper mapper = new ObjectMapper(); |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
protected TbCreateAlarmNodeConfiguration loadAlarmNodeConfig(TbNodeConfiguration configuration) throws TbNodeException { |
|
|
protected TbCreateAlarmNodeConfiguration loadAlarmNodeConfig(TbNodeConfiguration configuration) throws TbNodeException { |
|
|
return TbNodeUtils.convert(configuration, TbCreateAlarmNodeConfiguration.class); |
|
|
return TbNodeUtils.convert(configuration, TbCreateAlarmNodeConfiguration.class); |
|
|
@ -56,32 +61,59 @@ public class TbCreateAlarmNode extends TbAbstractAlarmNode<TbCreateAlarmNodeConf |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
protected ListenableFuture<AlarmResult> processAlarm(TbContext ctx, TbMsg msg) { |
|
|
protected ListenableFuture<AlarmResult> processAlarm(TbContext ctx, TbMsg msg) { |
|
|
ListenableFuture<Alarm> latest = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), config.getAlarmType()); |
|
|
String alarmType; |
|
|
return Futures.transformAsync(latest, a -> { |
|
|
final Alarm msgAlarm; |
|
|
if (a == null || a.getStatus().isCleared()) { |
|
|
|
|
|
return createNewAlarm(ctx, msg); |
|
|
if (!config.isUseMessageAlarmData()) { |
|
|
|
|
|
alarmType = config.getAlarmType(); |
|
|
|
|
|
msgAlarm = null; |
|
|
|
|
|
} else { |
|
|
|
|
|
try { |
|
|
|
|
|
msgAlarm = mapper.readValue(msg.getData(), Alarm.class); |
|
|
|
|
|
msgAlarm.setTenantId(ctx.getTenantId()); |
|
|
|
|
|
alarmType = msgAlarm.getType(); |
|
|
|
|
|
} catch (IOException e) { |
|
|
|
|
|
ctx.tellFailure(msg, e); |
|
|
|
|
|
return null; |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
ListenableFuture<Alarm> latest = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), alarmType); |
|
|
|
|
|
return Futures.transformAsync(latest, existingAlarm -> { |
|
|
|
|
|
if (existingAlarm == null || existingAlarm.getStatus().isCleared()) { |
|
|
|
|
|
return createNewAlarm(ctx, msg, msgAlarm); |
|
|
} else { |
|
|
} else { |
|
|
return updateAlarm(ctx, msg, a); |
|
|
return updateAlarm(ctx, msg, existingAlarm, msgAlarm); |
|
|
} |
|
|
} |
|
|
}, ctx.getDbCallbackExecutor()); |
|
|
}, ctx.getDbCallbackExecutor()); |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<AlarmResult> createNewAlarm(TbContext ctx, TbMsg msg) { |
|
|
private ListenableFuture<AlarmResult> createNewAlarm(TbContext ctx, TbMsg msg, Alarm msgAlarm) { |
|
|
ListenableFuture<Alarm> asyncAlarm = Futures.transform(buildAlarmDetails(ctx, msg, null), |
|
|
ListenableFuture<Alarm> asyncAlarm; |
|
|
details -> buildAlarm(msg, details, ctx.getTenantId())); |
|
|
if (msgAlarm != null ) { |
|
|
|
|
|
asyncAlarm = Futures.immediateCheckedFuture(msgAlarm); |
|
|
|
|
|
} else { |
|
|
|
|
|
asyncAlarm = Futures.transform(buildAlarmDetails(ctx, msg, null), |
|
|
|
|
|
details -> buildAlarm(msg, details, ctx.getTenantId())); |
|
|
|
|
|
} |
|
|
ListenableFuture<Alarm> asyncCreated = Futures.transform(asyncAlarm, |
|
|
ListenableFuture<Alarm> asyncCreated = Futures.transform(asyncAlarm, |
|
|
alarm -> ctx.getAlarmService().createOrUpdateAlarm(alarm), ctx.getDbCallbackExecutor()); |
|
|
alarm -> ctx.getAlarmService().createOrUpdateAlarm(alarm), ctx.getDbCallbackExecutor()); |
|
|
return Futures.transform(asyncCreated, alarm -> new AlarmResult(true, false, false, alarm)); |
|
|
return Futures.transform(asyncCreated, alarm -> new AlarmResult(true, false, false, alarm)); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<AlarmResult> updateAlarm(TbContext ctx, TbMsg msg, Alarm alarm) { |
|
|
private ListenableFuture<AlarmResult> updateAlarm(TbContext ctx, TbMsg msg, Alarm existingAlarm, Alarm msgAlarm) { |
|
|
ListenableFuture<Alarm> asyncUpdated = Futures.transform(buildAlarmDetails(ctx, msg, alarm.getDetails()), (Function<JsonNode, Alarm>) details -> { |
|
|
ListenableFuture<Alarm> asyncUpdated = Futures.transform(buildAlarmDetails(ctx, msg, existingAlarm.getDetails()), (Function<JsonNode, Alarm>) details -> { |
|
|
alarm.setSeverity(config.getSeverity()); |
|
|
if (msgAlarm != null) { |
|
|
alarm.setPropagate(config.isPropagate()); |
|
|
existingAlarm.setSeverity(msgAlarm.getSeverity()); |
|
|
alarm.setDetails(details); |
|
|
existingAlarm.setPropagate(msgAlarm.isPropagate()); |
|
|
alarm.setEndTs(System.currentTimeMillis()); |
|
|
} else { |
|
|
return ctx.getAlarmService().createOrUpdateAlarm(alarm); |
|
|
existingAlarm.setSeverity(config.getSeverity()); |
|
|
|
|
|
existingAlarm.setPropagate(config.isPropagate()); |
|
|
|
|
|
} |
|
|
|
|
|
existingAlarm.setDetails(details); |
|
|
|
|
|
existingAlarm.setEndTs(System.currentTimeMillis()); |
|
|
|
|
|
return ctx.getAlarmService().createOrUpdateAlarm(existingAlarm); |
|
|
}, ctx.getDbCallbackExecutor()); |
|
|
}, ctx.getDbCallbackExecutor()); |
|
|
|
|
|
|
|
|
return Futures.transform(asyncUpdated, a -> new AlarmResult(false, true, false, a)); |
|
|
return Futures.transform(asyncUpdated, a -> new AlarmResult(false, true, false, a)); |
|
|
|