|
|
|
@ -26,7 +26,6 @@ import org.thingsboard.server.cluster.TbClusterService; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
|
import org.thingsboard.server.common.data.alarm.Alarm; |
|
|
|
import org.thingsboard.server.common.data.alarm.AlarmInfo; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
@ -41,6 +40,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
|
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
|
|
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
|
|
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
|
|
|
import org.thingsboard.server.common.data.alarm.AlarmAssigneeUpdate; |
|
|
|
import org.thingsboard.server.dao.attributes.AttributesService; |
|
|
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.LocalSubscriptionServiceMsgProto; |
|
|
|
@ -269,7 +269,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
updateDeviceInactivityTimeout(tenantId, entityId, attributes); |
|
|
|
} else if (TbAttributeSubscriptionScope.SHARED_SCOPE.name().equalsIgnoreCase(scope) && notifyDevice) { |
|
|
|
clusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onUpdate(tenantId, |
|
|
|
new DeviceId(entityId.getId()), DataConstants.SHARED_SCOPE, new ArrayList<>(attributes)) |
|
|
|
new DeviceId(entityId.getId()), DataConstants.SHARED_SCOPE, new ArrayList<>(attributes)) |
|
|
|
, null); |
|
|
|
} |
|
|
|
} |
|
|
|
@ -293,7 +293,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onAlarmUpdate(TenantId tenantId, EntityId entityId, Alarm alarm, TbCallback callback) { |
|
|
|
public void onAlarmUpdate(TenantId tenantId, EntityId entityId, Alarm alarm, AlarmAssigneeUpdate assignee, TbCallback callback) { |
|
|
|
onLocalAlarmSubUpdate(entityId, |
|
|
|
s -> { |
|
|
|
if (TbSubscriptionType.ALARMS.equals(s.getType())) { |
|
|
|
@ -303,14 +303,13 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
} |
|
|
|
}, |
|
|
|
s -> alarm.getCreatedTime() >= s.getTs() || alarm.getAssignTs() >= s.getTs(), |
|
|
|
s -> alarm, |
|
|
|
false |
|
|
|
alarm, assignee, false |
|
|
|
); |
|
|
|
callback.onSuccess(); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onAlarmDeleted(TenantId tenantId, EntityId entityId, Alarm alarmInfo, TbCallback callback) { |
|
|
|
public void onAlarmDeleted(TenantId tenantId, EntityId entityId, Alarm alarm, TbCallback callback) { |
|
|
|
onLocalAlarmSubUpdate(entityId, |
|
|
|
s -> { |
|
|
|
if (TbSubscriptionType.ALARMS.equals(s.getType())) { |
|
|
|
@ -319,9 +318,8 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
return null; |
|
|
|
} |
|
|
|
}, |
|
|
|
s -> alarmInfo.getCreatedTime() >= s.getTs(), |
|
|
|
s -> alarmInfo, |
|
|
|
true |
|
|
|
s -> alarm.getCreatedTime() >= s.getTs(), |
|
|
|
alarm, null, true |
|
|
|
); |
|
|
|
callback.onSuccess(); |
|
|
|
} |
|
|
|
@ -355,7 +353,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
deleteDeviceInactivityTimeout(tenantId, entityId, keys); |
|
|
|
} else if (TbAttributeSubscriptionScope.SHARED_SCOPE.name().equalsIgnoreCase(scope) && notifyDevice) { |
|
|
|
clusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete(tenantId, |
|
|
|
new DeviceId(entityId.getId()), scope, keys), null); |
|
|
|
new DeviceId(entityId.getId()), scope, keys), null); |
|
|
|
} |
|
|
|
} |
|
|
|
callback.onSuccess(); |
|
|
|
@ -415,20 +413,21 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
private void onLocalAlarmSubUpdate(EntityId entityId, |
|
|
|
Function<TbSubscription, TbAlarmsSubscription> castFunction, |
|
|
|
Predicate<TbAlarmsSubscription> filterFunction, |
|
|
|
Function<TbAlarmsSubscription, Alarm> processFunction, |
|
|
|
Alarm alarm, AlarmAssigneeUpdate assignee, |
|
|
|
boolean deleted) { |
|
|
|
Set<TbSubscription> entitySubscriptions = subscriptionsByEntityId.get(entityId); |
|
|
|
if (alarm == null) { |
|
|
|
log.warn("[{}] empty alarm update!", entityId); |
|
|
|
return; |
|
|
|
} |
|
|
|
if (entitySubscriptions != null) { |
|
|
|
entitySubscriptions.stream().map(castFunction).filter(Objects::nonNull).filter(filterFunction).forEach(s -> { |
|
|
|
Alarm alarm = processFunction.apply(s); |
|
|
|
if (alarm != null) { |
|
|
|
if (serviceId.equals(s.getServiceId())) { |
|
|
|
AlarmSubscriptionUpdate update = new AlarmSubscriptionUpdate(s.getSubscriptionId(), alarm, deleted); |
|
|
|
localSubscriptionService.onSubscriptionUpdate(s.getSessionId(), update, TbCallback.EMPTY); |
|
|
|
} else { |
|
|
|
TopicPartitionInfo tpi = notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, s.getServiceId()); |
|
|
|
toCoreNotificationsProducer.send(tpi, toProto(s, alarm, deleted), null); |
|
|
|
} |
|
|
|
if (serviceId.equals(s.getServiceId())) { |
|
|
|
AlarmSubscriptionUpdate update = new AlarmSubscriptionUpdate(s.getSubscriptionId(), alarm, assignee, deleted); |
|
|
|
localSubscriptionService.onSubscriptionUpdate(s.getSessionId(), update, TbCallback.EMPTY); |
|
|
|
} else { |
|
|
|
TopicPartitionInfo tpi = notificationsTopicService.getNotificationsTopic(ServiceType.TB_CORE, s.getServiceId()); |
|
|
|
toCoreNotificationsProducer.send(tpi, toProto(s, alarm, assignee, deleted), null); |
|
|
|
} |
|
|
|
}); |
|
|
|
} else { |
|
|
|
@ -559,22 +558,25 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
|
}); |
|
|
|
|
|
|
|
ToCoreNotificationMsg toCoreMsg = ToCoreNotificationMsg.newBuilder().setToLocalSubscriptionServiceMsg( |
|
|
|
LocalSubscriptionServiceMsgProto.newBuilder().setSubUpdate(builder.build()).build()) |
|
|
|
LocalSubscriptionServiceMsgProto.newBuilder().setSubUpdate(builder.build()).build()) |
|
|
|
.build(); |
|
|
|
return new TbProtoQueueMsg<>(subscription.getEntityId().getId(), toCoreMsg); |
|
|
|
} |
|
|
|
|
|
|
|
private TbProtoQueueMsg<ToCoreNotificationMsg> toProto(TbSubscription subscription, Alarm alarm, boolean deleted) { |
|
|
|
private TbProtoQueueMsg<ToCoreNotificationMsg> toProto(TbSubscription subscription, Alarm alarm, AlarmAssigneeUpdate assignee, boolean deleted) { |
|
|
|
TbAlarmSubscriptionUpdateProto.Builder builder = TbAlarmSubscriptionUpdateProto.newBuilder(); |
|
|
|
|
|
|
|
builder.setSessionId(subscription.getSessionId()); |
|
|
|
builder.setSubscriptionId(subscription.getSubscriptionId()); |
|
|
|
builder.setAlarm(JacksonUtil.toString(alarm)); |
|
|
|
if (assignee != null) { |
|
|
|
builder.setAssignee(JacksonUtil.toString(assignee)); |
|
|
|
} |
|
|
|
builder.setDeleted(deleted); |
|
|
|
|
|
|
|
ToCoreNotificationMsg toCoreMsg = ToCoreNotificationMsg.newBuilder().setToLocalSubscriptionServiceMsg( |
|
|
|
LocalSubscriptionServiceMsgProto.newBuilder() |
|
|
|
.setAlarmSubUpdate(builder.build()).build()) |
|
|
|
LocalSubscriptionServiceMsgProto.newBuilder() |
|
|
|
.setAlarmSubUpdate(builder.build()).build()) |
|
|
|
.build(); |
|
|
|
return new TbProtoQueueMsg<>(subscription.getEntityId().getId(), toCoreMsg); |
|
|
|
} |
|
|
|
|