120 changed files with 957 additions and 662 deletions
@ -0,0 +1,51 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.sql.relation; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.data.jpa.repository.Modifying; |
|||
import org.thingsboard.server.dao.model.sql.RelationEntity; |
|||
|
|||
import javax.persistence.EntityManager; |
|||
import javax.persistence.PersistenceContext; |
|||
import javax.persistence.Query; |
|||
|
|||
@Slf4j |
|||
public abstract class AbstractRelationInsertRepository implements RelationInsertRepository { |
|||
|
|||
@PersistenceContext |
|||
protected EntityManager entityManager; |
|||
|
|||
protected Query getQuery(RelationEntity entity, String query) { |
|||
Query nativeQuery = entityManager.createNativeQuery(query, RelationEntity.class); |
|||
if (entity.getAdditionalInfo() == null) { |
|||
nativeQuery.setParameter("additionalInfo", null); |
|||
} else { |
|||
nativeQuery.setParameter("additionalInfo", entity.getAdditionalInfo().toString()); |
|||
} |
|||
return nativeQuery |
|||
.setParameter("fromId", entity.getFromId()) |
|||
.setParameter("fromType", entity.getFromType()) |
|||
.setParameter("toId", entity.getToId()) |
|||
.setParameter("toType", entity.getToType()) |
|||
.setParameter("relationTypeGroup", entity.getRelationTypeGroup()) |
|||
.setParameter("relationType", entity.getRelationType()); |
|||
} |
|||
|
|||
@Modifying |
|||
protected abstract RelationEntity processSaveOrUpdate(RelationEntity entity); |
|||
|
|||
} |
|||
@ -0,0 +1,47 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.sql.relation; |
|||
|
|||
import org.springframework.stereotype.Repository; |
|||
import org.springframework.transaction.annotation.Transactional; |
|||
import org.thingsboard.server.dao.model.sql.RelationCompositeKey; |
|||
import org.thingsboard.server.dao.model.sql.RelationEntity; |
|||
import org.thingsboard.server.dao.util.HsqlDao; |
|||
import org.thingsboard.server.dao.util.SqlDao; |
|||
|
|||
@HsqlDao |
|||
@SqlDao |
|||
@Repository |
|||
@Transactional |
|||
public class HsqlRelationInsertRepository extends AbstractRelationInsertRepository implements RelationInsertRepository { |
|||
|
|||
private static final String INSERT_ON_CONFLICT_DO_UPDATE = "MERGE INTO relation USING (VALUES :fromId, :fromType, :toId, :toType, :relationTypeGroup, :relationType, :additionalInfo) R " + |
|||
"(from_id, from_type, to_id, to_type, relation_type_group, relation_type, additional_info) " + |
|||
"ON (relation.from_id = R.from_id AND relation.from_type = R.from_type AND relation.relation_type_group = R.relation_type_group AND relation.relation_type = R.relation_type AND relation.to_id = R.to_id AND relation.to_type = R.to_type) " + |
|||
"WHEN MATCHED THEN UPDATE SET relation.additional_info = R.additional_info " + |
|||
"WHEN NOT MATCHED THEN INSERT (from_id, from_type, to_id, to_type, relation_type_group, relation_type, additional_info) VALUES (R.from_id, R.from_type, R.to_id, R.to_type, R.relation_type_group, R.relation_type, R.additional_info)"; |
|||
|
|||
@Override |
|||
public RelationEntity saveOrUpdate(RelationEntity entity) { |
|||
return processSaveOrUpdate(entity); |
|||
} |
|||
|
|||
@Override |
|||
protected RelationEntity processSaveOrUpdate(RelationEntity entity) { |
|||
getQuery(entity, INSERT_ON_CONFLICT_DO_UPDATE).executeUpdate(); |
|||
return entityManager.find(RelationEntity.class, new RelationCompositeKey(entity.toData())); |
|||
} |
|||
} |
|||
@ -0,0 +1,43 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.sql.relation; |
|||
|
|||
import org.springframework.stereotype.Repository; |
|||
import org.springframework.transaction.annotation.Transactional; |
|||
import org.thingsboard.server.dao.model.sql.RelationEntity; |
|||
import org.thingsboard.server.dao.util.PsqlDao; |
|||
import org.thingsboard.server.dao.util.SqlDao; |
|||
|
|||
@PsqlDao |
|||
@SqlDao |
|||
@Repository |
|||
@Transactional |
|||
public class PsqlRelationInsertRepository extends AbstractRelationInsertRepository implements RelationInsertRepository { |
|||
|
|||
private static final String INSERT_ON_CONFLICT_DO_UPDATE = "INSERT INTO relation (from_id, from_type, to_id, to_type, relation_type_group, relation_type, additional_info)" + |
|||
" VALUES (:fromId, :fromType, :toId, :toType, :relationTypeGroup, :relationType, :additionalInfo) " + |
|||
"ON CONFLICT (from_id, from_type, relation_type_group, relation_type, to_id, to_type) DO UPDATE SET additional_info = :additionalInfo returning *"; |
|||
|
|||
@Override |
|||
public RelationEntity saveOrUpdate(RelationEntity entity) { |
|||
return processSaveOrUpdate(entity); |
|||
} |
|||
|
|||
@Override |
|||
protected RelationEntity processSaveOrUpdate(RelationEntity entity) { |
|||
return (RelationEntity) getQuery(entity, INSERT_ON_CONFLICT_DO_UPDATE).getSingleResult(); |
|||
} |
|||
} |
|||
@ -0,0 +1,24 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.dao.sql.relation; |
|||
|
|||
import org.thingsboard.server.dao.model.sql.RelationEntity; |
|||
|
|||
public interface RelationInsertRepository { |
|||
|
|||
RelationEntity saveOrUpdate(RelationEntity entity); |
|||
|
|||
} |
|||
@ -1,17 +0,0 @@ |
|||
-- |
|||
-- Copyright © 2016-2020 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. |
|||
-- |
|||
|
|||
CREATE INDEX IF NOT EXISTS idx_tenant_ts_kv ON tenant_ts_kv(tenant_id, entity_id, key, ts); |
|||
@ -0,0 +1,24 @@ |
|||
#!/bin/bash |
|||
# |
|||
# Copyright © 2016-2020 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. |
|||
# |
|||
|
|||
mkdir -p tb-node/log/ && sudo chown -R 799:799 tb-node/log/ |
|||
|
|||
mkdir -p tb-transports/coap/log && sudo chown -R 799:799 tb-transports/coap/log |
|||
|
|||
mkdir -p tb-transports/http/log && sudo chown -R 799:799 tb-transports/http/log |
|||
|
|||
mkdir -p tb-transports/mqtt/log && sudo chown -R 799:799 tb-transports/mqtt/log |
|||
@ -0,0 +1,100 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.rule.engine.filter; |
|||
|
|||
import com.fasterxml.jackson.databind.ObjectMapper; |
|||
import com.google.common.util.concurrent.FutureCallback; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import com.google.common.util.concurrent.MoreExecutors; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.rule.engine.api.RuleNode; |
|||
import org.thingsboard.rule.engine.api.TbContext; |
|||
import org.thingsboard.rule.engine.api.TbNode; |
|||
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|||
import org.thingsboard.rule.engine.api.TbNodeException; |
|||
import org.thingsboard.rule.engine.api.util.TbNodeUtils; |
|||
import org.thingsboard.server.common.data.alarm.Alarm; |
|||
import org.thingsboard.server.common.data.alarm.AlarmStatus; |
|||
import org.thingsboard.server.common.data.plugin.ComponentType; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
|
|||
import javax.annotation.Nullable; |
|||
import java.io.IOException; |
|||
|
|||
@Slf4j |
|||
@RuleNode( |
|||
type = ComponentType.FILTER, |
|||
name = "checks alarm status", |
|||
configClazz = TbCheckAlarmStatusNodeConfig.class, |
|||
relationTypes = {"True", "False"}, |
|||
nodeDescription = "Checks alarm status.", |
|||
nodeDetails = "If the alarm status matches the specified one - msg is success if does not match - msg is failure.", |
|||
uiResources = {"static/rulenode/rulenode-core-config.js"}, |
|||
configDirective = "tbFilterNodeCheckAlarmStatusConfig") |
|||
public class TbCheckAlarmStatusNode implements TbNode { |
|||
private TbCheckAlarmStatusNodeConfig config; |
|||
private final ObjectMapper mapper = new ObjectMapper(); |
|||
|
|||
@Override |
|||
public void init(TbContext tbContext, TbNodeConfiguration configuration) throws TbNodeException { |
|||
this.config = TbNodeUtils.convert(configuration, TbCheckAlarmStatusNodeConfig.class); |
|||
} |
|||
|
|||
@Override |
|||
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException { |
|||
try { |
|||
Alarm alarm = mapper.readValue(msg.getData(), Alarm.class); |
|||
|
|||
ListenableFuture<Alarm> latest = ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), alarm.getId()); |
|||
|
|||
Futures.addCallback(latest, new FutureCallback<Alarm>() { |
|||
@Override |
|||
public void onSuccess(@Nullable Alarm result) { |
|||
if (result != null) { |
|||
boolean isPresent = false; |
|||
for (AlarmStatus alarmStatus : config.getAlarmStatusList()) { |
|||
if (alarm.getStatus() == alarmStatus) { |
|||
isPresent = true; |
|||
break; |
|||
} |
|||
} |
|||
|
|||
if (isPresent) { |
|||
ctx.tellNext(msg, "True"); |
|||
} else { |
|||
ctx.tellNext(msg, "False"); |
|||
} |
|||
} else { |
|||
ctx.tellFailure(msg, new TbNodeException("No such Alarm found.")); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onFailure(Throwable t) { |
|||
ctx.tellFailure(msg, t); |
|||
} |
|||
}, MoreExecutors.directExecutor()); |
|||
} catch (IOException e) { |
|||
log.error("Failed to parse alarm: [{}]", msg.getData()); |
|||
throw new TbNodeException(e); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void destroy() { |
|||
} |
|||
} |
|||
@ -0,0 +1,35 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.rule.engine.filter; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.rule.engine.api.NodeConfiguration; |
|||
import org.thingsboard.server.common.data.alarm.AlarmStatus; |
|||
|
|||
import java.util.Arrays; |
|||
import java.util.List; |
|||
|
|||
@Data |
|||
public class TbCheckAlarmStatusNodeConfig implements NodeConfiguration<TbCheckAlarmStatusNodeConfig> { |
|||
private List<AlarmStatus> alarmStatusList; |
|||
|
|||
@Override |
|||
public TbCheckAlarmStatusNodeConfig defaultConfiguration() { |
|||
TbCheckAlarmStatusNodeConfig config = new TbCheckAlarmStatusNodeConfig(); |
|||
config.setAlarmStatusList(Arrays.asList(AlarmStatus.ACTIVE_ACK, AlarmStatus.ACTIVE_UNACK)); |
|||
return config; |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue