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');