fix: 修复最大仓位金额未生效和交易receipt解析问题

1. 修复 CopyTradingFilterService 中最大仓位金额(maxPositionValue)未生效的问题
   - 当只配置了最大仓位金额但未配置 maxSpread 或 minOrderDepth 时,
     仓位检查被跳过导致过滤失效
   - 将仓位检查移到 needOrderbook 判断之前,确保始终执行仓位检查

2. 修复 OnChainWsService 中交易 receipt 为 JsonNull 时的空指针问题
   - 添加 JsonNull 检查,防止解析失败时出现异常
This commit is contained in:
WrBug
2026-01-14 05:04:31 +08:00
parent 92b75d1926
commit b0aa7c0862
2 changed files with 44 additions and 30 deletions
@@ -71,12 +71,20 @@ class CopyTradingFilterService(
} }
} }
// 3. 检查是否需要获取订单簿 // 3. 检查是否需要获取订单簿或需要执行仓位检查
// 只有在配置了需要订单簿的过滤条件时才获取 // 只有在配置了需要订单簿的过滤条件时才获取订单簿
val needOrderbook = copyTrading.maxSpread != null || copyTrading.minOrderDepth != null val needOrderbook = copyTrading.maxSpread != null || copyTrading.minOrderDepth != null
// 3.5. 如果不需要订单簿,则跳过订单簿相关的检查,但仍然需要检查仓位限制
if (!needOrderbook) { if (!needOrderbook) {
// 不需要订单簿,直接通过 // 仓位检查(如果配置了最大仓位限制且提供了跟单金额和市场ID)
if (copyOrderAmount != null && marketId != null) {
val positionCheck = checkPositionLimits(copyTrading, copyOrderAmount, marketId)
if (!positionCheck.isPassed) {
return positionCheck
}
}
// 通过所有检查
return FilterResult.passed() return FilterResult.passed()
} }
@@ -1,5 +1,6 @@
package com.wrbug.polymarketbot.service.copytrading.monitor package com.wrbug.polymarketbot.service.copytrading.monitor
import com.google.gson.JsonNull
import com.wrbug.polymarketbot.api.* import com.wrbug.polymarketbot.api.*
import com.wrbug.polymarketbot.entity.Leader import com.wrbug.polymarketbot.entity.Leader
import com.wrbug.polymarketbot.repository.LeaderRepository import com.wrbug.polymarketbot.repository.LeaderRepository
@@ -23,12 +24,12 @@ class OnChainWsService(
private val copyOrderTrackingService: CopyOrderTrackingService, private val copyOrderTrackingService: CopyOrderTrackingService,
private val leaderRepository: LeaderRepository private val leaderRepository: LeaderRepository
) { ) {
private val logger = LoggerFactory.getLogger(OnChainWsService::class.java) private val logger = LoggerFactory.getLogger(OnChainWsService::class.java)
// 存储需要监听的LeaderleaderId -> Leader // 存储需要监听的LeaderleaderId -> Leader
private val monitoredLeaders = ConcurrentHashMap<Long, Leader>() private val monitoredLeaders = ConcurrentHashMap<Long, Leader>()
/** /**
* 启动链上 WebSocket 监听 * 启动链上 WebSocket 监听
* 通过统一服务订阅所有 Leader * 通过统一服务订阅所有 Leader
@@ -40,14 +41,14 @@ class OnChainWsService(
stop() stop()
return return
} }
// 更新 Leader 列表 // 更新 Leader 列表
monitoredLeaders.clear() monitoredLeaders.clear()
leaders.forEach { leader -> leaders.forEach { leader ->
addLeader(leader) addLeader(leader)
} }
} }
/** /**
* 添加Leader监听 * 添加Leader监听
* 通过统一服务订阅该 Leader 的地址 * 通过统一服务订阅该 Leader 的地址
@@ -57,17 +58,17 @@ class OnChainWsService(
logger.warn("Leader ID为空,跳过: ${leader.leaderAddress}") logger.warn("Leader ID为空,跳过: ${leader.leaderAddress}")
return return
} }
val leaderId = leader.id!! val leaderId = leader.id!!
// 如果已经在监听列表中,不重复添加 // 如果已经在监听列表中,不重复添加
if (monitoredLeaders.containsKey(leaderId)) { if (monitoredLeaders.containsKey(leaderId)) {
logger.debug("Leader 已在监听列表中: ${leader.leaderName} (${leader.leaderAddress})") logger.debug("Leader 已在监听列表中: ${leader.leaderName} (${leader.leaderAddress})")
return return
} }
monitoredLeaders[leaderId] = leader monitoredLeaders[leaderId] = leader
// 通过统一服务订阅 // 通过统一服务订阅
val subscriptionId = "LEADER_$leaderId" val subscriptionId = "LEADER_$leaderId"
unifiedOnChainWsService.subscribe( unifiedOnChainWsService.subscribe(
@@ -79,40 +80,45 @@ class OnChainWsService(
handleLeaderTransaction(leaderId, txHash, httpClient, rpcApi) handleLeaderTransaction(leaderId, txHash, httpClient, rpcApi)
} }
) )
logger.info("添加 Leader 监听: ${leader.leaderName} (${leader.leaderAddress})") logger.info("添加 Leader 监听: ${leader.leaderName} (${leader.leaderAddress})")
} }
/** /**
* 处理 Leader 的交易 * 处理 Leader 的交易
*/ */
private suspend fun handleLeaderTransaction(leaderId: Long, txHash: String, httpClient: OkHttpClient, rpcApi: EthereumRpcApi) { private suspend fun handleLeaderTransaction(
leaderId: Long,
txHash: String,
httpClient: OkHttpClient,
rpcApi: EthereumRpcApi
) {
val leader = monitoredLeaders[leaderId] ?: return val leader = monitoredLeaders[leaderId] ?: return
logger.debug("开始处理 Leader 交易: leaderId=$leaderId, txHash=$txHash, leaderAddress=${leader.leaderAddress}") logger.debug("开始处理 Leader 交易: leaderId=$leaderId, txHash=$txHash, leaderAddress=${leader.leaderAddress}")
try { try {
// 获取交易 receipt // 获取交易 receipt
val receiptRequest = JsonRpcRequest( val receiptRequest = JsonRpcRequest(
method = "eth_getTransactionReceipt", method = "eth_getTransactionReceipt",
params = listOf(txHash) params = listOf(txHash)
) )
val receiptResponse = rpcApi.call(receiptRequest) val receiptResponse = rpcApi.call(receiptRequest)
if (!receiptResponse.isSuccessful || receiptResponse.body() == null) { if (!receiptResponse.isSuccessful || receiptResponse.body() == null) {
logger.warn("获取交易 receipt 失败: leaderId=$leaderId, txHash=$txHash, code=${receiptResponse.code()}") logger.warn("获取交易 receipt 失败: leaderId=$leaderId, txHash=$txHash, code=${receiptResponse.code()}")
return return
} }
val receiptRpcResponse = receiptResponse.body()!! val receiptRpcResponse = receiptResponse.body()!!
if (receiptRpcResponse.error != null || receiptRpcResponse.result == null) { if (receiptRpcResponse.error != null || receiptRpcResponse.result == null || receiptRpcResponse.result is JsonNull) {
logger.warn("交易 receipt 错误: leaderId=$leaderId, txHash=$txHash, error=${receiptRpcResponse.error}") logger.warn("交易 receipt 错误: leaderId=$leaderId, txHash=$txHash, error=${receiptRpcResponse.error}")
return return
} }
// 使用 Gson 解析 receipt JSON // 使用 Gson 解析 receipt JSON
val receiptJson = receiptRpcResponse.result.asJsonObject val receiptJson = receiptRpcResponse.result.asJsonObject
// 获取区块号和时间戳 // 获取区块号和时间戳
val blockNumber = receiptJson.get("blockNumber")?.asString val blockNumber = receiptJson.get("blockNumber")?.asString
val blockTimestamp = if (blockNumber != null) { val blockTimestamp = if (blockNumber != null) {
@@ -120,7 +126,7 @@ class OnChainWsService(
} else { } else {
null null
} }
// 解析 receipt 中的 Transfer 日志 // 解析 receipt 中的 Transfer 日志
val logs = receiptJson.getAsJsonArray("logs") ?: run { val logs = receiptJson.getAsJsonArray("logs") ?: run {
logger.warn("交易 receipt 中没有日志: leaderId=$leaderId, txHash=$txHash") logger.warn("交易 receipt 中没有日志: leaderId=$leaderId, txHash=$txHash")
@@ -128,7 +134,7 @@ class OnChainWsService(
} }
val (erc20Transfers, erc1155Transfers) = OnChainWsUtils.parseReceiptTransfers(logs) val (erc20Transfers, erc1155Transfers) = OnChainWsUtils.parseReceiptTransfers(logs)
logger.debug("解析交易日志: leaderId=$leaderId, txHash=$txHash, erc20Transfers=${erc20Transfers.size}, erc1155Transfers=${erc1155Transfers.size}") logger.debug("解析交易日志: leaderId=$leaderId, txHash=$txHash, erc20Transfers=${erc20Transfers.size}, erc1155Transfers=${erc1155Transfers.size}")
// 解析交易信息 // 解析交易信息
val trade = OnChainWsUtils.parseTradeFromTransfers( val trade = OnChainWsUtils.parseTradeFromTransfers(
txHash = txHash, txHash = txHash,
@@ -138,7 +144,7 @@ class OnChainWsService(
erc1155Transfers = erc1155Transfers, erc1155Transfers = erc1155Transfers,
retrofitFactory = retrofitFactory retrofitFactory = retrofitFactory
) )
if (trade != null) { if (trade != null) {
logger.info("成功解析交易: leaderId=$leaderId, txHash=$txHash, side=${trade.side}, market=${trade.market}, size=${trade.size}") logger.info("成功解析交易: leaderId=$leaderId, txHash=$txHash, side=${trade.side}, market=${trade.market}, size=${trade.size}")
// 调用 processTrade 处理交易 // 调用 processTrade 处理交易
@@ -154,21 +160,21 @@ class OnChainWsService(
logger.error("处理 Leader 交易失败: leaderId=$leaderId, txHash=$txHash, ${e.message}", e) logger.error("处理 Leader 交易失败: leaderId=$leaderId, txHash=$txHash, ${e.message}", e)
} }
} }
/** /**
* 移除Leader监听 * 移除Leader监听
* 取消该 Leader 的订阅 * 取消该 Leader 的订阅
*/ */
fun removeLeader(leaderId: Long) { fun removeLeader(leaderId: Long) {
monitoredLeaders.remove(leaderId) monitoredLeaders.remove(leaderId)
// 通过统一服务取消订阅 // 通过统一服务取消订阅
val subscriptionId = "LEADER_$leaderId" val subscriptionId = "LEADER_$leaderId"
unifiedOnChainWsService.unsubscribe(subscriptionId) unifiedOnChainWsService.unsubscribe(subscriptionId)
logger.info("移除 Leader 监听: leaderId=$leaderId") logger.info("移除 Leader 监听: leaderId=$leaderId")
} }
/** /**
* 停止监听 * 停止监听
*/ */
@@ -180,7 +186,7 @@ class OnChainWsService(
} }
monitoredLeaders.clear() monitoredLeaders.clear()
} }
@PreDestroy @PreDestroy
fun destroy() { fun destroy() {
stop() stop()