feat(backend): 自动赎回 Builder Relayer 限流与配额处理

- RelayClientService: 429 限流时指数退避重试;解析 quota exceeded/resets in N seconds 并记录冷却时间;暴露 isBuilderRelayerQuotaBlocked 与 getBuilderRelayerQuotaBlockedRemainingSeconds
- PositionCheckService: checkRedeemablePositions 防重入;配额冷却期内跳过自动赎回

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
WrBug
2026-02-16 04:32:52 +08:00
co-authored by Cursor
parent 6b56f62532
commit 24bb7bed40
2 changed files with 90 additions and 11 deletions
@@ -25,6 +25,7 @@ import com.wrbug.polymarketbot.service.common.MarketPriceService
import org.springframework.stereotype.Service import org.springframework.stereotype.Service
import java.math.BigDecimal import java.math.BigDecimal
import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.AtomicBoolean
/** /**
* 仓位检查服务 * 仓位检查服务
@@ -78,6 +79,9 @@ class PositionCheckService(
// 同步锁,确保订阅任务的启动和停止是线程安全的 // 同步锁,确保订阅任务的启动和停止是线程安全的
private val lock = Any() private val lock = Any()
// 防止 checkRedeemablePositions 重入:上一轮检查未完成时,新一轮轮询直接跳过
private val redeemCheckInProgress = AtomicBoolean(false)
/** /**
* 初始化服务(订阅 PositionPollingService 的事件,启动缓存清理任务) * 初始化服务(订阅 PositionPollingService 的事件,启动缓存清理任务)
*/ */
@@ -328,12 +332,17 @@ class PositionCheckService(
/** /**
* 逻辑1:处理待赎回仓位 * 逻辑1:处理待赎回仓位
https://clob.polymarket.com * 按照以下逻辑处理: * 按照以下逻辑处理:
* 1. 无待赎回仓位:跳过 * 1. 无待赎回仓位:跳过
* 2. (未配置apikey || autoredeem==false) && 有待赎回的仓位:发送通知事件 * 2. (未配置apikey || autoredeem==false) && 有待赎回的仓位:发送通知事件
* 3. (已配置) && 有待赎回的仓位:处理订单逻辑 * 3. (已配置) && 有待赎回的仓位:处理订单逻辑
* 防重入:上一轮检查未完成时,本轮直接跳过,避免并发赎回。
*/ */
private suspend fun checkRedeemablePositions(redeemablePositions: List<AccountPositionDto>) { private suspend fun checkRedeemablePositions(redeemablePositions: List<AccountPositionDto>) {
if (!redeemCheckInProgress.compareAndSet(false, true)) {
logger.debug("跳过本次待赎回仓位检查:上一次检查尚未完成")
return
}
try { try {
// 1. 无待赎回仓位:跳过 // 1. 无待赎回仓位:跳过
if (redeemablePositions.isEmpty()) { if (redeemablePositions.isEmpty()) {
@@ -374,6 +383,13 @@ class PositionCheckService(
return // 未配置时直接返回,不进行后续处理 return // 未配置时直接返回,不进行后续处理
} }
// Builder Relayer 配额冷却期内不再发起赎回(如 API 返回 quota exceeded, resets in N seconds
if (relayClientService.isBuilderRelayerQuotaBlocked()) {
val remaining = relayClientService.getBuilderRelayerQuotaBlockedRemainingSeconds()
logger.info("Builder Relayer 配额冷却中,跳过本次自动赎回,约 ${remaining} 秒后恢复")
return
}
// 3. (已配置) && 有待赎回的仓位:处理订单逻辑 // 3. (已配置) && 有待赎回的仓位:处理订单逻辑
// 自动赎回已开启且已配置 API Key,按账户分组进行赎回处理 // 自动赎回已开启且已配置 API Key,按账户分组进行赎回处理
// 先执行赎回,赎回成功后再查找订单并更新订单状态 // 先执行赎回,赎回成功后再查找订单并更新订单状态
@@ -451,6 +467,8 @@ class PositionCheckService(
} }
} catch (e: Exception) { } catch (e: Exception) {
logger.error("处理待赎回仓位异常: ${e.message}", e) logger.error("处理待赎回仓位异常: ${e.message}", e)
} finally {
redeemCheckInProgress.set(false)
} }
} }
@@ -8,9 +8,12 @@ import com.wrbug.polymarketbot.enums.WalletType
import com.wrbug.polymarketbot.util.EthereumUtils import com.wrbug.polymarketbot.util.EthereumUtils
import com.wrbug.polymarketbot.util.RetrofitFactory import com.wrbug.polymarketbot.util.RetrofitFactory
import com.wrbug.polymarketbot.util.createClient import com.wrbug.polymarketbot.util.createClient
import kotlinx.coroutines.delay
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.stereotype.Service import org.springframework.stereotype.Service
import retrofit2.Response
import java.math.BigInteger import java.math.BigInteger
import java.util.concurrent.atomic.AtomicLong
/** /**
* RelayClient 服务 * RelayClient 服务
@@ -61,6 +64,59 @@ class RelayClientService(
retrofitFactory.createEthereumRpcApi(rpcUrl) 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 状态为 429Too Many Requests,如 Cloudflare 1015)时等待后重试,避免瞬时限流导致赎回失败。
*/
private suspend fun <T> withBuilderRelayerRateLimitRetry(block: suspend () -> Response<T>): Response<T> {
var lastResponse: Response<T>? = 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 客户端(动态获取,因为配置可能更新) * 获取 Builder Relayer API 客户端(动态获取,因为配置可能更新)
*/ */
@@ -131,6 +187,7 @@ class RelayClientService(
Result.success(responseTime) Result.success(responseTime)
} else { } else {
val errorBody = response.errorBody()?.string() ?: "未知错误" val errorBody = response.errorBody()?.string() ?: "未知错误"
updateQuotaBlockedFromErrorBody(errorBody)
Result.failure(Exception("Builder Relayer API 调用失败: ${response.code()} - $errorBody")) Result.failure(Exception("Builder Relayer API 调用失败: ${response.code()} - $errorBody"))
} }
} catch (e: Exception) { } catch (e: Exception) {
@@ -395,9 +452,10 @@ class RelayClientService(
val credentials = org.web3j.crypto.Credentials.create(privateKeyBigInt.toString(16)) val credentials = org.web3j.crypto.Credentials.create(privateKeyBigInt.toString(16))
val fromAddress = credentials.address 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) { if (!relayPayloadResponse.isSuccessful || relayPayloadResponse.body() == null) {
val errorBody = relayPayloadResponse.errorBody()?.string() ?: "未知错误" val errorBody = relayPayloadResponse.errorBody()?.string() ?: "未知错误"
updateQuotaBlockedFromErrorBody(errorBody)
logger.error("获取 Relay Payload 失败: code=${relayPayloadResponse.code()}, body=$errorBody") logger.error("获取 Relay Payload 失败: code=${relayPayloadResponse.code()}, body=$errorBody")
return Result.failure(Exception("获取 Relay Payload 失败: ${relayPayloadResponse.code()} - $errorBody")) return Result.failure(Exception("获取 Relay Payload 失败: ${relayPayloadResponse.code()} - $errorBody"))
} }
@@ -461,9 +519,10 @@ class RelayClientService(
metadata = "Redeem positions via Builder Relayer PROXY" metadata = "Redeem positions via Builder Relayer PROXY"
) )
val response = relayerApi.submitTransaction(request) val response = withBuilderRelayerRateLimitRetry { relayerApi.submitTransaction(request) }
if (!response.isSuccessful || response.body() == null) { if (!response.isSuccessful || response.body() == null) {
val errorBody = response.errorBody()?.string() ?: "未知错误" val errorBody = response.errorBody()?.string() ?: "未知错误"
updateQuotaBlockedFromErrorBody(errorBody)
logger.error("Builder Relayer PROXY API 调用失败: code=${response.code()}, body=$errorBody") logger.error("Builder Relayer PROXY API 调用失败: code=${response.code()}, body=$errorBody")
return Result.failure(Exception("Builder Relayer PROXY 调用失败: ${response.code()} - $errorBody")) return Result.failure(Exception("Builder Relayer PROXY 调用失败: ${response.code()} - $errorBody"))
} }
@@ -625,10 +684,11 @@ class RelayClientService(
// safeTx.data 已经是带 0x 前缀的完整调用数据 // safeTx.data 已经是带 0x 前缀的完整调用数据
val redeemCallData = safeTx.data val redeemCallData = safeTx.data
// 获取 Proxy 的 nonce(通过 Builder Relayer API // 获取 Proxy 的 nonce(通过 Builder Relayer API,遇 429 限流时重试
val nonceResponse = relayerApi.getNonce(fromAddress, RELAYER_TYPE_SAFE) val nonceResponse = withBuilderRelayerRateLimitRetry { relayerApi.getNonce(fromAddress, RELAYER_TYPE_SAFE) }
if (!nonceResponse.isSuccessful || nonceResponse.body() == null) { if (!nonceResponse.isSuccessful || nonceResponse.body() == null) {
val errorBody = nonceResponse.errorBody()?.string() ?: "未知错误" val errorBody = nonceResponse.errorBody()?.string() ?: "未知错误"
updateQuotaBlockedFromErrorBody(errorBody)
logger.error("获取 nonce 失败: code=${nonceResponse.code()}, body=$errorBody") logger.error("获取 nonce 失败: code=${nonceResponse.code()}, body=$errorBody")
return Result.failure(Exception("获取 nonce 失败: ${nonceResponse.code()} - $errorBody")) return Result.failure(Exception("获取 nonce 失败: ${nonceResponse.code()} - $errorBody"))
} }
@@ -710,11 +770,12 @@ class RelayClientService(
} }
) )
// 调用 Builder Relayer API(认证头通过拦截器添加) // 调用 Builder Relayer API(认证头通过拦截器添加,遇 429 限流时重试
val response = relayerApi.submitTransaction(request) val response = withBuilderRelayerRateLimitRetry { relayerApi.submitTransaction(request) }
if (!response.isSuccessful || response.body() == null) { if (!response.isSuccessful || response.body() == null) {
val errorBody = response.errorBody()?.string() ?: "未知错误" val errorBody = response.errorBody()?.string() ?: "未知错误"
updateQuotaBlockedFromErrorBody(errorBody)
logger.error("Builder Relayer API 调用失败: code=${response.code()}, body=$errorBody") logger.error("Builder Relayer API 调用失败: code=${response.code()}, body=$errorBody")
return Result.failure(Exception("Builder Relayer API 调用失败: ${response.code()} - $errorBody")) return Result.failure(Exception("Builder Relayer API 调用失败: ${response.code()} - $errorBody"))
} }