import { writeFileSync } from 'node:fs'; import { fileURLToPath } from 'node:url'; import path from 'node:path'; const base = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../src/main/java/com/fs/company/service/easycall'); const helper = `package com.fs.company.service.easycall; import com.fs.common.core.redis.RedisCacheT; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Component; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.TimeUnit; /** * ${'\u5916\u547C\u9650\u6D41\u91CD\u8BD5\u961F\u5217\u4E13\u7528 Redis \u64CD\u4F5C\uFF08ZSET \u5F39\u51FA\u3001Hash \u8F7D\u8377\u3001\u5904\u7406\u9501\u3001\u7F51\u5173\u9650\u901F\u7B49\uFF09\u3002'} *

* ${'\u4E1A\u52A1\u8BED\u4E49\u4E0D\u8FDB'} {@link com.fs.common.core.redis.RedisCache} / {@link RedisCacheT} ${'\u901A\u7528\u5DE5\u5177\u7C7B\u3002'} *

* ${'\u4E0D\u4F7F\u7528 Lua \u811A\u672C\uFF1ARedisson \u4F20\u5165 ARGV \u65F6\u53EF\u80FD\u5E26\u5F15\u53F7/\u5B57\u8282\u6570\u7EC4\uFF0C\u89E6\u53D1'} ERR Lua redis() command arguments must be strings or integers。 */ @Component public class OutboundLimitRetryRedisHelper { @Autowired private RedisCacheT redisCacheT; /** * Pop due members via Java ZSet API (Redisson-safe; same as WxAddWxFrequencyRedisHelper). */ @SuppressWarnings({"unchecked", "rawtypes"}) public List popDueMembersByScore(String zsetKey, double maxScore, int count) { if (count <= 0) { return Collections.emptyList(); } RedisTemplate redisTemplate = redisCacheT.redisTemplate; Set members = redisTemplate.opsForZSet().rangeByScore(zsetKey, 0, maxScore, 0, count); if (members == null || members.isEmpty()) { return Collections.emptyList(); } List dueMembers = new ArrayList<>(members.size()); for (Object member : members) { Long removed = redisTemplate.opsForZSet().remove(zsetKey, member); if (removed != null && removed > 0) { dueMembers.add(String.valueOf(member)); } } return dueMembers; } public Double zScore(String zsetKey, Object member) { return redisCacheT.redisTemplate.opsForZSet().score(zsetKey, member); } public void zSetAdd(String zsetKey, Object member, double score) { redisCacheT.redisTemplate.opsForZSet().add(zsetKey, member, score); } public void hashPut(String hashKey, String field, String value) { redisCacheT.setCacheMapValue(hashKey, field, value); } public String hashGet(String hashKey, String field) { return redisCacheT.getCacheMapValue(hashKey, field); } public Map hashGetAll(String hashKey) { Map map = redisCacheT.getCacheMap(hashKey); return map != null ? map : Collections.emptyMap(); } public void hashDelete(String hashKey, String field) { redisCacheT.redisTemplate.opsForHash().delete(hashKey, field); } public boolean setIfAbsent(String key, String value, long timeout, TimeUnit unit) { return Boolean.TRUE.equals( redisCacheT.redisTemplate.opsForValue().setIfAbsent(key, value, timeout, unit)); } public void deleteKey(String key) { redisCacheT.deleteObject(key); } public String getString(String key) { return redisCacheT.getCacheObject(key); } public void setString(String key, String value, long timeout, TimeUnit unit) { redisCacheT.setCacheObject(key, value, timeout, unit); } public Long incr(String key, long delta) { return redisCacheT.redisTemplate.opsForValue().increment(key, delta); } public boolean expire(String key, long timeout, TimeUnit unit) { return redisCacheT.expire(key, timeout, unit); } } `; const support = `package com.fs.company.service.easycall; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.fs.common.utils.StringUtils; import com.fs.company.vo.easycall.EasyCallCommonAddCallListParam; import com.fs.company.vo.easycall.EasyCallPhoneItemVO; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import java.util.Date; import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; /** * ${'\u5916\u547C\u7EBF\u8DEF\u9650\u6D41\u91CD\u8BD5\u961F\u5217\uFF1A\u53BB\u91CD\u5165\u961F\u3001\u6296\u52A8\u5206\u6563\u3001\u539F\u5B50\u51FA\u961F\u3001\u5931\u8D25\u9000\u907F\u91CD\u5165\u961F\u3002'} */ @Slf4j @Component public class OutboundLimitRetrySupport { /** ${'\u4E0E EasyCallServiceImpl \u4FDD\u6301\u4E00\u81F4'} */ public static final String OUTBOUND_LIMIT_REDIS_PREFIX = "outbound:limit:"; public static final String RETRY_ZSET_KEY = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:zset"; public static final String RETRY_PAYLOAD_HASH_KEY = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:payload"; private static final String RETRY_PROCESSING_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:processing:"; private static final String RETRY_DISPATCHED_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:dispatched:"; /** ${'\u5DF2\u6210\u529F append+startTask \u7684 callBackUuid \u6807\u8BB0 TTL\uFF08\u9632 pop \u540E\u5D29\u6E83\u5BFC\u81F4\u91CD\u590D\u5916\u547C\uFF09'} */ private static final int DISPATCHED_TTL_HOURS = 24; private static final String RETRY_RATE_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:rate:"; /** ${'\u5230\u671F\u540E\u968F\u673A\u5EF6\u8FDF\u4E0A\u9650\uFF08\u6BEB\u79D2\uFF09\uFF0C\u5E73\u6ED1\u7A97\u53E3\u8FB9\u754C\u6D41\u91CF'} */ private static final int RETRY_JITTER_MAX_MS = 30_000; /** ${'\u5B64\u513F payload \u5355\u6B21\u6700\u591A\u6062\u590D\u6761\u6570'} */ private static final int MAX_ORPHAN_RECOVER_BATCH = 100; /** ${'\u5931\u8D25\u9000\u907F\u57FA\u7840\u5EF6\u8FDF\uFF08\u6BEB\u79D2\uFF09'} */ private static final long RETRY_BACKOFF_MS = 60_000L; /** ${'\u5355\u7F51\u5173\u6BCF\u79D2\u6700\u591A\u91CD\u8BD5\u5916\u547C\u6B21\u6570\uFF08\u591A\u5B9E\u4F8B\u5171\u4EAB\uFF09'} */ private static final int MAX_RETRY_PER_GATEWAY_PER_SEC = 30; /** ${'\u5904\u7406\u9501 TTL\uFF0C\u9632\u6B62\u5B9E\u4F8B\u5D29\u6E83\u540E\u957F\u671F\u5360\u9501'} */ private static final int PROCESSING_LOCK_SECONDS = 600; @Autowired private OutboundLimitRetryRedisHelper retryRedisHelper; /** * ${'\u9650\u6D41\u540E\u5165\u961F\uFF08member=callBackUuid\uFF0C\u540C UUID \u4E0D\u91CD\u590D\u5165\u961F\uFF0Cscore \u5E26\u6296\u52A8\uFF09\u3002'} */ public void enqueueRetry(Long companyId, Long gatewayId, EasyCallCommonAddCallListParam param, Date nextAvailableTime) { if (param == null || nextAvailableTime == null) { return; } try { String callBackUuid = extractCallBackUuid(param); String member = StringUtils.isNotBlank(callBackUuid) ? callBackUuid : UUID.randomUUID().toString(); JSONObject retryData = new JSONObject(); retryData.put("companyId", companyId); retryData.put("gatewayId", gatewayId); retryData.put("param", param); retryData.put("callBackUuid", member); retryData.put("nextAvailableTime", nextAvailableTime.getTime()); retryData.put("createTime", System.currentTimeMillis()); String existingPayload = retryRedisHelper.hashGet(RETRY_PAYLOAD_HASH_KEY, member); if (existingPayload != null) { JSONObject oldData = JSON.parseObject(existingPayload); if (oldData != null && oldData.containsKey("failCount")) { retryData.put("failCount", oldData.getIntValue("failCount")); } } double score = computeScoreWithJitter(nextAvailableTime.getTime()); Double existingScore = retryRedisHelper.zScore(RETRY_ZSET_KEY, member); if (existingScore != null && existingScore <= score) { retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString()); log.info("enqueueRetry: ${'\u5DF2\u5B58\u5728\u66F4\u65E9\u6216\u76F8\u540C\u8C03\u5EA6\uFF0C\u4EC5\u5237\u65B0 payload'} - member={}", member); return; } retryRedisHelper.zSetAdd(RETRY_ZSET_KEY, member, score); retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString()); log.info("enqueueRetry: member={}, score={}, companyId={}, gatewayId={}", member, (long) score, companyId, gatewayId); } catch (Exception e) { log.error("enqueueRetry ${'\u5F02\u5E38'} companyId={}, gatewayId={}", companyId, gatewayId, e); } } /** * ${'\u5F39\u51FA\u5DF2\u5230\u671F\u7684\u91CD\u8BD5 member\uFF08Java ZSet rangeByScore + remove\uFF0C\u4E0E\u52A0\u5FAE\u9891\u7387\u91CD\u8BD5\u540C\u6A21\u5F0F\uFF09\u3002'} */ public List popDueMembers(int batchSize) { long now = System.currentTimeMillis(); return retryRedisHelper.popDueMembersByScore(RETRY_ZSET_KEY, now, batchSize); } /** * ${'\u89E3\u6790\u4EFB\u52A1 payload\uFF1B\u517C\u5BB9\u65E7\u7248 member \u4E3A\u6574\u6BB5 JSON \u7684\u5199\u6CD5\u3002'} */ public JSONObject resolvePayload(String member) { if (StringUtils.isBlank(member)) { return null; } if (member.trim().startsWith("{")) { return JSON.parseObject(member); } String raw = retryRedisHelper.hashGet(RETRY_PAYLOAD_HASH_KEY, member); if (raw == null) { return null; } return JSON.parseObject(raw); } /** * ${'\u5C1D\u8BD5\u83B7\u53D6\u5904\u7406\u9501\uFF0C\u907F\u514D\u6781\u7AEF\u60C5\u51B5\u4E0B\u91CD\u590D\u5916\u547C\u3002'} */ public boolean tryAcquireProcessingLock(String callBackUuid) { if (StringUtils.isBlank(callBackUuid)) { return true; } String lockKey = RETRY_PROCESSING_PREFIX + callBackUuid; return retryRedisHelper.setIfAbsent(lockKey, "1", PROCESSING_LOCK_SECONDS, TimeUnit.SECONDS); } public void releaseProcessingLock(String callBackUuid) { if (StringUtils.isBlank(callBackUuid)) { return; } retryRedisHelper.deleteKey(RETRY_PROCESSING_PREFIX + callBackUuid); } /** * ${'\u662F\u5426\u5DF2\u6210\u529F\u5B8C\u6210 append+startTask\uFF08\u9632 pop \u540E\u8FDB\u7A0B\u5D29\u6E83\u518D\u6B21\u91CD\u8BD5\u5BFC\u81F4\u91CD\u590D\u5916\u547C\uFF09\u3002'} */ public boolean isDispatched(String callBackUuid) { if (StringUtils.isBlank(callBackUuid)) { return false; } return retryRedisHelper.getString(RETRY_DISPATCHED_PREFIX + callBackUuid) != null; } /** * ${'\u6807\u8BB0 callBackUuid \u5DF2\u5B8C\u6210\u6D3E\u53D1\uFF08append \u6210\u529F\u4E14 startTask \u5DF2\u89E6\u53D1\uFF09\u3002'} */ public void markDispatched(String callBackUuid) { if (StringUtils.isBlank(callBackUuid)) { return; } retryRedisHelper.setString(RETRY_DISPATCHED_PREFIX + callBackUuid, "1", DISPATCHED_TTL_HOURS, TimeUnit.HOURS); } /** * ${'\u6062\u590D ZSET \u5DF2\u5F39\u51FA\u4F46 payload \u4ECD\u6EDE\u7559 Hash \u7684\u5B64\u513F\u4EFB\u52A1\uFF08\u8FDB\u7A0B\u5D29\u6E83\u573A\u666F\uFF0C\u4E0D\u91CD\u8BD5\u4E22\u5355\uFF09\u3002'} * * @return ${'\u672C\u6B21\u6062\u590D\u5165\u961F\u6761\u6570'} */ public int recoverOrphanedPayloads() { Map all = retryRedisHelper.hashGetAll(RETRY_PAYLOAD_HASH_KEY); if (all == null || all.isEmpty()) { return 0; } int recovered = 0; for (Map.Entry entry : all.entrySet()) { if (recovered >= MAX_ORPHAN_RECOVER_BATCH) { break; } String member = entry.getKey(); if (retryRedisHelper.zScore(RETRY_ZSET_KEY, member) != null) { continue; } JSONObject data = JSON.parseObject(entry.getValue()); if (data == null) { ackSuccess(member); continue; } String callBackUuid = data.getString("callBackUuid"); EasyCallCommonAddCallListParam param = data.getObject("param", EasyCallCommonAddCallListParam.class); if (StringUtils.isBlank(callBackUuid) && param != null) { callBackUuid = extractCallBackUuid(param); } if (isDispatched(callBackUuid)) { ackSuccess(member); continue; } Long companyId = data.getLong("companyId"); Long gatewayId = data.getLong("gatewayId"); if (companyId == null || param == null) { ackSuccess(member); continue; } Date nextTime = new Date(System.currentTimeMillis() + RETRY_BACKOFF_MS); enqueueRetry(companyId, gatewayId, param, nextTime); recovered++; log.warn("recoverOrphanedPayloads: ${'\u5B64\u513F\u4EFB\u52A1\u91CD\u65B0\u5165\u961F'} member={}, callBackUuid={}", member, callBackUuid); } return recovered; } /** * ${'\u6309\u7F51\u5173\u4EE4\u724C\u6876\u9650\u901F\uFF0C\u8D85\u9650\u5219\u5EF6\u8FDF\u91CD\u5165\u961F\u3002'} * * @return true ${'\u8868\u793A\u5F53\u524D\u79D2\u989D\u5EA6\u5DF2\u7528\u5C3D\uFF0C\u8C03\u7528\u65B9\u5E94\u8DF3\u8FC7\u672C\u6B21\u5916\u547C'} */ public boolean isGatewayRateLimited(Long companyId, Long gatewayId) { if (companyId == null || gatewayId == null) { return false; } String bucketKey = RETRY_RATE_PREFIX + companyId + ":" + gatewayId; Long count = retryRedisHelper.incr(bucketKey, 1L); if (count == null) { count = 1L; } if (count == 1L) { retryRedisHelper.expire(bucketKey, 1, TimeUnit.SECONDS); } return count > MAX_RETRY_PER_GATEWAY_PER_SEC; } /** * ${'\u6210\u529F\u540E\u6E05\u7406 payload\uFF1Bmember \u5DF2\u5728 pop \u65F6\u4ECE ZSET \u79FB\u9664\u3002'} */ public void ackSuccess(String member) { if (StringUtils.isBlank(member) || member.trim().startsWith("{")) { return; } retryRedisHelper.hashDelete(RETRY_PAYLOAD_HASH_KEY, member); } /** * ${'\u5931\u8D25\u6216\u5F02\u5E38\uFF1A\u9000\u907F\u540E\u91CD\u65B0\u5165\u961F\uFF0C\u4E0D\u4E22\u5355\u3002'} */ public void requeueWithBackoff(JSONObject retryData, String member) { if (retryData == null) { return; } Long companyId = retryData.getLong("companyId"); Long gatewayId = retryData.getLong("gatewayId"); EasyCallCommonAddCallListParam param = retryData.getObject("param", EasyCallCommonAddCallListParam.class); if (param == null) { log.warn("requeueWithBackoff: param ${'\u4E3A\u7A7A\uFF0C\u65E0\u6CD5\u91CD\u5165\u961F'} member={}", member); return; } int failCount = retryData.getIntValue("failCount") + 1; retryData.put("failCount", failCount); if (StringUtils.isNotBlank(member) && !member.trim().startsWith("{")) { retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString()); } long delay = RETRY_BACKOFF_MS * Math.min(failCount, 5); Date nextTime = new Date(System.currentTimeMillis() + delay); enqueueRetry(companyId, gatewayId, param, nextTime); log.warn("requeueWithBackoff: member={}, failCount={}, nextDelayMs={}", member, failCount, delay); } public static String extractCallBackUuid(EasyCallCommonAddCallListParam param) { if (param == null || param.getPhoneList() == null || param.getPhoneList().isEmpty()) { return null; } EasyCallPhoneItemVO item = param.getPhoneList().get(0); if (item == null || item.getBizJson() == null) { return null; } JSONObject bizJson; if (item.getBizJson() instanceof JSONObject) { bizJson = (JSONObject) item.getBizJson(); } else { bizJson = JSON.parseObject(String.valueOf(item.getBizJson())); } return bizJson != null ? bizJson.getString("callBackUuid") : null; } private static double computeScoreWithJitter(long baseTimeMs) { int jitter = ThreadLocalRandom.current().nextInt(RETRY_JITTER_MAX_MS + 1); return baseTimeMs + jitter; } } `; writeFileSync(path.join(base, 'OutboundLimitRetryRedisHelper.java'), helper, 'utf8'); writeFileSync(path.join(base, 'OutboundLimitRetrySupport.java'), support, 'utf8'); console.log('UTF-8 files written');