From 24bb7bed40e26d586daf6b14465d7244ce86b954 Mon Sep 17 00:00:00 2001 From: WrBug Date: Mon, 16 Feb 2026 04:32:52 +0800 Subject: [PATCH] =?UTF-8?q?feat(backend):=20=E8=87=AA=E5=8A=A8=E8=B5=8E?= =?UTF-8?q?=E5=9B=9E=20Builder=20Relayer=20=E9=99=90=E6=B5=81=E4=B8=8E?= =?UTF-8?q?=E9=85=8D=E9=A2=9D=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - RelayClientService: 429 限流时指数退避重试;解析 quota exceeded/resets in N seconds 并记录冷却时间;暴露 isBuilderRelayerQuotaBlocked 与 getBuilderRelayerQuotaBlockedRemainingSeconds - PositionCheckService: checkRedeemablePositions 防重入;配额冷却期内跳过自动赎回 Co-authored-by: Cursor --- .../service/accounts/PositionCheckService.kt | 28 +++++-- .../service/system/RelayClientService.kt | 73 +++++++++++++++++-- 2 files changed, 90 insertions(+), 11 deletions(-) diff --git a/backend/src/main/kotlin/com/wrbug/polymarketbot/service/accounts/PositionCheckService.kt b/backend/src/main/kotlin/com/wrbug/polymarketbot/service/accounts/PositionCheckService.kt index f49ff74..6217462 100644 --- a/backend/src/main/kotlin/com/wrbug/polymarketbot/service/accounts/PositionCheckService.kt +++ b/backend/src/main/kotlin/com/wrbug/polymarketbot/service/accounts/PositionCheckService.kt @@ -25,6 +25,7 @@ import com.wrbug.polymarketbot.service.common.MarketPriceService import org.springframework.stereotype.Service import java.math.BigDecimal import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.atomic.AtomicBoolean /** * 仓位检查服务 @@ -77,7 +78,10 @@ class PositionCheckService( // 同步锁,确保订阅任务的启动和停止是线程安全的 private val lock = Any() - + + // 防止 checkRedeemablePositions 重入:上一轮检查未完成时,新一轮轮询直接跳过 + private val redeemCheckInProgress = AtomicBoolean(false) + /** * 初始化服务(订阅 PositionPollingService 的事件,启动缓存清理任务) */ @@ -328,18 +332,23 @@ class PositionCheckService( /** * 逻辑1:处理待赎回仓位 - https://clob.polymarket.com * 按照以下逻辑处理: + * 按照以下逻辑处理: * 1. 无待赎回仓位:跳过 * 2. (未配置apikey || autoredeem==false) && 有待赎回的仓位:发送通知事件 * 3. (已配置) && 有待赎回的仓位:处理订单逻辑 + * 防重入:上一轮检查未完成时,本轮直接跳过,避免并发赎回。 */ private suspend fun checkRedeemablePositions(redeemablePositions: List) { + if (!redeemCheckInProgress.compareAndSet(false, true)) { + logger.debug("跳过本次待赎回仓位检查:上一次检查尚未完成") + return + } try { // 1. 无待赎回仓位:跳过 if (redeemablePositions.isEmpty()) { return } - + // 检查系统级别的自动赎回配置 val autoRedeemEnabled = systemConfigService.isAutoRedeemEnabled() val apiKeyConfigured = relayClientService.isBuilderApiKeyConfigured() @@ -373,7 +382,14 @@ class PositionCheckService( } return // 未配置时直接返回,不进行后续处理 } - + + // Builder Relayer 配额冷却期内不再发起赎回(如 API 返回 quota exceeded, resets in N seconds) + if (relayClientService.isBuilderRelayerQuotaBlocked()) { + val remaining = relayClientService.getBuilderRelayerQuotaBlockedRemainingSeconds() + logger.info("Builder Relayer 配额冷却中,跳过本次自动赎回,约 ${remaining} 秒后恢复") + return + } + // 3. (已配置) && 有待赎回的仓位:处理订单逻辑 // 自动赎回已开启且已配置 API Key,按账户分组进行赎回处理 // 先执行赎回,赎回成功后再查找订单并更新订单状态 @@ -451,9 +467,11 @@ class PositionCheckService( } } catch (e: Exception) { logger.error("处理待赎回仓位异常: ${e.message}", e) + } finally { + redeemCheckInProgress.set(false) } } - + /** * 逻辑2:处理未卖出订单 * 检查所有未卖出的订单,匹配仓位 diff --git a/backend/src/main/kotlin/com/wrbug/polymarketbot/service/system/RelayClientService.kt b/backend/src/main/kotlin/com/wrbug/polymarketbot/service/system/RelayClientService.kt index e06df54..7cae48f 100644 --- a/backend/src/main/kotlin/com/wrbug/polymarketbot/service/system/RelayClientService.kt +++ b/backend/src/main/kotlin/com/wrbug/polymarketbot/service/system/RelayClientService.kt @@ -8,9 +8,12 @@ import com.wrbug.polymarketbot.enums.WalletType import com.wrbug.polymarketbot.util.EthereumUtils import com.wrbug.polymarketbot.util.RetrofitFactory import com.wrbug.polymarketbot.util.createClient +import kotlinx.coroutines.delay import org.slf4j.LoggerFactory import org.springframework.stereotype.Service +import retrofit2.Response import java.math.BigInteger +import java.util.concurrent.atomic.AtomicLong /** * RelayClient 服务 @@ -61,6 +64,59 @@ class RelayClientService( retrofitFactory.createEthereumRpcApi(rpcUrl) } + /** 遇到 429 限流时的重试次数 */ + private val builderRelayerRateLimitMaxAttempts = 3 + + /** 429 限流重试退避基数(毫秒),第 n 次重试等待 baseMs * 2^(n-1) */ + private val builderRelayerRateLimitBackoffMs = 2000L + + /** Builder Relayer 配额用尽后的冷却截止时间(毫秒时间戳),在此时间前不再发起赎回 */ + private val builderRelayerQuotaBlockedUntilMs = AtomicLong(0) + + /** + * 是否处于 Builder Relayer 配额冷却期(配额用尽后在该时间内不再发起赎回)。 + */ + fun isBuilderRelayerQuotaBlocked(): Boolean = System.currentTimeMillis() < builderRelayerQuotaBlockedUntilMs.get() + + /** + * 配额冷却剩余秒数,未在冷却期时返回 0。 + */ + fun getBuilderRelayerQuotaBlockedRemainingSeconds(): Long { + val remaining = (builderRelayerQuotaBlockedUntilMs.get() - System.currentTimeMillis()) / 1000 + return maxOf(0, remaining) + } + + /** + * 从 API 错误响应中解析 "quota exceeded... resets in N seconds",并设置配额冷却截止时间。 + */ + private fun updateQuotaBlockedFromErrorBody(errorBody: String) { + if (!errorBody.contains("quota exceeded", ignoreCase = true)) return + val regex = Regex("resets\\s+in\\s+(\\d+)\\s+seconds", RegexOption.IGNORE_CASE) + regex.find(errorBody)?.groupValues?.getOrNull(1)?.toLongOrNull()?.let { seconds -> + val untilMs = System.currentTimeMillis() + seconds * 1000 + builderRelayerQuotaBlockedUntilMs.set(untilMs) + logger.warn("Builder Relayer 配额已用尽,${seconds}秒内不再发起赎回") + } + } + + /** + * 对 Builder Relayer API 调用进行 429 限流重试(指数退避)。 + * 当 HTTP 状态为 429(Too Many Requests,如 Cloudflare 1015)时等待后重试,避免瞬时限流导致赎回失败。 + */ + private suspend fun withBuilderRelayerRateLimitRetry(block: suspend () -> Response): Response { + var lastResponse: Response? = null + for (attempt in 1..builderRelayerRateLimitMaxAttempts) { + val response = block() + lastResponse = response + if (response.code() != 429) return response + if (attempt == builderRelayerRateLimitMaxAttempts) return response + val delayMs = builderRelayerRateLimitBackoffMs * (1L shl (attempt - 1)) + logger.warn("Builder Relayer API 限流(429),${delayMs}ms 后重试 (${attempt}/${builderRelayerRateLimitMaxAttempts})") + delay(delayMs) + } + return lastResponse!! + } + /** * 获取 Builder Relayer API 客户端(动态获取,因为配置可能更新) */ @@ -131,6 +187,7 @@ class RelayClientService( Result.success(responseTime) } else { val errorBody = response.errorBody()?.string() ?: "未知错误" + updateQuotaBlockedFromErrorBody(errorBody) Result.failure(Exception("Builder Relayer API 调用失败: ${response.code()} - $errorBody")) } } catch (e: Exception) { @@ -395,9 +452,10 @@ class RelayClientService( val credentials = org.web3j.crypto.Credentials.create(privateKeyBigInt.toString(16)) val fromAddress = credentials.address - val relayPayloadResponse = relayerApi.getRelayPayload(fromAddress, RELAYER_TYPE_PROXY) + val relayPayloadResponse = withBuilderRelayerRateLimitRetry { relayerApi.getRelayPayload(fromAddress, RELAYER_TYPE_PROXY) } if (!relayPayloadResponse.isSuccessful || relayPayloadResponse.body() == null) { val errorBody = relayPayloadResponse.errorBody()?.string() ?: "未知错误" + updateQuotaBlockedFromErrorBody(errorBody) logger.error("获取 Relay Payload 失败: code=${relayPayloadResponse.code()}, body=$errorBody") return Result.failure(Exception("获取 Relay Payload 失败: ${relayPayloadResponse.code()} - $errorBody")) } @@ -461,9 +519,10 @@ class RelayClientService( metadata = "Redeem positions via Builder Relayer PROXY" ) - val response = relayerApi.submitTransaction(request) + val response = withBuilderRelayerRateLimitRetry { relayerApi.submitTransaction(request) } if (!response.isSuccessful || response.body() == null) { val errorBody = response.errorBody()?.string() ?: "未知错误" + updateQuotaBlockedFromErrorBody(errorBody) logger.error("Builder Relayer PROXY API 调用失败: code=${response.code()}, body=$errorBody") return Result.failure(Exception("Builder Relayer PROXY 调用失败: ${response.code()} - $errorBody")) } @@ -625,10 +684,11 @@ class RelayClientService( // safeTx.data 已经是带 0x 前缀的完整调用数据 val redeemCallData = safeTx.data - // 获取 Proxy 的 nonce(通过 Builder Relayer API) - val nonceResponse = relayerApi.getNonce(fromAddress, RELAYER_TYPE_SAFE) + // 获取 Proxy 的 nonce(通过 Builder Relayer API,遇 429 限流时重试) + val nonceResponse = withBuilderRelayerRateLimitRetry { relayerApi.getNonce(fromAddress, RELAYER_TYPE_SAFE) } if (!nonceResponse.isSuccessful || nonceResponse.body() == null) { val errorBody = nonceResponse.errorBody()?.string() ?: "未知错误" + updateQuotaBlockedFromErrorBody(errorBody) logger.error("获取 nonce 失败: code=${nonceResponse.code()}, body=$errorBody") return Result.failure(Exception("获取 nonce 失败: ${nonceResponse.code()} - $errorBody")) } @@ -710,11 +770,12 @@ class RelayClientService( } ) - // 调用 Builder Relayer API(认证头通过拦截器添加) - val response = relayerApi.submitTransaction(request) + // 调用 Builder Relayer API(认证头通过拦截器添加,遇 429 限流时重试) + val response = withBuilderRelayerRateLimitRetry { relayerApi.submitTransaction(request) } if (!response.isSuccessful || response.body() == null) { val errorBody = response.errorBody()?.string() ?: "未知错误" + updateQuotaBlockedFromErrorBody(errorBody) logger.error("Builder Relayer API 调用失败: code=${response.code()}, body=$errorBody") return Result.failure(Exception("Builder Relayer API 调用失败: ${response.code()} - $errorBody")) }