30 changed files with 315 additions and 159 deletions
@ -0,0 +1,72 @@ |
|||||
|
/** |
||||
|
* 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.notification.NotificationDao; |
||||
|
import org.thingsboard.server.dao.notification.NotificationRequestDao; |
||||
|
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.NOTIFICATION_TABLE_NAME; |
||||
|
|
||||
|
@Service |
||||
|
@ConditionalOnExpression("${sql.ttl.notifications.enabled:true} && ${sql.ttl.notifications.ttl:0} > 0") |
||||
|
@Slf4j |
||||
|
public class NotificationsCleanUpService extends AbstractCleanUpService { |
||||
|
|
||||
|
private final NotificationDao notificationDao; |
||||
|
private final NotificationRequestDao notificationRequestDao; |
||||
|
private final SqlPartitioningRepository partitioningRepository; |
||||
|
|
||||
|
@Value("${sql.ttl.notifications.ttl:2592000}") |
||||
|
private long ttlInSec; |
||||
|
@Value("${sql.notifications.partition_size:168}") |
||||
|
private int partitionSizeInHours; |
||||
|
|
||||
|
public NotificationsCleanUpService(PartitionService partitionService, |
||||
|
NotificationDao notificationDao, |
||||
|
NotificationRequestDao notificationRequestDao, |
||||
|
SqlPartitioningRepository partitioningRepository) { |
||||
|
super(partitionService); |
||||
|
this.notificationDao = notificationDao; |
||||
|
this.notificationRequestDao = notificationRequestDao; |
||||
|
this.partitioningRepository = partitioningRepository; |
||||
|
} |
||||
|
|
||||
|
@Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.notifications.checking_interval_ms:86400000})}", |
||||
|
fixedDelayString = "${sql.ttl.notifications.checking_interval_ms:86400000}") |
||||
|
public void cleanUp() { |
||||
|
long expTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttlInSec); |
||||
|
long partitionDurationMs = TimeUnit.HOURS.toMillis(partitionSizeInHours); |
||||
|
if (isSystemTenantPartitionMine()) { |
||||
|
long actualExpTime = partitioningRepository.getLastPartitionEnd(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); |
||||
|
|
||||
|
partitioningRepository.dropPartitionsBefore(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); |
||||
|
// select distinct request_id for period and manually delete them ?
|
||||
|
// sql trigger ? ----
|
||||
|
} else { |
||||
|
partitioningRepository.cleanupPartitionsCache(NOTIFICATION_TABLE_NAME, expTime, partitionDurationMs); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -1,35 +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.common.data.notification; |
|
||||
|
|
||||
import lombok.Data; |
|
||||
import lombok.EqualsAndHashCode; |
|
||||
|
|
||||
import java.util.Map; |
|
||||
|
|
||||
@Data |
|
||||
@EqualsAndHashCode(callSuper = true) |
|
||||
public class NotificationRequestInfo extends NotificationRequest { |
|
||||
|
|
||||
private int sent; |
|
||||
private int read; |
|
||||
private Map<String, NotificationStatus> statusesByRecipient; |
|
||||
|
|
||||
public NotificationRequestInfo(NotificationRequest notificationRequest) { |
|
||||
super(notificationRequest); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -0,0 +1,54 @@ |
|||||
|
/** |
||||
|
* 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.common.data.notification; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonCreator; |
||||
|
import com.fasterxml.jackson.annotation.JsonProperty; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.common.data.User; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
|
||||
|
@Data |
||||
|
public class NotificationRequestStats { |
||||
|
|
||||
|
private final Map<NotificationDeliveryMethod, AtomicInteger> sent; |
||||
|
private final Map<NotificationDeliveryMethod, Map<String, String>> errors; |
||||
|
|
||||
|
public NotificationRequestStats() { |
||||
|
this.sent = new ConcurrentHashMap<>(); |
||||
|
this.errors = new ConcurrentHashMap<>(); |
||||
|
} |
||||
|
|
||||
|
@JsonCreator |
||||
|
public NotificationRequestStats(@JsonProperty("sent") Map<NotificationDeliveryMethod, AtomicInteger> sent, |
||||
|
@JsonProperty("errors") Map<NotificationDeliveryMethod, Map<String, String>> errors) { |
||||
|
this.sent = sent; |
||||
|
this.errors = errors; |
||||
|
} |
||||
|
|
||||
|
public void reportSent(NotificationDeliveryMethod deliveryMethod) { |
||||
|
sent.computeIfAbsent(deliveryMethod, k -> new AtomicInteger()).incrementAndGet(); |
||||
|
} |
||||
|
|
||||
|
public void reportError(NotificationDeliveryMethod deliveryMethod, User recipient, Throwable error) { |
||||
|
String errorMessage = error.getMessage(); |
||||
|
errors.computeIfAbsent(deliveryMethod, k -> new ConcurrentHashMap<>()).put(recipient.getEmail(), errorMessage); |
||||
|
} |
||||
|
|
||||
|
} |
||||
Loading…
Reference in new issue