fix_outbound_limit_retry_utf8.mjs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403
  1. import { writeFileSync } from 'node:fs';
  2. import { fileURLToPath } from 'node:url';
  3. import path from 'node:path';
  4. const base = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../src/main/java/com/fs/company/service/easycall');
  5. const helper = `package com.fs.company.service.easycall;
  6. import com.fs.common.core.redis.RedisCacheT;
  7. import org.springframework.beans.factory.annotation.Autowired;
  8. import org.springframework.data.redis.core.RedisTemplate;
  9. import org.springframework.stereotype.Component;
  10. import java.util.ArrayList;
  11. import java.util.Collections;
  12. import java.util.List;
  13. import java.util.Map;
  14. import java.util.Set;
  15. import java.util.concurrent.TimeUnit;
  16. /**
  17. * ${'\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'}
  18. * <p>
  19. * ${'\u4E1A\u52A1\u8BED\u4E49\u4E0D\u8FDB'} {@link com.fs.common.core.redis.RedisCache} / {@link RedisCacheT} ${'\u901A\u7528\u5DE5\u5177\u7C7B\u3002'}
  20. * <p>
  21. * ${'\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。
  22. */
  23. @Component
  24. public class OutboundLimitRetryRedisHelper {
  25. @Autowired
  26. private RedisCacheT<String> redisCacheT;
  27. /**
  28. * Pop due members via Java ZSet API (Redisson-safe; same as WxAddWxFrequencyRedisHelper).
  29. */
  30. @SuppressWarnings({"unchecked", "rawtypes"})
  31. public List<String> popDueMembersByScore(String zsetKey, double maxScore, int count) {
  32. if (count <= 0) {
  33. return Collections.emptyList();
  34. }
  35. RedisTemplate redisTemplate = redisCacheT.redisTemplate;
  36. Set<Object> members = redisTemplate.opsForZSet().rangeByScore(zsetKey, 0, maxScore, 0, count);
  37. if (members == null || members.isEmpty()) {
  38. return Collections.emptyList();
  39. }
  40. List<String> dueMembers = new ArrayList<>(members.size());
  41. for (Object member : members) {
  42. Long removed = redisTemplate.opsForZSet().remove(zsetKey, member);
  43. if (removed != null && removed > 0) {
  44. dueMembers.add(String.valueOf(member));
  45. }
  46. }
  47. return dueMembers;
  48. }
  49. public Double zScore(String zsetKey, Object member) {
  50. return redisCacheT.redisTemplate.opsForZSet().score(zsetKey, member);
  51. }
  52. public void zSetAdd(String zsetKey, Object member, double score) {
  53. redisCacheT.redisTemplate.opsForZSet().add(zsetKey, member, score);
  54. }
  55. public void hashPut(String hashKey, String field, String value) {
  56. redisCacheT.setCacheMapValue(hashKey, field, value);
  57. }
  58. public String hashGet(String hashKey, String field) {
  59. return redisCacheT.getCacheMapValue(hashKey, field);
  60. }
  61. public Map<String, String> hashGetAll(String hashKey) {
  62. Map<String, String> map = redisCacheT.getCacheMap(hashKey);
  63. return map != null ? map : Collections.emptyMap();
  64. }
  65. public void hashDelete(String hashKey, String field) {
  66. redisCacheT.redisTemplate.opsForHash().delete(hashKey, field);
  67. }
  68. public boolean setIfAbsent(String key, String value, long timeout, TimeUnit unit) {
  69. return Boolean.TRUE.equals(
  70. redisCacheT.redisTemplate.opsForValue().setIfAbsent(key, value, timeout, unit));
  71. }
  72. public void deleteKey(String key) {
  73. redisCacheT.deleteObject(key);
  74. }
  75. public String getString(String key) {
  76. return redisCacheT.getCacheObject(key);
  77. }
  78. public void setString(String key, String value, long timeout, TimeUnit unit) {
  79. redisCacheT.setCacheObject(key, value, timeout, unit);
  80. }
  81. public Long incr(String key, long delta) {
  82. return redisCacheT.redisTemplate.opsForValue().increment(key, delta);
  83. }
  84. public boolean expire(String key, long timeout, TimeUnit unit) {
  85. return redisCacheT.expire(key, timeout, unit);
  86. }
  87. }
  88. `;
  89. const support = `package com.fs.company.service.easycall;
  90. import com.alibaba.fastjson.JSON;
  91. import com.alibaba.fastjson.JSONObject;
  92. import com.fs.common.utils.StringUtils;
  93. import com.fs.company.vo.easycall.EasyCallCommonAddCallListParam;
  94. import com.fs.company.vo.easycall.EasyCallPhoneItemVO;
  95. import lombok.extern.slf4j.Slf4j;
  96. import org.springframework.beans.factory.annotation.Autowired;
  97. import org.springframework.stereotype.Component;
  98. import java.util.Date;
  99. import java.util.List;
  100. import java.util.Map;
  101. import java.util.UUID;
  102. import java.util.concurrent.ThreadLocalRandom;
  103. import java.util.concurrent.TimeUnit;
  104. /**
  105. * ${'\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'}
  106. */
  107. @Slf4j
  108. @Component
  109. public class OutboundLimitRetrySupport {
  110. /** ${'\u4E0E EasyCallServiceImpl \u4FDD\u6301\u4E00\u81F4'} */
  111. public static final String OUTBOUND_LIMIT_REDIS_PREFIX = "outbound:limit:";
  112. public static final String RETRY_ZSET_KEY = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:zset";
  113. public static final String RETRY_PAYLOAD_HASH_KEY = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:payload";
  114. private static final String RETRY_PROCESSING_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:processing:";
  115. private static final String RETRY_DISPATCHED_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:dispatched:";
  116. /** ${'\u5DF2\u6210\u529F append+startTask \u7684 callBackUuid \u6807\u8BB0 TTL\uFF08\u9632 pop \u540E\u5D29\u6E83\u5BFC\u81F4\u91CD\u590D\u5916\u547C\uFF09'} */
  117. private static final int DISPATCHED_TTL_HOURS = 24;
  118. private static final String RETRY_RATE_PREFIX = OUTBOUND_LIMIT_REDIS_PREFIX + "retry:rate:";
  119. /** ${'\u5230\u671F\u540E\u968F\u673A\u5EF6\u8FDF\u4E0A\u9650\uFF08\u6BEB\u79D2\uFF09\uFF0C\u5E73\u6ED1\u7A97\u53E3\u8FB9\u754C\u6D41\u91CF'} */
  120. private static final int RETRY_JITTER_MAX_MS = 30_000;
  121. /** ${'\u5B64\u513F payload \u5355\u6B21\u6700\u591A\u6062\u590D\u6761\u6570'} */
  122. private static final int MAX_ORPHAN_RECOVER_BATCH = 100;
  123. /** ${'\u5931\u8D25\u9000\u907F\u57FA\u7840\u5EF6\u8FDF\uFF08\u6BEB\u79D2\uFF09'} */
  124. private static final long RETRY_BACKOFF_MS = 60_000L;
  125. /** ${'\u5355\u7F51\u5173\u6BCF\u79D2\u6700\u591A\u91CD\u8BD5\u5916\u547C\u6B21\u6570\uFF08\u591A\u5B9E\u4F8B\u5171\u4EAB\uFF09'} */
  126. private static final int MAX_RETRY_PER_GATEWAY_PER_SEC = 30;
  127. /** ${'\u5904\u7406\u9501 TTL\uFF0C\u9632\u6B62\u5B9E\u4F8B\u5D29\u6E83\u540E\u957F\u671F\u5360\u9501'} */
  128. private static final int PROCESSING_LOCK_SECONDS = 600;
  129. @Autowired
  130. private OutboundLimitRetryRedisHelper retryRedisHelper;
  131. /**
  132. * ${'\u9650\u6D41\u540E\u5165\u961F\uFF08member=callBackUuid\uFF0C\u540C UUID \u4E0D\u91CD\u590D\u5165\u961F\uFF0Cscore \u5E26\u6296\u52A8\uFF09\u3002'}
  133. */
  134. public void enqueueRetry(Long companyId, Long gatewayId, EasyCallCommonAddCallListParam param, Date nextAvailableTime) {
  135. if (param == null || nextAvailableTime == null) {
  136. return;
  137. }
  138. try {
  139. String callBackUuid = extractCallBackUuid(param);
  140. String member = StringUtils.isNotBlank(callBackUuid) ? callBackUuid : UUID.randomUUID().toString();
  141. JSONObject retryData = new JSONObject();
  142. retryData.put("companyId", companyId);
  143. retryData.put("gatewayId", gatewayId);
  144. retryData.put("param", param);
  145. retryData.put("callBackUuid", member);
  146. retryData.put("nextAvailableTime", nextAvailableTime.getTime());
  147. retryData.put("createTime", System.currentTimeMillis());
  148. String existingPayload = retryRedisHelper.hashGet(RETRY_PAYLOAD_HASH_KEY, member);
  149. if (existingPayload != null) {
  150. JSONObject oldData = JSON.parseObject(existingPayload);
  151. if (oldData != null && oldData.containsKey("failCount")) {
  152. retryData.put("failCount", oldData.getIntValue("failCount"));
  153. }
  154. }
  155. double score = computeScoreWithJitter(nextAvailableTime.getTime());
  156. Double existingScore = retryRedisHelper.zScore(RETRY_ZSET_KEY, member);
  157. if (existingScore != null && existingScore <= score) {
  158. retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString());
  159. log.info("enqueueRetry: ${'\u5DF2\u5B58\u5728\u66F4\u65E9\u6216\u76F8\u540C\u8C03\u5EA6\uFF0C\u4EC5\u5237\u65B0 payload'} - member={}", member);
  160. return;
  161. }
  162. retryRedisHelper.zSetAdd(RETRY_ZSET_KEY, member, score);
  163. retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString());
  164. log.info("enqueueRetry: member={}, score={}, companyId={}, gatewayId={}", member, (long) score, companyId, gatewayId);
  165. } catch (Exception e) {
  166. log.error("enqueueRetry ${'\u5F02\u5E38'} companyId={}, gatewayId={}", companyId, gatewayId, e);
  167. }
  168. }
  169. /**
  170. * ${'\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'}
  171. */
  172. public List<String> popDueMembers(int batchSize) {
  173. long now = System.currentTimeMillis();
  174. return retryRedisHelper.popDueMembersByScore(RETRY_ZSET_KEY, now, batchSize);
  175. }
  176. /**
  177. * ${'\u89E3\u6790\u4EFB\u52A1 payload\uFF1B\u517C\u5BB9\u65E7\u7248 member \u4E3A\u6574\u6BB5 JSON \u7684\u5199\u6CD5\u3002'}
  178. */
  179. public JSONObject resolvePayload(String member) {
  180. if (StringUtils.isBlank(member)) {
  181. return null;
  182. }
  183. if (member.trim().startsWith("{")) {
  184. return JSON.parseObject(member);
  185. }
  186. String raw = retryRedisHelper.hashGet(RETRY_PAYLOAD_HASH_KEY, member);
  187. if (raw == null) {
  188. return null;
  189. }
  190. return JSON.parseObject(raw);
  191. }
  192. /**
  193. * ${'\u5C1D\u8BD5\u83B7\u53D6\u5904\u7406\u9501\uFF0C\u907F\u514D\u6781\u7AEF\u60C5\u51B5\u4E0B\u91CD\u590D\u5916\u547C\u3002'}
  194. */
  195. public boolean tryAcquireProcessingLock(String callBackUuid) {
  196. if (StringUtils.isBlank(callBackUuid)) {
  197. return true;
  198. }
  199. String lockKey = RETRY_PROCESSING_PREFIX + callBackUuid;
  200. return retryRedisHelper.setIfAbsent(lockKey, "1", PROCESSING_LOCK_SECONDS, TimeUnit.SECONDS);
  201. }
  202. public void releaseProcessingLock(String callBackUuid) {
  203. if (StringUtils.isBlank(callBackUuid)) {
  204. return;
  205. }
  206. retryRedisHelper.deleteKey(RETRY_PROCESSING_PREFIX + callBackUuid);
  207. }
  208. /**
  209. * ${'\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'}
  210. */
  211. public boolean isDispatched(String callBackUuid) {
  212. if (StringUtils.isBlank(callBackUuid)) {
  213. return false;
  214. }
  215. return retryRedisHelper.getString(RETRY_DISPATCHED_PREFIX + callBackUuid) != null;
  216. }
  217. /**
  218. * ${'\u6807\u8BB0 callBackUuid \u5DF2\u5B8C\u6210\u6D3E\u53D1\uFF08append \u6210\u529F\u4E14 startTask \u5DF2\u89E6\u53D1\uFF09\u3002'}
  219. */
  220. public void markDispatched(String callBackUuid) {
  221. if (StringUtils.isBlank(callBackUuid)) {
  222. return;
  223. }
  224. retryRedisHelper.setString(RETRY_DISPATCHED_PREFIX + callBackUuid, "1", DISPATCHED_TTL_HOURS, TimeUnit.HOURS);
  225. }
  226. /**
  227. * ${'\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'}
  228. *
  229. * @return ${'\u672C\u6B21\u6062\u590D\u5165\u961F\u6761\u6570'}
  230. */
  231. public int recoverOrphanedPayloads() {
  232. Map<String, String> all = retryRedisHelper.hashGetAll(RETRY_PAYLOAD_HASH_KEY);
  233. if (all == null || all.isEmpty()) {
  234. return 0;
  235. }
  236. int recovered = 0;
  237. for (Map.Entry<String, String> entry : all.entrySet()) {
  238. if (recovered >= MAX_ORPHAN_RECOVER_BATCH) {
  239. break;
  240. }
  241. String member = entry.getKey();
  242. if (retryRedisHelper.zScore(RETRY_ZSET_KEY, member) != null) {
  243. continue;
  244. }
  245. JSONObject data = JSON.parseObject(entry.getValue());
  246. if (data == null) {
  247. ackSuccess(member);
  248. continue;
  249. }
  250. String callBackUuid = data.getString("callBackUuid");
  251. EasyCallCommonAddCallListParam param = data.getObject("param", EasyCallCommonAddCallListParam.class);
  252. if (StringUtils.isBlank(callBackUuid) && param != null) {
  253. callBackUuid = extractCallBackUuid(param);
  254. }
  255. if (isDispatched(callBackUuid)) {
  256. ackSuccess(member);
  257. continue;
  258. }
  259. Long companyId = data.getLong("companyId");
  260. Long gatewayId = data.getLong("gatewayId");
  261. if (companyId == null || param == null) {
  262. ackSuccess(member);
  263. continue;
  264. }
  265. Date nextTime = new Date(System.currentTimeMillis() + RETRY_BACKOFF_MS);
  266. enqueueRetry(companyId, gatewayId, param, nextTime);
  267. recovered++;
  268. log.warn("recoverOrphanedPayloads: ${'\u5B64\u513F\u4EFB\u52A1\u91CD\u65B0\u5165\u961F'} member={}, callBackUuid={}", member, callBackUuid);
  269. }
  270. return recovered;
  271. }
  272. /**
  273. * ${'\u6309\u7F51\u5173\u4EE4\u724C\u6876\u9650\u901F\uFF0C\u8D85\u9650\u5219\u5EF6\u8FDF\u91CD\u5165\u961F\u3002'}
  274. *
  275. * @return true ${'\u8868\u793A\u5F53\u524D\u79D2\u989D\u5EA6\u5DF2\u7528\u5C3D\uFF0C\u8C03\u7528\u65B9\u5E94\u8DF3\u8FC7\u672C\u6B21\u5916\u547C'}
  276. */
  277. public boolean isGatewayRateLimited(Long companyId, Long gatewayId) {
  278. if (companyId == null || gatewayId == null) {
  279. return false;
  280. }
  281. String bucketKey = RETRY_RATE_PREFIX + companyId + ":" + gatewayId;
  282. Long count = retryRedisHelper.incr(bucketKey, 1L);
  283. if (count == null) {
  284. count = 1L;
  285. }
  286. if (count == 1L) {
  287. retryRedisHelper.expire(bucketKey, 1, TimeUnit.SECONDS);
  288. }
  289. return count > MAX_RETRY_PER_GATEWAY_PER_SEC;
  290. }
  291. /**
  292. * ${'\u6210\u529F\u540E\u6E05\u7406 payload\uFF1Bmember \u5DF2\u5728 pop \u65F6\u4ECE ZSET \u79FB\u9664\u3002'}
  293. */
  294. public void ackSuccess(String member) {
  295. if (StringUtils.isBlank(member) || member.trim().startsWith("{")) {
  296. return;
  297. }
  298. retryRedisHelper.hashDelete(RETRY_PAYLOAD_HASH_KEY, member);
  299. }
  300. /**
  301. * ${'\u5931\u8D25\u6216\u5F02\u5E38\uFF1A\u9000\u907F\u540E\u91CD\u65B0\u5165\u961F\uFF0C\u4E0D\u4E22\u5355\u3002'}
  302. */
  303. public void requeueWithBackoff(JSONObject retryData, String member) {
  304. if (retryData == null) {
  305. return;
  306. }
  307. Long companyId = retryData.getLong("companyId");
  308. Long gatewayId = retryData.getLong("gatewayId");
  309. EasyCallCommonAddCallListParam param = retryData.getObject("param", EasyCallCommonAddCallListParam.class);
  310. if (param == null) {
  311. log.warn("requeueWithBackoff: param ${'\u4E3A\u7A7A\uFF0C\u65E0\u6CD5\u91CD\u5165\u961F'} member={}", member);
  312. return;
  313. }
  314. int failCount = retryData.getIntValue("failCount") + 1;
  315. retryData.put("failCount", failCount);
  316. if (StringUtils.isNotBlank(member) && !member.trim().startsWith("{")) {
  317. retryRedisHelper.hashPut(RETRY_PAYLOAD_HASH_KEY, member, retryData.toJSONString());
  318. }
  319. long delay = RETRY_BACKOFF_MS * Math.min(failCount, 5);
  320. Date nextTime = new Date(System.currentTimeMillis() + delay);
  321. enqueueRetry(companyId, gatewayId, param, nextTime);
  322. log.warn("requeueWithBackoff: member={}, failCount={}, nextDelayMs={}", member, failCount, delay);
  323. }
  324. public static String extractCallBackUuid(EasyCallCommonAddCallListParam param) {
  325. if (param == null || param.getPhoneList() == null || param.getPhoneList().isEmpty()) {
  326. return null;
  327. }
  328. EasyCallPhoneItemVO item = param.getPhoneList().get(0);
  329. if (item == null || item.getBizJson() == null) {
  330. return null;
  331. }
  332. JSONObject bizJson;
  333. if (item.getBizJson() instanceof JSONObject) {
  334. bizJson = (JSONObject) item.getBizJson();
  335. } else {
  336. bizJson = JSON.parseObject(String.valueOf(item.getBizJson()));
  337. }
  338. return bizJson != null ? bizJson.getString("callBackUuid") : null;
  339. }
  340. private static double computeScoreWithJitter(long baseTimeMs) {
  341. int jitter = ThreadLocalRandom.current().nextInt(RETRY_JITTER_MAX_MS + 1);
  342. return baseTimeMs + jitter;
  343. }
  344. }
  345. `;
  346. writeFileSync(path.join(base, 'OutboundLimitRetryRedisHelper.java'), helper, 'utf8');
  347. writeFileSync(path.join(base, 'OutboundLimitRetrySupport.java'), support, 'utf8');
  348. console.log('UTF-8 files written');