From 34981c30a78e8bbd7791131059a9210f9928b62c Mon Sep 17 00:00:00 2001
From: 刘光辉 <347230014@qq.com>
Date: 星期四, 17 九月 2026 09:24:09 +0800
Subject: [PATCH] Merge remote-tracking branch 'origin/master' into master

---
 jnpf-lims/jnpf-lims-biz/src/main/java/jnpf/limsService/impl/LimsMessageRetryService.java |  148 +++++++++++++++++++++++++++++++++++++++++++++++++
 1 files changed, 148 insertions(+), 0 deletions(-)

diff --git a/jnpf-lims/jnpf-lims-biz/src/main/java/jnpf/limsService/impl/LimsMessageRetryService.java b/jnpf-lims/jnpf-lims-biz/src/main/java/jnpf/limsService/impl/LimsMessageRetryService.java
new file mode 100644
index 0000000..6c747e1
--- /dev/null
+++ b/jnpf-lims/jnpf-lims-biz/src/main/java/jnpf/limsService/impl/LimsMessageRetryService.java
@@ -0,0 +1,148 @@
+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);
+    }
+}

--
Gitblit v1.8.0