409 changed files with 12956 additions and 5511 deletions
@ -0,0 +1,75 @@ |
|||||
|
-- |
||||
|
-- Copyright © 2016-2022 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. |
||||
|
-- |
||||
|
|
||||
|
DO |
||||
|
$$ |
||||
|
DECLARE table_partition RECORD; |
||||
|
BEGIN |
||||
|
-- in case of running the upgrade script a second time: |
||||
|
IF NOT (SELECT exists(SELECT FROM pg_tables WHERE tablename = 'old_audit_log')) THEN |
||||
|
ALTER TABLE audit_log RENAME TO old_audit_log; |
||||
|
ALTER INDEX IF EXISTS idx_audit_log_tenant_id_and_created_time RENAME TO idx_old_audit_log_tenant_id_and_created_time; |
||||
|
|
||||
|
FOR table_partition IN SELECT tablename AS name, split_part(tablename, '_', 3) AS partition_ts |
||||
|
FROM pg_tables WHERE tablename LIKE 'audit_log_%' |
||||
|
LOOP |
||||
|
EXECUTE format('ALTER TABLE %s RENAME TO old_audit_log_%s', table_partition.name, table_partition.partition_ts); |
||||
|
END LOOP; |
||||
|
ELSE |
||||
|
RAISE NOTICE 'Table old_audit_log already exists, leaving as is'; |
||||
|
END IF; |
||||
|
END; |
||||
|
$$; |
||||
|
|
||||
|
CREATE TABLE IF NOT EXISTS audit_log ( |
||||
|
id uuid NOT NULL, |
||||
|
created_time bigint NOT NULL, |
||||
|
tenant_id uuid, |
||||
|
customer_id uuid, |
||||
|
entity_id uuid, |
||||
|
entity_type varchar(255), |
||||
|
entity_name varchar(255), |
||||
|
user_id uuid, |
||||
|
user_name varchar(255), |
||||
|
action_type varchar(255), |
||||
|
action_data varchar(1000000), |
||||
|
action_status varchar(255), |
||||
|
action_failure_details varchar(1000000) |
||||
|
) PARTITION BY RANGE (created_time); |
||||
|
CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_id_and_created_time ON audit_log(tenant_id, created_time DESC); |
||||
|
|
||||
|
CREATE OR REPLACE PROCEDURE migrate_audit_logs(IN start_time_ms BIGINT, IN end_time_ms BIGINT, IN partition_size_ms BIGINT) |
||||
|
LANGUAGE plpgsql AS |
||||
|
$$ |
||||
|
DECLARE |
||||
|
p RECORD; |
||||
|
partition_end_ts BIGINT; |
||||
|
BEGIN |
||||
|
FOR p IN SELECT DISTINCT (created_time - created_time % partition_size_ms) AS partition_ts FROM old_audit_log |
||||
|
WHERE created_time >= start_time_ms AND created_time < end_time_ms |
||||
|
LOOP |
||||
|
partition_end_ts = p.partition_ts + partition_size_ms; |
||||
|
RAISE NOTICE '[audit_log] Partition to create : [%-%]', p.partition_ts, partition_end_ts; |
||||
|
EXECUTE format('CREATE TABLE IF NOT EXISTS audit_log_%s PARTITION OF audit_log ' || |
||||
|
'FOR VALUES FROM ( %s ) TO ( %s )', p.partition_ts, p.partition_ts, partition_end_ts); |
||||
|
END LOOP; |
||||
|
|
||||
|
INSERT INTO audit_log |
||||
|
SELECT id, created_time, tenant_id, customer_id, entity_id, entity_type, entity_name, user_id, user_name, action_type, action_data, action_status, action_failure_details |
||||
|
FROM old_audit_log |
||||
|
WHERE created_time >= start_time_ms AND created_time < end_time_ms; |
||||
|
END; |
||||
|
$$; |
||||
@ -0,0 +1,46 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.edge.rpc.constructor; |
||||
|
|
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.edge.Edge; |
||||
|
import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
|
||||
|
@Component |
||||
|
@TbCoreComponent |
||||
|
public class EdgeMsgConstructor { |
||||
|
|
||||
|
public EdgeConfiguration constructEdgeConfiguration(Edge edge) { |
||||
|
EdgeConfiguration.Builder builder = EdgeConfiguration.newBuilder() |
||||
|
.setEdgeIdMSB(edge.getId().getId().getMostSignificantBits()) |
||||
|
.setEdgeIdLSB(edge.getId().getId().getLeastSignificantBits()) |
||||
|
.setTenantIdMSB(edge.getTenantId().getId().getMostSignificantBits()) |
||||
|
.setTenantIdLSB(edge.getTenantId().getId().getLeastSignificantBits()) |
||||
|
.setName(edge.getName()) |
||||
|
.setType(edge.getType()) |
||||
|
.setRoutingKey(edge.getRoutingKey()) |
||||
|
.setSecret(edge.getSecret()) |
||||
|
.setAdditionalInfo(JacksonUtil.toString(edge.getAdditionalInfo())) |
||||
|
.setCloudType("CE"); |
||||
|
if (edge.getCustomerId() != null) { |
||||
|
builder.setCustomerIdMSB(edge.getCustomerId().getId().getMostSignificantBits()) |
||||
|
.setCustomerIdLSB(edge.getCustomerId().getId().getLeastSignificantBits()); |
||||
|
} |
||||
|
return builder.build(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,47 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.edge.rpc.fetch; |
||||
|
|
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.EdgeUtils; |
||||
|
import org.thingsboard.server.common.data.edge.Edge; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEvent; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEventType; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.dao.device.DeviceService; |
||||
|
|
||||
|
@AllArgsConstructor |
||||
|
@Slf4j |
||||
|
public class DevicesEdgeEventFetcher extends BasePageableEdgeEventFetcher<Device> { |
||||
|
|
||||
|
private final DeviceService deviceService; |
||||
|
|
||||
|
@Override |
||||
|
PageData<Device> fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { |
||||
|
return deviceService.findDevicesByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, Device device) { |
||||
|
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, |
||||
|
EdgeEventActionType.ADDED, device.getId(), null); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,47 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.edge.rpc.fetch; |
||||
|
|
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.common.data.EdgeUtils; |
||||
|
import org.thingsboard.server.common.data.EntityView; |
||||
|
import org.thingsboard.server.common.data.edge.Edge; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEvent; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEventType; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.dao.entityview.EntityViewService; |
||||
|
|
||||
|
@AllArgsConstructor |
||||
|
@Slf4j |
||||
|
public class EntityViewsEdgeEventFetcher extends BasePageableEdgeEventFetcher<EntityView> { |
||||
|
|
||||
|
private final EntityViewService entityViewService; |
||||
|
|
||||
|
@Override |
||||
|
PageData<EntityView> fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { |
||||
|
return entityViewService.findEntityViewsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, EntityView entityView) { |
||||
|
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ENTITY_VIEW, |
||||
|
EdgeEventActionType.ADDED, entityView.getId(), null); |
||||
|
} |
||||
|
} |
||||
@ -1,235 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.edge.rpc.processor; |
|
||||
|
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.stereotype.Component; |
|
||||
import org.thingsboard.server.common.data.Device; |
|
||||
import org.thingsboard.server.common.data.EdgeUtils; |
|
||||
import org.thingsboard.server.common.data.edge.Edge; |
|
||||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|
||||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|
||||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|
||||
import org.thingsboard.server.common.data.id.CustomerId; |
|
||||
import org.thingsboard.server.common.data.id.DeviceId; |
|
||||
import org.thingsboard.server.common.data.id.EdgeId; |
|
||||
import org.thingsboard.server.common.data.id.EntityId; |
|
||||
import org.thingsboard.server.common.data.id.EntityIdFactory; |
|
||||
import org.thingsboard.server.common.data.id.RuleChainId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.data.page.PageData; |
|
||||
import org.thingsboard.server.common.data.page.PageLink; |
|
||||
import org.thingsboard.server.common.data.rule.RuleChain; |
|
||||
import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo; |
|
||||
import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; |
|
||||
import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; |
|
||||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|
||||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|
||||
import org.thingsboard.server.gen.transport.TransportProtos; |
|
||||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|
||||
|
|
||||
import java.util.ArrayList; |
|
||||
import java.util.List; |
|
||||
import java.util.UUID; |
|
||||
|
|
||||
@Component |
|
||||
@Slf4j |
|
||||
@TbCoreComponent |
|
||||
public class EntityEdgeProcessor extends BaseEdgeProcessor { |
|
||||
|
|
||||
public DownlinkMsg processEntityMergeRequestMessageToEdge(Edge edge, EdgeEvent edgeEvent) { |
|
||||
DownlinkMsg downlinkMsg = null; |
|
||||
if (EdgeEventType.DEVICE.equals(edgeEvent.getType())) { |
|
||||
DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); |
|
||||
Device device = deviceService.findDeviceById(edge.getTenantId(), deviceId); |
|
||||
CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(device, edge); |
|
||||
String conflictName = null; |
|
||||
if(edgeEvent.getBody() != null) { |
|
||||
conflictName = edgeEvent.getBody().get("conflictName").asText(); |
|
||||
} |
|
||||
DeviceUpdateMsg deviceUpdateMsg = deviceMsgConstructor |
|
||||
.constructDeviceUpdatedMsg(UpdateMsgType.ENTITY_MERGE_RPC_MESSAGE, device, customerId, conflictName); |
|
||||
downlinkMsg = DownlinkMsg.newBuilder() |
|
||||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|
||||
.addDeviceUpdateMsg(deviceUpdateMsg) |
|
||||
.build(); |
|
||||
} |
|
||||
return downlinkMsg; |
|
||||
} |
|
||||
|
|
||||
public DownlinkMsg processCredentialsRequestMessageToEdge(EdgeEvent edgeEvent) { |
|
||||
DownlinkMsg downlinkMsg = null; |
|
||||
if (EdgeEventType.DEVICE.equals(edgeEvent.getType())) { |
|
||||
DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); |
|
||||
DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = DeviceCredentialsRequestMsg.newBuilder() |
|
||||
.setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) |
|
||||
.setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) |
|
||||
.build(); |
|
||||
DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() |
|
||||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|
||||
.addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsg); |
|
||||
downlinkMsg = builder.build(); |
|
||||
} |
|
||||
return downlinkMsg; |
|
||||
} |
|
||||
|
|
||||
public ListenableFuture<Void> processEntityNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { |
|
||||
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); |
|
||||
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); |
|
||||
EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, |
|
||||
new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); |
|
||||
EdgeId edgeId = safeGetEdgeId(edgeNotificationMsg); |
|
||||
switch (actionType) { |
|
||||
case ADDED: // used only for USER entity
|
|
||||
case UPDATED: |
|
||||
case CREDENTIALS_UPDATED: |
|
||||
return pushNotificationToAllRelatedEdges(tenantId, entityId, type, actionType); |
|
||||
case ASSIGNED_TO_CUSTOMER: |
|
||||
case UNASSIGNED_FROM_CUSTOMER: |
|
||||
return pushNotificationToAllRelatedCustomerEdges(tenantId, edgeNotificationMsg, entityId, actionType, type); |
|
||||
case DELETED: |
|
||||
if (edgeId != null) { |
|
||||
return saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); |
|
||||
} else { |
|
||||
return pushNotificationToAllRelatedEdges(tenantId, entityId, type, actionType); |
|
||||
} |
|
||||
case ASSIGNED_TO_EDGE: |
|
||||
case UNASSIGNED_FROM_EDGE: |
|
||||
ListenableFuture<Void> future = saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); |
|
||||
return Futures.transformAsync(future, unused -> { |
|
||||
if (type.equals(EdgeEventType.RULE_CHAIN)) { |
|
||||
return updateDependentRuleChains(tenantId, new RuleChainId(entityId.getId()), edgeId); |
|
||||
} else { |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
}, dbCallbackExecutorService); |
|
||||
default: |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private ListenableFuture<Void> pushNotificationToAllRelatedCustomerEdges(TenantId tenantId, |
|
||||
TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, |
|
||||
EntityId entityId, |
|
||||
EdgeEventActionType actionType, |
|
||||
EdgeEventType type) { |
|
||||
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
||||
PageData<EdgeId> pageData; |
|
||||
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
||||
do { |
|
||||
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); |
|
||||
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
||||
for (EdgeId relatedEdgeId : pageData.getData()) { |
|
||||
try { |
|
||||
CustomerId customerId = mapper.readValue(edgeNotificationMsg.getBody(), CustomerId.class); |
|
||||
ListenableFuture<Edge> future = edgeService.findEdgeByIdAsync(tenantId, relatedEdgeId); |
|
||||
futures.add(Futures.transformAsync(future, edge -> { |
|
||||
if (edge != null && edge.getCustomerId() != null && |
|
||||
!edge.getCustomerId().isNullUid() && edge.getCustomerId().equals(customerId)) { |
|
||||
return saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null); |
|
||||
} else { |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
}, dbCallbackExecutorService)); |
|
||||
} catch (Exception e) { |
|
||||
log.error("Can't parse customer id from entity body [{}]", edgeNotificationMsg, e); |
|
||||
return Futures.immediateFailedFuture(e); |
|
||||
} |
|
||||
} |
|
||||
if (pageData.hasNext()) { |
|
||||
pageLink = pageLink.nextPageLink(); |
|
||||
} |
|
||||
} |
|
||||
} while (pageData != null && pageData.hasNext()); |
|
||||
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
|
||||
} |
|
||||
|
|
||||
private EdgeId safeGetEdgeId(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { |
|
||||
if (edgeNotificationMsg.getEdgeIdMSB() != 0 && edgeNotificationMsg.getEdgeIdLSB() != 0) { |
|
||||
return new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB())); |
|
||||
} else { |
|
||||
return null; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private ListenableFuture<Void> pushNotificationToAllRelatedEdges(TenantId tenantId, EntityId entityId, EdgeEventType type, EdgeEventActionType actionType) { |
|
||||
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
||||
PageData<EdgeId> pageData; |
|
||||
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
||||
do { |
|
||||
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); |
|
||||
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
||||
for (EdgeId relatedEdgeId : pageData.getData()) { |
|
||||
futures.add(saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null)); |
|
||||
} |
|
||||
if (pageData.hasNext()) { |
|
||||
pageLink = pageLink.nextPageLink(); |
|
||||
} |
|
||||
} |
|
||||
} while (pageData != null && pageData.hasNext()); |
|
||||
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
|
||||
} |
|
||||
|
|
||||
private ListenableFuture<Void> updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) { |
|
||||
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
||||
PageData<RuleChain> pageData; |
|
||||
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
||||
do { |
|
||||
pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink); |
|
||||
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
||||
for (RuleChain ruleChain : pageData.getData()) { |
|
||||
if (!ruleChain.getId().equals(processingRuleChainId)) { |
|
||||
List<RuleChainConnectionInfo> connectionInfos = |
|
||||
ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections(); |
|
||||
if (connectionInfos != null && !connectionInfos.isEmpty()) { |
|
||||
for (RuleChainConnectionInfo connectionInfo : connectionInfos) { |
|
||||
if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) { |
|
||||
futures.add(saveEdgeEvent(tenantId, |
|
||||
edgeId, |
|
||||
EdgeEventType.RULE_CHAIN_METADATA, |
|
||||
EdgeEventActionType.UPDATED, |
|
||||
ruleChain.getId(), |
|
||||
null)); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
if (pageData.hasNext()) { |
|
||||
pageLink = pageLink.nextPageLink(); |
|
||||
} |
|
||||
} |
|
||||
} while (pageData != null && pageData.hasNext()); |
|
||||
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
|
||||
} |
|
||||
|
|
||||
public ListenableFuture<Void> processEntityNotificationForAllEdges(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { |
|
||||
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); |
|
||||
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); |
|
||||
EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); |
|
||||
switch (actionType) { |
|
||||
case ADDED: |
|
||||
case UPDATED: |
|
||||
case DELETED: |
|
||||
return processActionForAllEdges(tenantId, type, actionType, entityId); |
|
||||
default: |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@ -1,173 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.script; |
|
||||
|
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
||||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|
||||
import org.thingsboard.server.common.data.id.CustomerId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import java.util.Map; |
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.ConcurrentHashMap; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.ScheduledExecutorService; |
|
||||
import java.util.concurrent.TimeoutException; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
|
|
||||
/** |
|
||||
* Created by ashvayka on 26.09.18. |
|
||||
*/ |
|
||||
@Slf4j |
|
||||
public abstract class AbstractJsInvokeService implements JsInvokeService { |
|
||||
|
|
||||
private final TbApiUsageStateService apiUsageStateService; |
|
||||
private final TbApiUsageClient apiUsageClient; |
|
||||
protected ScheduledExecutorService timeoutExecutorService; |
|
||||
protected Map<UUID, String> scriptIdToNameMap = new ConcurrentHashMap<>(); |
|
||||
protected Map<UUID, DisableListInfo> disabledFunctions = new ConcurrentHashMap<>(); |
|
||||
|
|
||||
protected AbstractJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient) { |
|
||||
this.apiUsageStateService = apiUsageStateService; |
|
||||
this.apiUsageClient = apiUsageClient; |
|
||||
} |
|
||||
|
|
||||
public void init(long maxRequestsTimeout) { |
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
timeoutExecutorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("nashorn-js-timeout")); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public void stop() { |
|
||||
if (timeoutExecutorService != null) { |
|
||||
timeoutExecutorService.shutdownNow(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<UUID> eval(TenantId tenantId, JsScriptType scriptType, String scriptBody, String... argNames) { |
|
||||
if (apiUsageStateService.getApiUsageState(tenantId).isJsExecEnabled()) { |
|
||||
UUID scriptId = UUID.randomUUID(); |
|
||||
String functionName = "invokeInternal_" + scriptId.toString().replace('-', '_'); |
|
||||
String jsScript = generateJsScript(scriptType, functionName, scriptBody, argNames); |
|
||||
return doEval(scriptId, functionName, jsScript); |
|
||||
} else { |
|
||||
return Futures.immediateFailedFuture(new RuntimeException("JS Execution is disabled due to API limits!")); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Object> invokeFunction(TenantId tenantId, CustomerId customerId, UUID scriptId, Object... args) { |
|
||||
if (apiUsageStateService.getApiUsageState(tenantId).isJsExecEnabled()) { |
|
||||
String functionName = scriptIdToNameMap.get(scriptId); |
|
||||
if (functionName == null) { |
|
||||
return Futures.immediateFailedFuture(new RuntimeException("No compiled script found for scriptId: [" + scriptId + "]!")); |
|
||||
} |
|
||||
if (!isDisabled(scriptId)) { |
|
||||
apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.JS_EXEC_COUNT, 1); |
|
||||
return doInvokeFunction(scriptId, functionName, args); |
|
||||
} else { |
|
||||
String message = "Script invocation is blocked due to maximum error count " |
|
||||
+ getMaxErrors() + ", scriptId " + scriptId + "!"; |
|
||||
log.warn(message); |
|
||||
return Futures.immediateFailedFuture(new RuntimeException(message)); |
|
||||
} |
|
||||
} else { |
|
||||
return Futures.immediateFailedFuture(new RuntimeException("JS Execution is disabled due to API limits!")); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public ListenableFuture<Void> release(UUID scriptId) { |
|
||||
String functionName = scriptIdToNameMap.get(scriptId); |
|
||||
if (functionName != null) { |
|
||||
try { |
|
||||
scriptIdToNameMap.remove(scriptId); |
|
||||
disabledFunctions.remove(scriptId); |
|
||||
doRelease(scriptId, functionName); |
|
||||
} catch (Exception e) { |
|
||||
return Futures.immediateFailedFuture(e); |
|
||||
} |
|
||||
} |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
|
|
||||
protected abstract ListenableFuture<UUID> doEval(UUID scriptId, String functionName, String scriptBody); |
|
||||
|
|
||||
protected abstract ListenableFuture<Object> doInvokeFunction(UUID scriptId, String functionName, Object[] args); |
|
||||
|
|
||||
protected abstract void doRelease(UUID scriptId, String functionName) throws Exception; |
|
||||
|
|
||||
protected abstract int getMaxErrors(); |
|
||||
|
|
||||
protected abstract long getMaxBlacklistDuration(); |
|
||||
|
|
||||
protected void onScriptExecutionError(UUID scriptId, Throwable t, String scriptBody) { |
|
||||
DisableListInfo disableListInfo = disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()); |
|
||||
log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}, scriptBody {}", |
|
||||
disableListInfo.get(), scriptId, t, t.getCause(), scriptBody); |
|
||||
disableListInfo.incrementAndGet(); |
|
||||
} |
|
||||
|
|
||||
private String generateJsScript(JsScriptType scriptType, String functionName, String scriptBody, String... argNames) { |
|
||||
if (scriptType == JsScriptType.RULE_NODE_SCRIPT) { |
|
||||
return RuleNodeScriptFactory.generateRuleNodeScript(functionName, scriptBody, argNames); |
|
||||
} |
|
||||
throw new RuntimeException("No script factory implemented for scriptType: " + scriptType); |
|
||||
} |
|
||||
|
|
||||
private boolean isDisabled(UUID scriptId) { |
|
||||
DisableListInfo errorCount = disabledFunctions.get(scriptId); |
|
||||
if (errorCount != null) { |
|
||||
if (errorCount.getExpirationTime() <= System.currentTimeMillis()) { |
|
||||
disabledFunctions.remove(scriptId); |
|
||||
return false; |
|
||||
} else { |
|
||||
return errorCount.get() >= getMaxErrors(); |
|
||||
} |
|
||||
} else { |
|
||||
return false; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private class DisableListInfo { |
|
||||
private final AtomicInteger counter; |
|
||||
private long expirationTime; |
|
||||
|
|
||||
private DisableListInfo() { |
|
||||
this.counter = new AtomicInteger(0); |
|
||||
} |
|
||||
|
|
||||
public int get() { |
|
||||
return counter.get(); |
|
||||
} |
|
||||
|
|
||||
public int incrementAndGet() { |
|
||||
int result = counter.incrementAndGet(); |
|
||||
expirationTime = System.currentTimeMillis() + getMaxBlacklistDuration(); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
public long getExpirationTime() { |
|
||||
return expirationTime; |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,185 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.script; |
|
||||
|
|
||||
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 delight.nashornsandbox.NashornSandbox; |
|
||||
import delight.nashornsandbox.NashornSandboxes; |
|
||||
import lombok.Getter; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.scheduling.annotation.Scheduled; |
|
||||
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import javax.annotation.PostConstruct; |
|
||||
import javax.annotation.PreDestroy; |
|
||||
import javax.script.Invocable; |
|
||||
import javax.script.ScriptEngine; |
|
||||
import javax.script.ScriptEngineManager; |
|
||||
import javax.script.ScriptException; |
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.ExecutionException; |
|
||||
import java.util.concurrent.ExecutorService; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
import java.util.concurrent.locks.ReentrantLock; |
|
||||
|
|
||||
@Slf4j |
|
||||
public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeService { |
|
||||
|
|
||||
private NashornSandbox sandbox; |
|
||||
private ScriptEngine engine; |
|
||||
private ExecutorService monitorExecutorService; |
|
||||
|
|
||||
private final AtomicInteger jsPushedMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsInvokeMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsEvalMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsFailedMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger jsTimeoutMsgs = new AtomicInteger(0); |
|
||||
private final FutureCallback<UUID> evalCallback = new JsStatCallback<>(jsEvalMsgs, jsTimeoutMsgs, jsFailedMsgs); |
|
||||
private final FutureCallback<Object> invokeCallback = new JsStatCallback<>(jsInvokeMsgs, jsTimeoutMsgs, jsFailedMsgs); |
|
||||
|
|
||||
private final ReentrantLock evalLock = new ReentrantLock(); |
|
||||
|
|
||||
@Getter |
|
||||
private final JsExecutorService jsExecutor; |
|
||||
|
|
||||
@Value("${js.local.max_requests_timeout:0}") |
|
||||
private long maxRequestsTimeout; |
|
||||
|
|
||||
@Value("${js.local.stats.enabled:false}") |
|
||||
private boolean statsEnabled; |
|
||||
|
|
||||
public AbstractNashornJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient, JsExecutorService jsExecutor) { |
|
||||
super(apiUsageStateService, apiUsageClient); |
|
||||
this.jsExecutor = jsExecutor; |
|
||||
} |
|
||||
|
|
||||
@Scheduled(fixedDelayString = "${js.local.stats.print_interval_ms:10000}") |
|
||||
public void printStats() { |
|
||||
if (statsEnabled) { |
|
||||
int pushedMsgs = jsPushedMsgs.getAndSet(0); |
|
||||
int invokeMsgs = jsInvokeMsgs.getAndSet(0); |
|
||||
int evalMsgs = jsEvalMsgs.getAndSet(0); |
|
||||
int failed = jsFailedMsgs.getAndSet(0); |
|
||||
int timedOut = jsTimeoutMsgs.getAndSet(0); |
|
||||
if (pushedMsgs > 0 || invokeMsgs > 0 || evalMsgs > 0 || failed > 0 || timedOut > 0) { |
|
||||
log.info("Nashorn JS Invoke Stats: pushed [{}] received [{}] invoke [{}] eval [{}] failed [{}] timedOut [{}]", |
|
||||
pushedMsgs, invokeMsgs + evalMsgs, invokeMsgs, evalMsgs, failed, timedOut); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PostConstruct |
|
||||
public void init() { |
|
||||
super.init(maxRequestsTimeout); |
|
||||
if (useJsSandbox()) { |
|
||||
sandbox = NashornSandboxes.create(); |
|
||||
monitorExecutorService = ThingsBoardExecutors.newWorkStealingPool(getMonitorThreadPoolSize(), "nashorn-js-monitor"); |
|
||||
sandbox.setExecutor(monitorExecutorService); |
|
||||
sandbox.setMaxCPUTime(getMaxCpuTime()); |
|
||||
sandbox.allowNoBraces(false); |
|
||||
sandbox.allowLoadFunctions(true); |
|
||||
sandbox.setMaxPreparedStatements(30); |
|
||||
} else { |
|
||||
ScriptEngineManager factory = new ScriptEngineManager(); |
|
||||
engine = factory.getEngineByName("nashorn"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@PreDestroy |
|
||||
public void stop() { |
|
||||
super.stop(); |
|
||||
if (monitorExecutorService != null) { |
|
||||
monitorExecutorService.shutdownNow(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
protected abstract boolean useJsSandbox(); |
|
||||
|
|
||||
protected abstract int getMonitorThreadPoolSize(); |
|
||||
|
|
||||
protected abstract long getMaxCpuTime(); |
|
||||
|
|
||||
@Override |
|
||||
protected ListenableFuture<UUID> doEval(UUID scriptId, String functionName, String jsScript) { |
|
||||
jsPushedMsgs.incrementAndGet(); |
|
||||
ListenableFuture<UUID> result = jsExecutor.executeAsync(() -> { |
|
||||
try { |
|
||||
evalLock.lock(); |
|
||||
try { |
|
||||
if (useJsSandbox()) { |
|
||||
sandbox.eval(jsScript); |
|
||||
} else { |
|
||||
engine.eval(jsScript); |
|
||||
} |
|
||||
} finally { |
|
||||
evalLock.unlock(); |
|
||||
} |
|
||||
scriptIdToNameMap.put(scriptId, functionName); |
|
||||
return scriptId; |
|
||||
} catch (Exception e) { |
|
||||
log.debug("Failed to compile JS script: {}", e.getMessage(), e); |
|
||||
throw new ExecutionException(e); |
|
||||
} |
|
||||
}); |
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
result = Futures.withTimeout(result, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
Futures.addCallback(result, evalCallback, MoreExecutors.directExecutor()); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected ListenableFuture<Object> doInvokeFunction(UUID scriptId, String functionName, Object[] args) { |
|
||||
jsPushedMsgs.incrementAndGet(); |
|
||||
ListenableFuture<Object> result = jsExecutor.executeAsync(() -> { |
|
||||
try { |
|
||||
if (useJsSandbox()) { |
|
||||
return sandbox.getSandboxedInvocable().invokeFunction(functionName, args); |
|
||||
} else { |
|
||||
return ((Invocable) engine).invokeFunction(functionName, args); |
|
||||
} |
|
||||
} catch (ScriptException e) { |
|
||||
throw new ExecutionException(e); |
|
||||
} catch (Exception e) { |
|
||||
onScriptExecutionError(scriptId, e, functionName); |
|
||||
throw new ExecutionException(e); |
|
||||
} |
|
||||
}); |
|
||||
|
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
result = Futures.withTimeout(result, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
Futures.addCallback(result, invokeCallback, MoreExecutors.directExecutor()); |
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
protected void doRelease(UUID scriptId, String functionName) throws ScriptException { |
|
||||
if (useJsSandbox()) { |
|
||||
sandbox.eval(functionName + " = undefined;"); |
|
||||
} else { |
|
||||
engine.eval(functionName + " = undefined;"); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -1,75 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.script; |
|
||||
|
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; |
|
||||
import org.springframework.stereotype.Service; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
|
|
||||
@Slf4j |
|
||||
@ConditionalOnProperty(prefix = "js", value = "evaluator", havingValue = "local", matchIfMissing = true) |
|
||||
@Service |
|
||||
public class NashornJsInvokeService extends AbstractNashornJsInvokeService { |
|
||||
|
|
||||
@Value("${js.local.use_js_sandbox}") |
|
||||
private boolean useJsSandbox; |
|
||||
|
|
||||
@Value("${js.local.monitor_thread_pool_size}") |
|
||||
private int monitorThreadPoolSize; |
|
||||
|
|
||||
@Value("${js.local.max_cpu_time}") |
|
||||
private long maxCpuTime; |
|
||||
|
|
||||
@Value("${js.local.max_errors}") |
|
||||
private int maxErrors; |
|
||||
|
|
||||
@Value("${js.local.max_black_list_duration_sec:60}") |
|
||||
private int maxBlackListDurationSec; |
|
||||
|
|
||||
public NashornJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient, JsExecutorService jsExecutor) { |
|
||||
super(apiUsageStateService, apiUsageClient, jsExecutor); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected boolean useJsSandbox() { |
|
||||
return useJsSandbox; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected int getMonitorThreadPoolSize() { |
|
||||
return monitorThreadPoolSize; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected long getMaxCpuTime() { |
|
||||
return maxCpuTime; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected int getMaxErrors() { |
|
||||
return maxErrors; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected long getMaxBlacklistDuration() { |
|
||||
return TimeUnit.SECONDS.toMillis(maxBlackListDurationSec); |
|
||||
} |
|
||||
} |
|
||||
@ -1,261 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 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.script; |
|
||||
|
|
||||
import com.google.common.util.concurrent.FutureCallback; |
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import lombok.Getter; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Autowired; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; |
|
||||
import org.springframework.scheduling.annotation.Scheduled; |
|
||||
import org.springframework.stereotype.Service; |
|
||||
import org.springframework.util.StopWatch; |
|
||||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
||||
import org.thingsboard.server.gen.js.JsInvokeProtos; |
|
||||
import org.thingsboard.server.queue.TbQueueRequestTemplate; |
|
||||
import org.thingsboard.server.queue.common.TbProtoJsQueueMsg; |
|
||||
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
|
||||
import org.thingsboard.server.queue.usagestats.TbApiUsageClient; |
|
||||
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; |
|
||||
|
|
||||
import javax.annotation.Nullable; |
|
||||
import javax.annotation.PostConstruct; |
|
||||
import javax.annotation.PreDestroy; |
|
||||
import java.util.Map; |
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.ConcurrentHashMap; |
|
||||
import java.util.concurrent.ExecutorService; |
|
||||
import java.util.concurrent.Executors; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
import java.util.concurrent.TimeoutException; |
|
||||
import java.util.concurrent.atomic.AtomicInteger; |
|
||||
|
|
||||
@Slf4j |
|
||||
@ConditionalOnExpression("'${js.evaluator:null}'=='remote' && ('${service.type:null}'=='monolith' || '${service.type:null}'=='tb-core' || '${service.type:null}'=='tb-rule-engine')") |
|
||||
@Service |
|
||||
public class RemoteJsInvokeService extends AbstractJsInvokeService { |
|
||||
|
|
||||
@Value("${queue.js.max_eval_requests_timeout}") |
|
||||
private long maxEvalRequestsTimeout; |
|
||||
|
|
||||
@Value("${queue.js.max_requests_timeout}") |
|
||||
private long maxRequestsTimeout; |
|
||||
|
|
||||
@Value("${queue.js.max_exec_requests_timeout:2000}") |
|
||||
private long maxExecRequestsTimeout; |
|
||||
|
|
||||
@Getter |
|
||||
@Value("${js.remote.max_errors}") |
|
||||
private int maxErrors; |
|
||||
|
|
||||
@Value("${js.remote.max_black_list_duration_sec:60}") |
|
||||
private int maxBlackListDurationSec; |
|
||||
|
|
||||
@Value("${js.remote.stats.enabled:false}") |
|
||||
private boolean statsEnabled; |
|
||||
|
|
||||
private final AtomicInteger queuePushedMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger queueInvokeMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger queueEvalMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger queueFailedMsgs = new AtomicInteger(0); |
|
||||
private final AtomicInteger queueTimeoutMsgs = new AtomicInteger(0); |
|
||||
private final ExecutorService callbackExecutor = Executors.newFixedThreadPool( |
|
||||
Runtime.getRuntime().availableProcessors(), ThingsBoardThreadFactory.forName("js-executor-remote-callback")); |
|
||||
|
|
||||
public RemoteJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient) { |
|
||||
super(apiUsageStateService, apiUsageClient); |
|
||||
} |
|
||||
|
|
||||
@Scheduled(fixedDelayString = "${js.remote.stats.print_interval_ms}") |
|
||||
public void printStats() { |
|
||||
if (statsEnabled) { |
|
||||
int pushedMsgs = queuePushedMsgs.getAndSet(0); |
|
||||
int invokeMsgs = queueInvokeMsgs.getAndSet(0); |
|
||||
int evalMsgs = queueEvalMsgs.getAndSet(0); |
|
||||
int failed = queueFailedMsgs.getAndSet(0); |
|
||||
int timedOut = queueTimeoutMsgs.getAndSet(0); |
|
||||
if (pushedMsgs > 0 || invokeMsgs > 0 || evalMsgs > 0 || failed > 0 || timedOut > 0) { |
|
||||
log.info("Queue JS Invoke Stats: pushed [{}] received [{}] invoke [{}] eval [{}] failed [{}] timedOut [{}]", |
|
||||
pushedMsgs, invokeMsgs + evalMsgs, invokeMsgs, evalMsgs, failed, timedOut); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Autowired |
|
||||
private TbQueueRequestTemplate<TbProtoJsQueueMsg<JsInvokeProtos.RemoteJsRequest>, TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> requestTemplate; |
|
||||
|
|
||||
private Map<UUID, String> scriptIdToBodysMap = new ConcurrentHashMap<>(); |
|
||||
|
|
||||
@PostConstruct |
|
||||
public void init() { |
|
||||
super.init(maxRequestsTimeout); |
|
||||
requestTemplate.init(); |
|
||||
} |
|
||||
|
|
||||
@PreDestroy |
|
||||
public void destroy() { |
|
||||
super.stop(); |
|
||||
if (requestTemplate != null) { |
|
||||
requestTemplate.stop(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected ListenableFuture<UUID> doEval(UUID scriptId, String functionName, String scriptBody) { |
|
||||
JsInvokeProtos.JsCompileRequest jsRequest = JsInvokeProtos.JsCompileRequest.newBuilder() |
|
||||
.setScriptIdMSB(scriptId.getMostSignificantBits()) |
|
||||
.setScriptIdLSB(scriptId.getLeastSignificantBits()) |
|
||||
.setFunctionName(functionName) |
|
||||
.setScriptBody(scriptBody).build(); |
|
||||
|
|
||||
JsInvokeProtos.RemoteJsRequest jsRequestWrapper = JsInvokeProtos.RemoteJsRequest.newBuilder() |
|
||||
.setCompileRequest(jsRequest) |
|
||||
.build(); |
|
||||
|
|
||||
log.trace("Post compile request for scriptId [{}]", scriptId); |
|
||||
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); |
|
||||
if (maxEvalRequestsTimeout > 0) { |
|
||||
future = Futures.withTimeout(future, maxEvalRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
queuePushedMsgs.incrementAndGet(); |
|
||||
Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() { |
|
||||
@Override |
|
||||
public void onSuccess(@Nullable TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse> result) { |
|
||||
queueEvalMsgs.incrementAndGet(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(Throwable t) { |
|
||||
if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { |
|
||||
queueTimeoutMsgs.incrementAndGet(); |
|
||||
} |
|
||||
queueFailedMsgs.incrementAndGet(); |
|
||||
} |
|
||||
}, callbackExecutor); |
|
||||
return Futures.transform(future, response -> { |
|
||||
JsInvokeProtos.JsCompileResponse compilationResult = response.getValue().getCompileResponse(); |
|
||||
UUID compiledScriptId = new UUID(compilationResult.getScriptIdMSB(), compilationResult.getScriptIdLSB()); |
|
||||
if (compilationResult.getSuccess()) { |
|
||||
scriptIdToNameMap.put(scriptId, functionName); |
|
||||
scriptIdToBodysMap.put(scriptId, scriptBody); |
|
||||
return compiledScriptId; |
|
||||
} else { |
|
||||
log.debug("[{}] Failed to compile script due to [{}]: {}", compiledScriptId, compilationResult.getErrorCode().name(), compilationResult.getErrorDetails()); |
|
||||
throw new RuntimeException(compilationResult.getErrorDetails()); |
|
||||
} |
|
||||
}, callbackExecutor); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected ListenableFuture<Object> doInvokeFunction(UUID scriptId, String functionName, Object[] args) { |
|
||||
log.trace("doInvokeFunction js-request for uuid {} with timeout {}ms", scriptId, maxRequestsTimeout); |
|
||||
final String scriptBody = scriptIdToBodysMap.get(scriptId); |
|
||||
if (scriptBody == null) { |
|
||||
return Futures.immediateFailedFuture(new RuntimeException("No script body found for scriptId: [" + scriptId + "]!")); |
|
||||
} |
|
||||
JsInvokeProtos.JsInvokeRequest.Builder jsRequestBuilder = JsInvokeProtos.JsInvokeRequest.newBuilder() |
|
||||
.setScriptIdMSB(scriptId.getMostSignificantBits()) |
|
||||
.setScriptIdLSB(scriptId.getLeastSignificantBits()) |
|
||||
.setFunctionName(functionName) |
|
||||
.setTimeout((int) maxExecRequestsTimeout) |
|
||||
.setScriptBody(scriptBody); |
|
||||
|
|
||||
for (Object arg : args) { |
|
||||
jsRequestBuilder.addArgs(arg.toString()); |
|
||||
} |
|
||||
|
|
||||
JsInvokeProtos.RemoteJsRequest jsRequestWrapper = JsInvokeProtos.RemoteJsRequest.newBuilder() |
|
||||
.setInvokeRequest(jsRequestBuilder.build()) |
|
||||
.build(); |
|
||||
|
|
||||
StopWatch stopWatch = new StopWatch(); |
|
||||
stopWatch.start(); |
|
||||
|
|
||||
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); |
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
queuePushedMsgs.incrementAndGet(); |
|
||||
Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() { |
|
||||
@Override |
|
||||
public void onSuccess(@Nullable TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse> result) { |
|
||||
queueInvokeMsgs.incrementAndGet(); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void onFailure(Throwable t) { |
|
||||
if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { |
|
||||
queueTimeoutMsgs.incrementAndGet(); |
|
||||
} |
|
||||
queueFailedMsgs.incrementAndGet(); |
|
||||
} |
|
||||
}, callbackExecutor); |
|
||||
return Futures.transform(future, response -> { |
|
||||
stopWatch.stop(); |
|
||||
log.trace("doInvokeFunction js-response took {}ms for uuid {}", stopWatch.getTotalTimeMillis(), response.getKey()); |
|
||||
JsInvokeProtos.JsInvokeResponse invokeResult = response.getValue().getInvokeResponse(); |
|
||||
if (invokeResult.getSuccess()) { |
|
||||
return invokeResult.getResult(); |
|
||||
} else { |
|
||||
final RuntimeException e = new RuntimeException(invokeResult.getErrorDetails()); |
|
||||
if (JsInvokeProtos.JsInvokeErrorCode.TIMEOUT_ERROR.equals(invokeResult.getErrorCode())) { |
|
||||
onScriptExecutionError(scriptId, e, scriptBody); |
|
||||
queueTimeoutMsgs.incrementAndGet(); |
|
||||
} else if (JsInvokeProtos.JsInvokeErrorCode.COMPILATION_ERROR.equals(invokeResult.getErrorCode())) { |
|
||||
onScriptExecutionError(scriptId, e, scriptBody); |
|
||||
} |
|
||||
queueFailedMsgs.incrementAndGet(); |
|
||||
log.debug("[{}] Failed to invoke function due to [{}]: {}", scriptId, invokeResult.getErrorCode().name(), invokeResult.getErrorDetails()); |
|
||||
throw e; |
|
||||
} |
|
||||
}, callbackExecutor); |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected void doRelease(UUID scriptId, String functionName) throws Exception { |
|
||||
JsInvokeProtos.JsReleaseRequest jsRequest = JsInvokeProtos.JsReleaseRequest.newBuilder() |
|
||||
.setScriptIdMSB(scriptId.getMostSignificantBits()) |
|
||||
.setScriptIdLSB(scriptId.getLeastSignificantBits()) |
|
||||
.setFunctionName(functionName).build(); |
|
||||
|
|
||||
JsInvokeProtos.RemoteJsRequest jsRequestWrapper = JsInvokeProtos.RemoteJsRequest.newBuilder() |
|
||||
.setReleaseRequest(jsRequest) |
|
||||
.build(); |
|
||||
|
|
||||
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); |
|
||||
if (maxRequestsTimeout > 0) { |
|
||||
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
||||
} |
|
||||
JsInvokeProtos.RemoteJsResponse response = future.get().getValue(); |
|
||||
|
|
||||
JsInvokeProtos.JsReleaseResponse compilationResult = response.getReleaseResponse(); |
|
||||
UUID compiledScriptId = new UUID(compilationResult.getScriptIdMSB(), compilationResult.getScriptIdLSB()); |
|
||||
if (compilationResult.getSuccess()) { |
|
||||
scriptIdToBodysMap.remove(scriptId); |
|
||||
} else { |
|
||||
log.debug("[{}] Failed to release script due", compiledScriptId); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
protected long getMaxBlacklistDuration() { |
|
||||
return TimeUnit.SECONDS.toMillis(maxBlackListDurationSec); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -0,0 +1,171 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.type.TypeReference; |
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
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.common.util.JacksonUtil; |
||||
|
import org.thingsboard.script.api.RuleNodeScriptFactory; |
||||
|
import org.thingsboard.script.api.mvel.MvelInvokeService; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
|
||||
|
import javax.script.ScriptException; |
||||
|
import java.util.ArrayList; |
||||
|
import java.util.Collection; |
||||
|
import java.util.Collections; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.HashSet; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.Set; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
|
||||
|
@Slf4j |
||||
|
public class RuleNodeMvelScriptEngine extends RuleNodeScriptEngine<MvelInvokeService, Object> { |
||||
|
|
||||
|
public RuleNodeMvelScriptEngine(TenantId tenantId, MvelInvokeService scriptInvokeService, String script, String... argNames) { |
||||
|
super(tenantId, scriptInvokeService, script, argNames); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<Boolean> executeFilterTransform(Object result) { |
||||
|
if (result instanceof Boolean) { |
||||
|
return Futures.immediateFuture((Boolean) result); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<List<TbMsg>> executeUpdateTransform(TbMsg msg, Object result) { |
||||
|
if (result instanceof Map) { |
||||
|
return Futures.immediateFuture(Collections.singletonList(unbindMsg((Map) result, msg))); |
||||
|
} else if (result instanceof Collection) { |
||||
|
List<TbMsg> res = new ArrayList<>(); |
||||
|
for (Object resObject : (Collection) result) { |
||||
|
if (resObject instanceof Map) { |
||||
|
res.add(unbindMsg((Map) result, msg)); |
||||
|
} else { |
||||
|
return wrongResultType(resObject); |
||||
|
} |
||||
|
} |
||||
|
return Futures.immediateFuture(res); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<TbMsg> executeGenerateTransform(TbMsg prevMsg, Object result) { |
||||
|
if (result instanceof Map) { |
||||
|
return Futures.immediateFuture(unbindMsg((Map) result, prevMsg)); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<String> executeToStringTransform(Object result) { |
||||
|
if (result instanceof String) { |
||||
|
return Futures.immediateFuture((String) result); |
||||
|
} else { |
||||
|
return Futures.immediateFuture(JacksonUtil.toString(result)); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ListenableFuture<Set<String>> executeSwitchTransform(Object result) { |
||||
|
if (result instanceof String) { |
||||
|
return Futures.immediateFuture(Collections.singleton((String) result)); |
||||
|
} else if (result instanceof Collection) { |
||||
|
Set<String> res = new HashSet<>(); |
||||
|
for (Object resObject : (Collection) result) { |
||||
|
if (resObject instanceof String) { |
||||
|
res.add((String) resObject); |
||||
|
} else { |
||||
|
return wrongResultType(resObject); |
||||
|
} |
||||
|
} |
||||
|
return Futures.immediateFuture(res); |
||||
|
} |
||||
|
return wrongResultType(result); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<JsonNode> executeJsonAsync(TbMsg msg) { |
||||
|
return Futures.transform(executeScriptAsync(msg), JacksonUtil::valueToTree, MoreExecutors.directExecutor()); |
||||
|
|
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected Object convertResult(Object result) { |
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected Object[] prepareArgs(TbMsg msg) { |
||||
|
Object[] args = new Object[3]; |
||||
|
if (msg.getData() != null) { |
||||
|
args[0] = JacksonUtil.fromString(msg.getData(), Map.class); |
||||
|
} else { |
||||
|
args[0] = new HashMap<>(); |
||||
|
} |
||||
|
args[1] = new HashMap<>(msg.getMetaData().getData()); |
||||
|
args[2] = msg.getType(); |
||||
|
return args; |
||||
|
} |
||||
|
|
||||
|
private static TbMsg unbindMsg(Map msgData, TbMsg msg) { |
||||
|
String data = null; |
||||
|
Map<String, String> metadata = null; |
||||
|
String messageType = null; |
||||
|
if (msgData.containsKey(RuleNodeScriptFactory.MSG)) { |
||||
|
data = JacksonUtil.toString(msgData.get(RuleNodeScriptFactory.MSG)); |
||||
|
} |
||||
|
if (msgData.containsKey(RuleNodeScriptFactory.METADATA)) { |
||||
|
Object msgMetadataObj = msgData.get(RuleNodeScriptFactory.METADATA); |
||||
|
if (msgMetadataObj instanceof Map) { |
||||
|
metadata = ((Map<?, ?>) msgMetadataObj).entrySet().stream().filter(e -> e.getValue() != null) |
||||
|
.collect(Collectors.toMap(e -> e.getKey().toString(), e -> e.getValue().toString())); |
||||
|
} else { |
||||
|
metadata = JacksonUtil.convertValue(msgMetadataObj, new TypeReference<>() { |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
if (msgData.containsKey(RuleNodeScriptFactory.MSG_TYPE)) { |
||||
|
messageType = msgData.get(RuleNodeScriptFactory.MSG_TYPE).toString(); |
||||
|
} |
||||
|
String newData = data != null ? data : msg.getData(); |
||||
|
TbMsgMetaData newMetadata = metadata != null ? new TbMsgMetaData(metadata) : msg.getMetaData().copy(); |
||||
|
String newMessageType = !StringUtils.isEmpty(messageType) ? messageType : msg.getType(); |
||||
|
return TbMsg.transformMsg(msg, newMessageType, msg.getOriginator(), newMetadata, newData); |
||||
|
} |
||||
|
|
||||
|
private static <T> ListenableFuture<T> wrongResultType(Object result) { |
||||
|
String className = toClassName(result); |
||||
|
log.warn("Wrong result type: {}", className); |
||||
|
return Futures.immediateFailedFuture(new ScriptException("Wrong result type: " + className)); |
||||
|
} |
||||
|
|
||||
|
private static String toClassName(Object result) { |
||||
|
return result != null ? result.getClass().getSimpleName() : "null"; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,133 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.script; |
||||
|
|
||||
|
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.ScriptEngine; |
||||
|
import org.thingsboard.script.api.ScriptInvokeService; |
||||
|
import org.thingsboard.script.api.ScriptType; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
|
||||
|
import javax.script.ScriptException; |
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
|
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class RuleNodeScriptEngine<T extends ScriptInvokeService, R> implements ScriptEngine { |
||||
|
|
||||
|
private final T scriptInvokeService; |
||||
|
|
||||
|
private final UUID scriptId; |
||||
|
private final TenantId tenantId; |
||||
|
|
||||
|
public RuleNodeScriptEngine(TenantId tenantId, T scriptInvokeService, String script, String... argNames) { |
||||
|
this.tenantId = tenantId; |
||||
|
this.scriptInvokeService = scriptInvokeService; |
||||
|
try { |
||||
|
this.scriptId = this.scriptInvokeService.eval(tenantId, ScriptType.RULE_NODE_SCRIPT, script, argNames).get(); |
||||
|
} catch (Exception e) { |
||||
|
Throwable t = e; |
||||
|
if (e instanceof ExecutionException) { |
||||
|
t = e.getCause(); |
||||
|
} |
||||
|
throw new IllegalArgumentException("Can't compile script: " + t.getMessage(), t); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected abstract Object[] prepareArgs(TbMsg msg); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<List<TbMsg>> executeUpdateAsync(TbMsg msg) { |
||||
|
ListenableFuture<R> result = executeScriptAsync(msg); |
||||
|
return Futures.transformAsync(result, |
||||
|
json -> executeUpdateTransform(msg, json), |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<List<TbMsg>> executeUpdateTransform(TbMsg msg, R result); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<TbMsg> executeGenerateAsync(TbMsg prevMsg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(prevMsg), |
||||
|
result -> executeGenerateTransform(prevMsg, result), |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<TbMsg> executeGenerateTransform(TbMsg prevMsg, R result); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<String> executeToStringAsync(TbMsg msg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(msg), this::executeToStringTransform, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Boolean> executeFilterAsync(TbMsg msg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(msg), |
||||
|
this::executeFilterTransform, |
||||
|
MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected abstract ListenableFuture<String> executeToStringTransform(R result); |
||||
|
|
||||
|
protected abstract ListenableFuture<Boolean> executeFilterTransform(R result); |
||||
|
|
||||
|
protected abstract ListenableFuture<Set<String>> executeSwitchTransform(R result); |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<Set<String>> executeSwitchAsync(TbMsg msg) { |
||||
|
return Futures.transformAsync(executeScriptAsync(msg), |
||||
|
this::executeSwitchTransform, |
||||
|
MoreExecutors.directExecutor()); //usually runs in a callbackExecutor
|
||||
|
} |
||||
|
|
||||
|
ListenableFuture<R> executeScriptAsync(TbMsg msg) { |
||||
|
log.trace("execute script async, msg {}", msg); |
||||
|
Object[] inArgs = prepareArgs(msg); |
||||
|
return executeScriptAsync(msg.getCustomerId(), inArgs[0], inArgs[1], inArgs[2]); |
||||
|
} |
||||
|
|
||||
|
ListenableFuture<R> executeScriptAsync(CustomerId customerId, Object... args) { |
||||
|
return Futures.transformAsync(scriptInvokeService.invokeScript(tenantId, customerId, this.scriptId, args), |
||||
|
o -> { |
||||
|
try { |
||||
|
return Futures.immediateFuture(convertResult(o)); |
||||
|
} catch (Exception e) { |
||||
|
if (e.getCause() instanceof ScriptException) { |
||||
|
return Futures.immediateFailedFuture(e.getCause()); |
||||
|
} else if (e.getCause() instanceof RuntimeException) { |
||||
|
return Futures.immediateFailedFuture(new ScriptException(e.getCause().getMessage())); |
||||
|
} else { |
||||
|
return Futures.immediateFailedFuture(new ScriptException(e)); |
||||
|
} |
||||
|
} |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
public void destroy() { |
||||
|
scriptInvokeService.release(this.scriptId); |
||||
|
} |
||||
|
|
||||
|
protected abstract R convertResult(Object result); |
||||
|
} |
||||
@ -0,0 +1,61 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 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.ttl; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; |
||||
|
import org.springframework.scheduling.annotation.Scheduled; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.dao.audit.AuditLogDao; |
||||
|
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository; |
||||
|
import org.thingsboard.server.queue.discovery.PartitionService; |
||||
|
|
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
|
||||
|
import static org.thingsboard.server.dao.model.ModelConstants.AUDIT_LOG_COLUMN_FAMILY_NAME; |
||||
|
|
||||
|
@Service |
||||
|
@ConditionalOnExpression("${sql.ttl.audit_logs.enabled:true} && ${sql.ttl.audit_logs.ttl:0} > 0") |
||||
|
@Slf4j |
||||
|
public class AuditLogsCleanUpService extends AbstractCleanUpService { |
||||
|
|
||||
|
private final AuditLogDao auditLogDao; |
||||
|
private final SqlPartitioningRepository partitioningRepository; |
||||
|
|
||||
|
@Value("${sql.ttl.audit_logs.ttl:0}") |
||||
|
private long ttlInSec; |
||||
|
@Value("${sql.audit_logs.partition_size:168}") |
||||
|
private int partitionSizeInHours; |
||||
|
|
||||
|
public AuditLogsCleanUpService(PartitionService partitionService, AuditLogDao auditLogDao, SqlPartitioningRepository partitioningRepository) { |
||||
|
super(partitionService); |
||||
|
this.auditLogDao = auditLogDao; |
||||
|
this.partitioningRepository = partitioningRepository; |
||||
|
} |
||||
|
|
||||
|
@Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.audit_logs.checking_interval_ms})}", |
||||
|
fixedDelayString = "${sql.ttl.audit_logs.checking_interval_ms}") |
||||
|
public void cleanUp() { |
||||
|
long auditLogsExpTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttlInSec); |
||||
|
if (isSystemTenantPartitionMine()) { |
||||
|
auditLogDao.cleanUpAuditLogs(auditLogsExpTime); |
||||
|
} else { |
||||
|
partitioningRepository.cleanupPartitionsCache(AUDIT_LOG_COLUMN_FAMILY_NAME, auditLogsExpTime, TimeUnit.HOURS.toMillis(partitionSizeInHours)); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue