package jnpf.limsService.impl;
|
|
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
import jnpf.base.UserInfo;
|
import jnpf.limsEntity.LimsMessageDeliveryEntity;
|
import jnpf.limsEntity.LimsMessageDeliveryStatus;
|
import jnpf.limsEntity.LimsMessageEventEntity;
|
import jnpf.limsEntity.LimsMessageEventStatus;
|
import jnpf.limsMapper.LimsMessageDeliveryMapper;
|
import jnpf.limsMapper.LimsMessageEventMapper;
|
import jnpf.limsService.client.LimsMessageClient;
|
import jnpf.message.model.SentMessageForm;
|
import jnpf.util.JsonUtil;
|
import lombok.RequiredArgsConstructor;
|
import lombok.extern.slf4j.Slf4j;
|
import org.springframework.beans.factory.annotation.Value;
|
import org.springframework.scheduling.annotation.Scheduled;
|
import org.springframework.stereotype.Service;
|
import org.springframework.util.StringUtils;
|
|
import java.util.Collections;
|
import java.util.Date;
|
import java.util.HashMap;
|
import java.util.List;
|
import java.util.Map;
|
|
@Slf4j
|
@Service
|
@RequiredArgsConstructor
|
public class LimsMessageRetryService {
|
private final LimsMessageEventMapper eventMapper;
|
private final LimsMessageDeliveryMapper deliveryMapper;
|
private final LimsMessageClient messageClient;
|
|
@Value("${lims.message.enabled:false}")
|
private boolean messageEnabled;
|
@Value("${lims.message.water-sampling-enabled:false}")
|
private boolean waterSamplingMessageEnabled;
|
@Value("${lims.message.retry.batch-size:100}")
|
private int batchSize;
|
@Value("${lims.message.retry.max-attempts:5}")
|
private int maxAttempts;
|
@Value("${lims.message.retry.stale-minutes:15}")
|
private int staleMinutes;
|
|
@Scheduled(cron = "${lims.message.retry.cron:0 */5 * * * ?}")
|
public void retryFailedDeliveries() {
|
if (!messageEnabled && !waterSamplingMessageEnabled) {
|
return;
|
}
|
Date staleBefore = new Date(System.currentTimeMillis() - staleMinutes * 60_000L);
|
List<LimsMessageDeliveryEntity> deliveries = deliveryMapper.selectRetryable(
|
staleBefore, batchSize, maxAttempts);
|
for (LimsMessageDeliveryEntity delivery : deliveries) {
|
if (delivery.getRetryCount() != null && delivery.getRetryCount() >= maxAttempts) {
|
continue;
|
}
|
if (deliveryMapper.claimRetry(delivery.getId(), staleBefore) == 0) {
|
continue;
|
}
|
retryOne(delivery);
|
}
|
}
|
|
@SuppressWarnings("unchecked")
|
private void retryOne(LimsMessageDeliveryEntity delivery) {
|
LimsMessageEventEntity event = eventMapper.selectEvent(delivery.getEventId());
|
if (event == null) {
|
fail(delivery, "MESSAGE_EVENT_NOT_FOUND");
|
return;
|
}
|
Map<String, Object> parameters = StringUtils.hasText(event.getParameterSnapshot())
|
? JsonUtil.getJsonToBean(event.getParameterSnapshot(), Map.class) : new HashMap<>();
|
parameters.put("deliveryKey", delivery.getDeliveryKey());
|
try {
|
List<String> errors = messageClient.sendScheduleMessage(buildForm(event, delivery, parameters));
|
if (errors != null && !errors.isEmpty()) {
|
fail(delivery, String.join("; ", errors));
|
} else {
|
delivery.setStatus(LimsMessageDeliveryStatus.SENT);
|
delivery.setSentTime(new Date());
|
delivery.setErrorCode(null);
|
delivery.setErrorMessage(null);
|
delivery.setLastModifyTime(new Date());
|
deliveryMapper.updateById(delivery);
|
}
|
} catch (Exception e) {
|
fail(delivery, e.getMessage());
|
log.error("[lims-message-retry] failed delivery={}", delivery.getId(), e);
|
}
|
aggregate(event);
|
}
|
|
private SentMessageForm buildForm(LimsMessageEventEntity event, LimsMessageDeliveryEntity delivery,
|
Map<String, Object> parameters) {
|
UserInfo userInfo = new UserInfo();
|
userInfo.setUserId("system");
|
userInfo.setUserName("LIMS");
|
userInfo.setTenantId(event.getTenantId());
|
SentMessageForm form = new SentMessageForm();
|
form.setToUserIds(Collections.singletonList(delivery.getRecipientUserId()));
|
form.setTemplateId(event.getEventCode());
|
form.setParameterMap(parameters);
|
form.setContentMsg(new HashMap<>());
|
form.setUserInfo(userInfo);
|
return form;
|
}
|
|
private void fail(LimsMessageDeliveryEntity delivery, String error) {
|
int retries = delivery.getRetryCount() == null ? 1 : delivery.getRetryCount() + 1;
|
delivery.setStatus(LimsMessageDeliveryStatus.FAILED);
|
delivery.setRetryCount(retries);
|
delivery.setNextRetryTime(retries >= maxAttempts ? null
|
: new Date(System.currentTimeMillis() + retryDelayMillis(retries)));
|
delivery.setErrorCode("SEND_FAILED");
|
delivery.setErrorMessage(limit(error));
|
delivery.setLastModifyTime(new Date());
|
deliveryMapper.updateById(delivery);
|
}
|
|
private void aggregate(LimsMessageEventEntity event) {
|
List<LimsMessageDeliveryEntity> all = deliveryMapper.selectList(
|
new LambdaQueryWrapper<LimsMessageDeliveryEntity>()
|
.eq(LimsMessageDeliveryEntity::getEventId, event.getId()));
|
int sent = (int) all.stream().filter(item -> item.getStatus() == LimsMessageDeliveryStatus.SENT).count();
|
int failed = all.size() - sent;
|
event.setRecipientCount(all.size());
|
event.setSentCount(sent);
|
event.setFailedCount(failed);
|
event.setStatus(failed == 0 ? LimsMessageEventStatus.COMPLETED
|
: sent == 0 ? LimsMessageEventStatus.FAILED : LimsMessageEventStatus.PARTIAL_FAILED);
|
event.setLastModifyTime(new Date());
|
if (failed == 0) {
|
event.setCompletedTime(new Date());
|
}
|
eventMapper.updateById(event);
|
}
|
|
private long retryDelayMillis(int retries) {
|
int minutes = Math.min(60, 1 << Math.min(retries, 6));
|
return minutes * 60_000L;
|
}
|
|
private String limit(String value) {
|
String message = StringUtils.hasText(value) ? value : "Unknown send error";
|
return message.length() <= 1000 ? message : message.substring(0, 1000);
|
}
|
}
|