289 changed files with 6997 additions and 1517 deletions
@ -0,0 +1,136 @@ |
|||||
|
-- |
||||
|
-- Copyright © 2016-2024 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 new attribute_kv table schema |
||||
|
DO |
||||
|
$$ |
||||
|
BEGIN |
||||
|
-- in case of running the upgrade script a second time: |
||||
|
IF EXISTS(SELECT 1 FROM information_schema.columns WHERE table_name = 'attribute_kv' and column_name='entity_type') THEN |
||||
|
DROP VIEW IF EXISTS device_info_view; |
||||
|
DROP VIEW IF EXISTS device_info_active_attribute_view; |
||||
|
ALTER INDEX IF EXISTS idx_attribute_kv_by_key_and_last_update_ts RENAME TO idx_attribute_kv_by_key_and_last_update_ts_old; |
||||
|
IF EXISTS(SELECT 1 FROM pg_constraint WHERE conname = 'attribute_kv_pkey') THEN |
||||
|
ALTER TABLE attribute_kv RENAME CONSTRAINT attribute_kv_pkey TO attribute_kv_pkey_old; |
||||
|
END IF; |
||||
|
ALTER TABLE attribute_kv RENAME TO attribute_kv_old; |
||||
|
CREATE TABLE IF NOT EXISTS attribute_kv |
||||
|
( |
||||
|
entity_id uuid, |
||||
|
attribute_type int, |
||||
|
attribute_key int, |
||||
|
bool_v boolean, |
||||
|
str_v varchar(10000000), |
||||
|
long_v bigint, |
||||
|
dbl_v double precision, |
||||
|
json_v json, |
||||
|
last_update_ts bigint, |
||||
|
CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_id, attribute_type, attribute_key) |
||||
|
); |
||||
|
END IF; |
||||
|
END; |
||||
|
$$; |
||||
|
|
||||
|
-- rename ts_kv_dictionary table to key_dictionary or create table if not exists |
||||
|
DO |
||||
|
$$ |
||||
|
BEGIN |
||||
|
IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'ts_kv_dictionary') THEN |
||||
|
ALTER TABLE ts_kv_dictionary RENAME CONSTRAINT ts_key_id_pkey TO key_dictionary_id_pkey; |
||||
|
ALTER TABLE ts_kv_dictionary RENAME TO key_dictionary; |
||||
|
ELSE CREATE TABLE IF NOT EXISTS key_dictionary( |
||||
|
key varchar(255) NOT NULL, |
||||
|
key_id serial UNIQUE, |
||||
|
CONSTRAINT key_dictionary_id_pkey PRIMARY KEY (key) |
||||
|
); |
||||
|
END IF; |
||||
|
END; |
||||
|
$$; |
||||
|
|
||||
|
-- insert keys into key_dictionary |
||||
|
DO |
||||
|
$$ |
||||
|
BEGIN |
||||
|
IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'attribute_kv_old') THEN |
||||
|
INSERT INTO key_dictionary(key) SELECT DISTINCT attribute_key FROM attribute_kv_old ON CONFLICT DO NOTHING; |
||||
|
END IF; |
||||
|
END; |
||||
|
$$; |
||||
|
|
||||
|
-- migrate attributes from attribute_kv_old to attribute_kv |
||||
|
DO |
||||
|
$$ |
||||
|
DECLARE |
||||
|
row_num_old integer; |
||||
|
row_num integer; |
||||
|
BEGIN |
||||
|
IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'attribute_kv_old') THEN |
||||
|
INSERT INTO attribute_kv(entity_id, attribute_type, attribute_key, bool_v, str_v, long_v, dbl_v, json_v, last_update_ts) |
||||
|
SELECT a.entity_id, CASE |
||||
|
WHEN a.attribute_type = 'CLIENT_SCOPE' THEN 1 |
||||
|
WHEN a.attribute_type = 'SERVER_SCOPE' THEN 2 |
||||
|
WHEN a.attribute_type = 'SHARED_SCOPE' THEN 3 |
||||
|
ELSE 0 |
||||
|
END, |
||||
|
k.key_id, a.bool_v, a.str_v, a.long_v, a.dbl_v, a.json_v, a.last_update_ts |
||||
|
FROM attribute_kv_old a INNER JOIN key_dictionary k ON (a.attribute_key = k.key); |
||||
|
SELECT COUNT(*) INTO row_num_old FROM attribute_kv_old; |
||||
|
SELECT COUNT(*) INTO row_num FROM attribute_kv; |
||||
|
RAISE NOTICE 'Migrated % of % rows', row_num, row_num_old; |
||||
|
|
||||
|
IF row_num != 0 THEN |
||||
|
DROP TABLE IF EXISTS attribute_kv_old; |
||||
|
ELSE |
||||
|
RAISE EXCEPTION 'Table attribute_kv is empty'; |
||||
|
END IF; |
||||
|
|
||||
|
CREATE INDEX IF NOT EXISTS idx_attribute_kv_by_key_and_last_update_ts ON attribute_kv(entity_id, attribute_key, last_update_ts desc); |
||||
|
END IF; |
||||
|
EXCEPTION |
||||
|
WHEN others THEN |
||||
|
ROLLBACK; |
||||
|
RAISE EXCEPTION 'Error during COPY: %', SQLERRM; |
||||
|
END |
||||
|
$$; |
||||
|
|
||||
|
-- OAUTH2 PARAMS ALTER TABLE START |
||||
|
|
||||
|
ALTER TABLE oauth2_params |
||||
|
ADD COLUMN IF NOT EXISTS edge_enabled boolean DEFAULT false; |
||||
|
|
||||
|
-- OAUTH2 PARAMS ALTER TABLE END |
||||
|
|
||||
|
-- QUEUE STATS UPDATE START |
||||
|
|
||||
|
CREATE TABLE IF NOT EXISTS queue_stats ( |
||||
|
id uuid NOT NULL CONSTRAINT queue_stats_pkey PRIMARY KEY, |
||||
|
created_time bigint NOT NULL, |
||||
|
tenant_id uuid NOT NULL, |
||||
|
queue_name varchar(255) NOT NULL, |
||||
|
service_id varchar(255) NOT NULL, |
||||
|
CONSTRAINT queue_stats_name_unq_key UNIQUE (tenant_id, queue_name, service_id) |
||||
|
); |
||||
|
|
||||
|
INSERT INTO queue_stats |
||||
|
SELECT id, created_time, tenant_id, substring(name FROM 1 FOR position('_' IN name) - 1) AS queue_name, |
||||
|
substring(name FROM position('_' IN name) + 1) AS service_id |
||||
|
FROM asset |
||||
|
WHERE type = 'TbServiceQueue' and name LIKE '%\_%'; |
||||
|
|
||||
|
DELETE FROM asset WHERE type='TbServiceQueue'; |
||||
|
DELETE FROM asset_profile WHERE name ='TbServiceQueue'; |
||||
|
|
||||
|
-- QUEUE STATS UPDATE END |
||||
@ -0,0 +1,43 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.notification; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.id.NotificationRuleId; |
||||
|
import org.thingsboard.server.common.data.id.NotificationTargetId; |
||||
|
import org.thingsboard.server.common.data.id.NotificationTemplateId; |
||||
|
import org.thingsboard.server.common.data.notification.rule.NotificationRule; |
||||
|
import org.thingsboard.server.common.data.notification.targets.NotificationTarget; |
||||
|
import org.thingsboard.server.common.data.notification.template.NotificationTemplate; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationRuleUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationTargetUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationTemplateUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
||||
|
|
||||
|
public interface NotificationMsgConstructor { |
||||
|
|
||||
|
NotificationRuleUpdateMsg constructNotificationRuleUpdateMsg(UpdateMsgType msgType, NotificationRule notificationRule); |
||||
|
|
||||
|
NotificationRuleUpdateMsg constructNotificationRuleDeleteMsg(NotificationRuleId notificationRuleId); |
||||
|
|
||||
|
NotificationTargetUpdateMsg constructNotificationTargetUpdateMsg(UpdateMsgType msgType, NotificationTarget notificationTarget); |
||||
|
|
||||
|
NotificationTargetUpdateMsg constructNotificationTargetDeleteMsg(NotificationTargetId notificationTargetId); |
||||
|
|
||||
|
NotificationTemplateUpdateMsg constructNotificationTemplateUpdateMsg(UpdateMsgType msgType, NotificationTemplate notificationTemplate); |
||||
|
|
||||
|
NotificationTemplateUpdateMsg constructNotificationTemplateDeleteMsg(NotificationTemplateId notificationTemplateId); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,75 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.notification; |
||||
|
|
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.id.NotificationRuleId; |
||||
|
import org.thingsboard.server.common.data.id.NotificationTargetId; |
||||
|
import org.thingsboard.server.common.data.id.NotificationTemplateId; |
||||
|
import org.thingsboard.server.common.data.notification.rule.NotificationRule; |
||||
|
import org.thingsboard.server.common.data.notification.targets.NotificationTarget; |
||||
|
import org.thingsboard.server.common.data.notification.template.NotificationTemplate; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationRuleUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationTargetUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationTemplateUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
|
||||
|
@Component |
||||
|
@TbCoreComponent |
||||
|
public class NotificationMsgConstructorImpl implements NotificationMsgConstructor { |
||||
|
|
||||
|
@Override |
||||
|
public NotificationRuleUpdateMsg constructNotificationRuleUpdateMsg(UpdateMsgType msgType, NotificationRule notificationRule) { |
||||
|
return NotificationRuleUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(notificationRule)).build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public NotificationRuleUpdateMsg constructNotificationRuleDeleteMsg(NotificationRuleId notificationRuleId) { |
||||
|
return NotificationRuleUpdateMsg.newBuilder() |
||||
|
.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) |
||||
|
.setIdMSB(notificationRuleId.getId().getMostSignificantBits()) |
||||
|
.setIdLSB(notificationRuleId.getId().getLeastSignificantBits()).build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public NotificationTargetUpdateMsg constructNotificationTargetUpdateMsg(UpdateMsgType msgType, NotificationTarget notificationTarget) { |
||||
|
return NotificationTargetUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(notificationTarget)).build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public NotificationTargetUpdateMsg constructNotificationTargetDeleteMsg(NotificationTargetId notificationTargetId) { |
||||
|
return NotificationTargetUpdateMsg.newBuilder() |
||||
|
.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) |
||||
|
.setIdMSB(notificationTargetId.getId().getMostSignificantBits()) |
||||
|
.setIdLSB(notificationTargetId.getId().getLeastSignificantBits()).build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public NotificationTemplateUpdateMsg constructNotificationTemplateUpdateMsg(UpdateMsgType msgType, NotificationTemplate notificationTemplate) { |
||||
|
return NotificationTemplateUpdateMsg.newBuilder().setMsgType(msgType).setEntity(JacksonUtil.toString(notificationTemplate)).build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public NotificationTemplateUpdateMsg constructNotificationTemplateDeleteMsg(NotificationTemplateId notificationTemplateId) { |
||||
|
return NotificationTemplateUpdateMsg.newBuilder() |
||||
|
.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) |
||||
|
.setIdMSB(notificationTemplateId.getId().getMostSignificantBits()) |
||||
|
.setIdLSB(notificationTemplateId.getId().getLeastSignificantBits()).build(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,48 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.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.notification.rule.NotificationRule; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.dao.notification.NotificationRuleService; |
||||
|
|
||||
|
@AllArgsConstructor |
||||
|
@Slf4j |
||||
|
public class NotificationRuleEdgeEventFetcher extends BasePageableEdgeEventFetcher<NotificationRule>{ |
||||
|
|
||||
|
private NotificationRuleService notificationRuleService; |
||||
|
|
||||
|
@Override |
||||
|
PageData<NotificationRule> fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) { |
||||
|
return notificationRuleService.findNotificationRulesByTenantId(tenantId, pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, NotificationRule notificationRule) { |
||||
|
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.NOTIFICATION_RULE, |
||||
|
EdgeEventActionType.ADDED, notificationRule.getId(), null); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,48 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.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.notification.targets.NotificationTarget; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.dao.notification.NotificationTargetService; |
||||
|
|
||||
|
@AllArgsConstructor |
||||
|
@Slf4j |
||||
|
public class NotificationTargetEdgeEventFetcher extends BasePageableEdgeEventFetcher<NotificationTarget> { |
||||
|
|
||||
|
private NotificationTargetService notificationTargetService; |
||||
|
|
||||
|
@Override |
||||
|
PageData<NotificationTarget> fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) { |
||||
|
return notificationTargetService.findNotificationTargetsByTenantId(tenantId, pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, NotificationTarget notificationTarget) { |
||||
|
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.NOTIFICATION_TARGET, |
||||
|
EdgeEventActionType.ADDED, notificationTarget.getId(), null); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,51 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.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.notification.NotificationType; |
||||
|
import org.thingsboard.server.common.data.notification.template.NotificationTemplate; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.page.PageLink; |
||||
|
import org.thingsboard.server.dao.notification.NotificationTemplateService; |
||||
|
|
||||
|
import java.util.List; |
||||
|
|
||||
|
@AllArgsConstructor |
||||
|
@Slf4j |
||||
|
public class NotificationTemplateEdgeEventFetcher extends BasePageableEdgeEventFetcher<NotificationTemplate> { |
||||
|
|
||||
|
private NotificationTemplateService notificationTemplateService; |
||||
|
|
||||
|
@Override |
||||
|
PageData<NotificationTemplate> fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) { |
||||
|
return notificationTemplateService.findNotificationTemplatesByTenantIdAndNotificationTypes(tenantId, List.of(NotificationType.values()), pageLink); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, NotificationTemplate notificationTemplate) { |
||||
|
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.NOTIFICATION_TEMPLATE, |
||||
|
EdgeEventActionType.ADDED, notificationTemplate.getId(), null); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,119 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.notification; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.common.data.EdgeUtils; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEvent; |
||||
|
import org.thingsboard.server.common.data.id.NotificationRuleId; |
||||
|
import org.thingsboard.server.common.data.id.NotificationTargetId; |
||||
|
import org.thingsboard.server.common.data.id.NotificationTemplateId; |
||||
|
import org.thingsboard.server.common.data.notification.rule.NotificationRule; |
||||
|
import org.thingsboard.server.common.data.notification.targets.NotificationTarget; |
||||
|
import org.thingsboard.server.common.data.notification.template.NotificationTemplate; |
||||
|
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationRuleUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationTargetUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.NotificationTemplateUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Component |
||||
|
@TbCoreComponent |
||||
|
public class NotificationEdgeProcessor extends BaseEdgeProcessor { |
||||
|
|
||||
|
public DownlinkMsg convertNotificationRuleToDownlink(EdgeEvent edgeEvent) { |
||||
|
NotificationRuleId notificationRuleId = new NotificationRuleId(edgeEvent.getEntityId()); |
||||
|
DownlinkMsg downlinkMsg = null; |
||||
|
switch (edgeEvent.getAction()) { |
||||
|
case ADDED, UPDATED -> { |
||||
|
NotificationRule notificationRule = notificationRuleService.findNotificationRuleById(edgeEvent.getTenantId(), notificationRuleId); |
||||
|
if (notificationRule != null) { |
||||
|
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
||||
|
NotificationRuleUpdateMsg notificationRuleUpdateMsg = notificationMsgConstructor.constructNotificationRuleUpdateMsg(msgType, notificationRule); |
||||
|
downlinkMsg = DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addNotificationRuleUpdateMsg(notificationRuleUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
|
case DELETED -> { |
||||
|
NotificationRuleUpdateMsg notificationRuleUpdateMsg = notificationMsgConstructor.constructNotificationRuleDeleteMsg(notificationRuleId); |
||||
|
downlinkMsg = DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addNotificationRuleUpdateMsg(notificationRuleUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
|
return downlinkMsg; |
||||
|
} |
||||
|
|
||||
|
public DownlinkMsg convertNotificationTargetToDownlink(EdgeEvent edgeEvent) { |
||||
|
NotificationTargetId notificationTargetId = new NotificationTargetId(edgeEvent.getEntityId()); |
||||
|
DownlinkMsg downlinkMsg = null; |
||||
|
switch (edgeEvent.getAction()) { |
||||
|
case ADDED, UPDATED -> { |
||||
|
NotificationTarget notificationTarget = notificationTargetService.findNotificationTargetById(edgeEvent.getTenantId(), notificationTargetId); |
||||
|
if (notificationTarget != null) { |
||||
|
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
||||
|
NotificationTargetUpdateMsg notificationTargetUpdateMsg = notificationMsgConstructor.constructNotificationTargetUpdateMsg(msgType, notificationTarget); |
||||
|
downlinkMsg = DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addNotificationTargetUpdateMsg(notificationTargetUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
|
case DELETED -> { |
||||
|
NotificationTargetUpdateMsg notificationTargetUpdateMsg = notificationMsgConstructor.constructNotificationTargetDeleteMsg(notificationTargetId); |
||||
|
downlinkMsg = DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addNotificationTargetUpdateMsg(notificationTargetUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
|
return downlinkMsg; |
||||
|
} |
||||
|
|
||||
|
public DownlinkMsg convertNotificationTemplateToDownlink(EdgeEvent edgeEvent) { |
||||
|
NotificationTemplateId notificationTemplateId = new NotificationTemplateId(edgeEvent.getEntityId()); |
||||
|
DownlinkMsg downlinkMsg = null; |
||||
|
switch (edgeEvent.getAction()) { |
||||
|
case ADDED, UPDATED -> { |
||||
|
NotificationTemplate notificationTemplate = notificationTemplateService.findNotificationTemplateById(edgeEvent.getTenantId(), notificationTemplateId); |
||||
|
if (notificationTemplate != null) { |
||||
|
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
||||
|
NotificationTemplateUpdateMsg notificationTemplateUpdateMsg = notificationMsgConstructor.constructNotificationTemplateUpdateMsg(msgType, notificationTemplate); |
||||
|
downlinkMsg = DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addNotificationTemplateUpdateMsg(notificationTemplateUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
|
case DELETED -> { |
||||
|
NotificationTemplateUpdateMsg notificationTemplateUpdateMsg = notificationMsgConstructor.constructNotificationTemplateDeleteMsg(notificationTemplateId); |
||||
|
downlinkMsg = DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addNotificationTemplateUpdateMsg(notificationTemplateUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
|
return downlinkMsg; |
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue