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