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 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 parameters = StringUtils.hasText(event.getParameterSnapshot()) ? JsonUtil.getJsonToBean(event.getParameterSnapshot(), Map.class) : new HashMap<>(); parameters.put("deliveryKey", delivery.getDeliveryKey()); try { List 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 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 all = deliveryMapper.selectList( new LambdaQueryWrapper() .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); } }