刘光辉
15 小时以前 34981c30a78e8bbd7791131059a9210f9928b62c
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
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);
    }
}