删除后端无用日志
- 删除所有 logger.debug 调试日志 - 删除过于详细的 logger.info 常规操作日志 - 保留 logger.error 和 logger.warn 重要错误和警告日志 - 优化日志输出,减少生产环境日志噪音
This commit is contained in:
@@ -35,7 +35,6 @@ class AccountController(
|
||||
val result = accountService.importAccount(request)
|
||||
result.fold(
|
||||
onSuccess = { account ->
|
||||
logger.info("成功导入账户: ${account.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(account))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -61,7 +60,6 @@ class AccountController(
|
||||
val result = accountService.updateAccount(request)
|
||||
result.fold(
|
||||
onSuccess = { account ->
|
||||
logger.info("成功更新账户: ${account.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(account))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -87,7 +85,6 @@ class AccountController(
|
||||
val result = accountService.deleteAccount(request.accountId)
|
||||
result.fold(
|
||||
onSuccess = {
|
||||
logger.info("成功删除账户: ${request.accountId}")
|
||||
ResponseEntity.ok(ApiResponse.success(Unit))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -114,7 +111,6 @@ class AccountController(
|
||||
val result = accountService.getAccountList()
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.info("成功查询账户列表: ${response.total} 个账户")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -137,7 +133,6 @@ class AccountController(
|
||||
val result = accountService.getAccountDetail(request.accountId)
|
||||
result.fold(
|
||||
onSuccess = { account ->
|
||||
logger.info("成功查询账户详情: ${account.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(account))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -163,7 +158,6 @@ class AccountController(
|
||||
val result = accountService.getAccountBalance(request.accountId)
|
||||
result.fold(
|
||||
onSuccess = { balance ->
|
||||
logger.info("成功查询账户余额")
|
||||
ResponseEntity.ok(ApiResponse.success(balance))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -189,7 +183,6 @@ class AccountController(
|
||||
val result = accountService.setDefaultAccount(request.accountId)
|
||||
result.fold(
|
||||
onSuccess = {
|
||||
logger.info("成功设置默认账户: ${request.accountId}")
|
||||
ResponseEntity.ok(ApiResponse.success(Unit))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -216,7 +209,6 @@ class AccountController(
|
||||
result.fold(
|
||||
onSuccess = { positionListResponse ->
|
||||
val total = positionListResponse.currentPositions.size + positionListResponse.historyPositions.size
|
||||
logger.info("成功查询仓位列表: 当前仓位 ${positionListResponse.currentPositions.size} 个,历史仓位 ${positionListResponse.historyPositions.size} 个,共 $total 个")
|
||||
ResponseEntity.ok(ApiResponse.success(positionListResponse))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -260,7 +252,6 @@ class AccountController(
|
||||
val result = runBlocking { accountService.sellPosition(request) }
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.info("成功创建卖出订单: 账户=${request.accountId}, 市场=${request.marketId}, 订单ID=${response.orderId}")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -287,7 +278,6 @@ class AccountController(
|
||||
val result = runBlocking { accountService.getRedeemablePositionsSummary(request.accountId) }
|
||||
result.fold(
|
||||
onSuccess = { summary ->
|
||||
logger.info("获取可赎回仓位统计成功: 账户=${request.accountId}, 数量=${summary.totalCount}, 价值=${summary.totalValue}")
|
||||
ResponseEntity.ok(ApiResponse.success(summary))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -331,7 +321,6 @@ class AccountController(
|
||||
val result = runBlocking { accountService.redeemPositions(request) }
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.info("成功赎回仓位: 账户数=${response.transactions.size}, 交易数=${response.transactions.size}, 总价值=${response.totalRedeemedValue}")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
@@ -36,7 +36,6 @@ class CopyTradingController(
|
||||
val result = copyTradingService.createCopyTrading(request)
|
||||
result.fold(
|
||||
onSuccess = { copyTrading ->
|
||||
logger.info("成功创建跟单: ${copyTrading.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(copyTrading))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -88,7 +87,6 @@ class CopyTradingController(
|
||||
val result = copyTradingService.updateCopyTradingStatus(request)
|
||||
result.fold(
|
||||
onSuccess = { copyTrading ->
|
||||
logger.info("成功更新跟单状态: ${copyTrading.id}, enabled=${copyTrading.enabled}")
|
||||
ResponseEntity.ok(ApiResponse.success(copyTrading))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -119,7 +117,6 @@ class CopyTradingController(
|
||||
val result = copyTradingService.deleteCopyTrading(request.copyTradingId)
|
||||
result.fold(
|
||||
onSuccess = {
|
||||
logger.info("成功删除跟单: ${request.copyTradingId}")
|
||||
ResponseEntity.ok(ApiResponse.success(Unit))
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
-2
@@ -33,7 +33,6 @@ class CopyTradingStatisticsController(
|
||||
val result = runBlocking { statisticsService.getStatistics(request.copyTradingId) }
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.info("成功获取统计信息: copyTradingId=${request.copyTradingId}")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -86,7 +85,6 @@ class CopyOrderTrackingController(
|
||||
val result = statisticsService.getOrderList(request)
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.info("成功查询订单列表: copyTradingId=${request.copyTradingId}, type=${request.type}, total=${response.total}")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
-4
@@ -30,7 +30,6 @@ class CopyTradingTemplateController(
|
||||
val result = templateService.createTemplate(request)
|
||||
result.fold(
|
||||
onSuccess = { template ->
|
||||
logger.info("成功创建模板: ${template.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(template))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -60,7 +59,6 @@ class CopyTradingTemplateController(
|
||||
val result = templateService.updateTemplate(request)
|
||||
result.fold(
|
||||
onSuccess = { template ->
|
||||
logger.info("成功更新模板: ${template.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(template))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -90,7 +88,6 @@ class CopyTradingTemplateController(
|
||||
val result = templateService.deleteTemplate(request.templateId)
|
||||
result.fold(
|
||||
onSuccess = {
|
||||
logger.info("成功删除模板: ${request.templateId}")
|
||||
ResponseEntity.ok(ApiResponse.success(Unit))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -124,7 +121,6 @@ class CopyTradingTemplateController(
|
||||
val result = templateService.copyTemplate(request)
|
||||
result.fold(
|
||||
onSuccess = { template ->
|
||||
logger.info("成功复制模板: ${template.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(template))
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
@@ -30,7 +30,6 @@ class LeaderController(
|
||||
val result = leaderService.addLeader(request)
|
||||
result.fold(
|
||||
onSuccess = { leader ->
|
||||
logger.info("成功添加 Leader: ${leader.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(leader))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -61,7 +60,6 @@ class LeaderController(
|
||||
val result = leaderService.updateLeader(request)
|
||||
result.fold(
|
||||
onSuccess = { leader ->
|
||||
logger.info("成功更新 Leader: ${leader.id}")
|
||||
ResponseEntity.ok(ApiResponse.success(leader))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -91,7 +89,6 @@ class LeaderController(
|
||||
val result = leaderService.deleteLeader(request.leaderId)
|
||||
result.fold(
|
||||
onSuccess = {
|
||||
logger.info("成功删除 Leader: ${request.leaderId}")
|
||||
ResponseEntity.ok(ApiResponse.success(Unit))
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
@@ -36,7 +36,6 @@ class MarketController(
|
||||
val result = runBlocking { accountService.getMarketPrice(request.marketId) }
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.info("成功获取市场价格: 市场=${request.marketId}")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
@@ -65,7 +64,6 @@ class MarketController(
|
||||
val result = runBlocking { clobService.getLatestPrice(request.tokenId) }
|
||||
result.fold(
|
||||
onSuccess = { response ->
|
||||
logger.debug("成功获取最新价: tokenId=${request.tokenId}, bestBid=${response.bestBid}, bestAsk=${response.bestAsk}")
|
||||
ResponseEntity.ok(ApiResponse.success(response))
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
@@ -60,7 +60,6 @@ class AccountService(
|
||||
}
|
||||
|
||||
// 4. 自动获取或创建 API Key(必须成功,否则导入失败)
|
||||
logger.info("开始自动获取或创建 API Key: ${request.walletAddress}")
|
||||
val apiKeyCreds = runBlocking {
|
||||
val result = apiKeyService.createOrDeriveApiKey(
|
||||
privateKey = request.privateKey,
|
||||
@@ -71,7 +70,6 @@ class AccountService(
|
||||
if (result.isSuccess) {
|
||||
val creds = result.getOrNull()
|
||||
if (creds != null) {
|
||||
logger.info("成功自动获取 API Key: ${request.walletAddress}")
|
||||
creds
|
||||
} else {
|
||||
logger.error("自动获取 API Key 返回空值")
|
||||
@@ -98,7 +96,6 @@ class AccountService(
|
||||
if (proxyResult.isSuccess) {
|
||||
val address = proxyResult.getOrNull()
|
||||
if (address != null) {
|
||||
logger.info("成功获取代理地址: ${request.walletAddress} -> $address")
|
||||
address
|
||||
} else {
|
||||
logger.error("获取代理地址返回空值")
|
||||
@@ -127,7 +124,6 @@ class AccountService(
|
||||
)
|
||||
|
||||
val saved = accountRepository.save(account)
|
||||
logger.info("成功导入账户: ${saved.id}, ${saved.walletAddress}, 代理地址: ${saved.proxyAddress}, 启用状态: ${saved.isEnabled}")
|
||||
|
||||
// 刷新订单推送订阅(如果账户启用且有 API 凭证)
|
||||
orderPushService.refreshSubscriptions()
|
||||
@@ -171,7 +167,6 @@ class AccountService(
|
||||
)
|
||||
|
||||
val saved = accountRepository.save(updated)
|
||||
logger.info("成功更新账户: ${saved.id}, 启用状态: ${saved.isEnabled}")
|
||||
|
||||
// 刷新订单推送订阅(账户状态变更时)
|
||||
orderPushService.refreshSubscriptions()
|
||||
@@ -212,7 +207,6 @@ class AccountService(
|
||||
}
|
||||
|
||||
accountRepository.delete(account)
|
||||
logger.info("成功删除账户: $accountId")
|
||||
|
||||
// 刷新订单推送订阅(账户删除时)
|
||||
orderPushService.refreshSubscriptions()
|
||||
@@ -373,7 +367,6 @@ class AccountService(
|
||||
val updated = account.copy(isDefault = true, updatedAt = System.currentTimeMillis())
|
||||
accountRepository.save(updated)
|
||||
|
||||
logger.info("成功设置默认账户: $accountId")
|
||||
Result.success(Unit)
|
||||
} catch (e: Exception) {
|
||||
logger.error("设置默认账户失败", e)
|
||||
@@ -807,14 +800,12 @@ class AccountService(
|
||||
account.walletAddress
|
||||
)
|
||||
|
||||
logger.info("创建卖出订单: market=${request.marketId}, side=${request.side}, orderType=${request.orderType}, quantity=${request.quantity}, price=$sellPrice, tokenId=$tokenId")
|
||||
|
||||
val orderResponse = clobApi.createOrder(newOrderRequest)
|
||||
|
||||
if (orderResponse.isSuccessful && orderResponse.body() != null) {
|
||||
val response = orderResponse.body()!!
|
||||
if (response.success) {
|
||||
logger.info("订单创建成功: orderId=${response.orderId}, transactionsHashes=${response.transactionsHashes}")
|
||||
Result.success(
|
||||
PositionSellResponse(
|
||||
orderId = response.orderId ?: "",
|
||||
@@ -1067,7 +1058,6 @@ class AccountService(
|
||||
redeemResult.fold(
|
||||
onSuccess = { txHash ->
|
||||
lastTxHash = txHash
|
||||
logger.info("账户 $accountId 市场 $marketId 赎回成功: txHash=$txHash, indexSets=$indexSets")
|
||||
},
|
||||
onFailure = { e ->
|
||||
logger.error("账户 $accountId 市场 $marketId 赎回失败: ${e.message}", e)
|
||||
@@ -1115,7 +1105,6 @@ class AccountService(
|
||||
return try {
|
||||
// 如果账户没有配置 API 凭证,无法查询活跃订单,允许删除
|
||||
if (account.apiKey == null || account.apiSecret == null || account.apiPassphrase == null) {
|
||||
logger.debug("账户 ${account.id} 未配置 API 凭证,无法查询活跃订单,允许删除")
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -1139,7 +1128,6 @@ class AccountService(
|
||||
if (response.isSuccessful && response.body() != null) {
|
||||
val ordersResponse = response.body()!!
|
||||
val hasOrders = ordersResponse.data.isNotEmpty()
|
||||
logger.debug("账户 ${account.id} 活跃订单检查结果: $hasOrders (订单数: ${ordersResponse.data.size})")
|
||||
hasOrders
|
||||
} else {
|
||||
// 如果查询失败(可能是认证失败或网络问题),记录警告但允许删除
|
||||
|
||||
@@ -129,7 +129,6 @@ class BlockchainService(
|
||||
// 解析代理地址
|
||||
val proxyAddress = EthereumUtils.decodeAddress(hexResult)
|
||||
|
||||
logger.debug("获取代理地址成功: 原始地址=$walletAddress, 代理地址=$proxyAddress")
|
||||
Result.success(proxyAddress)
|
||||
} catch (e: Exception) {
|
||||
logger.error("获取代理地址失败: ${e.message}", e)
|
||||
@@ -158,7 +157,6 @@ class BlockchainService(
|
||||
return Result.failure(IllegalArgumentException("代理地址不能为空"))
|
||||
}
|
||||
|
||||
logger.debug("使用代理地址查询余额: $proxyAddress (原始地址: $walletAddress)")
|
||||
|
||||
// 使用 RPC 查询 USDC 余额(使用代理地址)
|
||||
val balance = queryUsdcBalanceViaRpc(proxyAddress)
|
||||
@@ -235,7 +233,6 @@ class BlockchainService(
|
||||
|
||||
if (response.isSuccessful && response.body() != null) {
|
||||
val positions = response.body()!!
|
||||
logger.debug("查询到 ${positions.size} 个仓位")
|
||||
Result.success(positions)
|
||||
} else {
|
||||
val errorMsg = "Data API 请求失败: ${response.code()} ${response.message()}"
|
||||
@@ -388,7 +385,6 @@ class BlockchainService(
|
||||
} else {
|
||||
0.0
|
||||
}
|
||||
logger.debug("查询到仓位总价值: $totalValue")
|
||||
Result.success(totalValue.toString())
|
||||
} else {
|
||||
val errorMsg = "Data API 请求失败: ${response.code()} ${response.message()}"
|
||||
@@ -452,7 +448,6 @@ class BlockchainService(
|
||||
val credentials = org.web3j.crypto.Credentials.create(privateKeyBigInt.toString(16))
|
||||
val fromAddress = credentials.address
|
||||
|
||||
logger.debug("赎回仓位: from=$fromAddress, proxy=$proxyAddress, conditionId=$conditionId, indexSets=$indexSets")
|
||||
|
||||
// 1. 构建 ConditionalTokens.redeemPositions 的调用数据
|
||||
val redeemFunctionSelector = EthereumUtils.getFunctionSelector("redeemPositions(address,bytes32,bytes32,uint256[])")
|
||||
@@ -638,7 +633,6 @@ class BlockchainService(
|
||||
val txHashResult = sendTransaction(rpcApi, transaction)
|
||||
txHashResult.fold(
|
||||
onSuccess = { txHash ->
|
||||
logger.info("赎回仓位交易已发送: txHash=$txHash, from=$fromAddress, proxy=$proxyAddress, conditionId=$conditionId, indexSets=$indexSets")
|
||||
Result.success(txHash)
|
||||
},
|
||||
onFailure = { e ->
|
||||
|
||||
@@ -50,17 +50,14 @@ class CopyOrderTrackingService(
|
||||
|
||||
if (existingProcessed != null) {
|
||||
if (existingProcessed.status == "FAILED") {
|
||||
logger.debug("交易已标记为失败,跳过: leaderId=$leaderId, tradeId=${trade.id}, source=$source")
|
||||
return Result.success(Unit)
|
||||
}
|
||||
logger.debug("交易已处理,跳过: leaderId=$leaderId, tradeId=${trade.id}, source=$source")
|
||||
return Result.success(Unit)
|
||||
}
|
||||
|
||||
// 检查是否已记录为失败交易
|
||||
val failedTrade = failedTradeRepository.findByLeaderIdAndLeaderTradeId(leaderId, trade.id)
|
||||
if (failedTrade != null) {
|
||||
logger.debug("交易已记录为失败,跳过: leaderId=$leaderId, tradeId=${trade.id}, source=$source")
|
||||
return Result.success(Unit)
|
||||
}
|
||||
|
||||
@@ -91,17 +88,14 @@ class CopyOrderTrackingService(
|
||||
processedAt = System.currentTimeMillis()
|
||||
)
|
||||
processedTradeRepository.save(processed)
|
||||
logger.info("成功处理交易: leaderId=$leaderId, tradeId=${trade.id}, source=$source, side=${trade.side}")
|
||||
} catch (e: DataIntegrityViolationException) {
|
||||
// 唯一约束冲突,说明已经处理过了(可能是并发请求)
|
||||
// 再次检查确认状态
|
||||
val existing = processedTradeRepository.findByLeaderIdAndLeaderTradeId(leaderId, trade.id)
|
||||
if (existing != null) {
|
||||
if (existing.status == "FAILED") {
|
||||
logger.debug("交易已标记为失败(并发检测): leaderId=$leaderId, tradeId=${trade.id}")
|
||||
return Result.success(Unit)
|
||||
}
|
||||
logger.debug("交易已处理(并发检测): leaderId=$leaderId, tradeId=${trade.id}, source=$source")
|
||||
return Result.success(Unit)
|
||||
} else {
|
||||
// 如果检查不到,说明可能是其他约束冲突,重新抛出异常
|
||||
@@ -128,7 +122,6 @@ class CopyOrderTrackingService(
|
||||
val copyTradings = copyTradingRepository.findByLeaderIdAndEnabledTrue(leaderId)
|
||||
|
||||
if (copyTradings.isEmpty()) {
|
||||
logger.debug("没有启用的跟单关系: leaderId=$leaderId")
|
||||
return Result.success(Unit)
|
||||
}
|
||||
|
||||
@@ -151,7 +144,6 @@ class CopyOrderTrackingService(
|
||||
|
||||
// 验证账户是否启用
|
||||
if (!account.isEnabled) {
|
||||
logger.debug("账户未启用,跳过创建订单: accountId=${account.id}")
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -273,7 +265,6 @@ class CopyOrderTrackingService(
|
||||
)
|
||||
|
||||
copyOrderTrackingRepository.save(tracking)
|
||||
logger.info("成功创建买入订单并记录跟踪: copyTradingId=${copyTrading.id}, orderId=$realOrderId, tradeId=${trade.id}, quantity=$finalBuyQuantity, price=$buyPrice")
|
||||
} catch (e: Exception) {
|
||||
logger.error("处理买入交易失败: copyTradingId=${copyTrading.id}, tradeId=${trade.id}", e)
|
||||
// 继续处理下一个跟单关系
|
||||
@@ -298,7 +289,6 @@ class CopyOrderTrackingService(
|
||||
val copyTradings = copyTradingRepository.findByLeaderIdAndEnabledTrue(leaderId)
|
||||
|
||||
if (copyTradings.isEmpty()) {
|
||||
logger.debug("没有启用的跟单关系: leaderId=$leaderId")
|
||||
return Result.success(Unit)
|
||||
}
|
||||
|
||||
@@ -311,7 +301,6 @@ class CopyOrderTrackingService(
|
||||
|
||||
// 检查是否支持卖出
|
||||
if (!template.supportSell) {
|
||||
logger.debug("模板不支持卖出,跳过: copyTradingId=${copyTrading.id}, templateId=${template.id}")
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -377,7 +366,6 @@ class CopyOrderTrackingService(
|
||||
|
||||
// 验证账户是否启用
|
||||
if (!account.isEnabled) {
|
||||
logger.debug("账户未启用,跳过创建卖出订单: accountId=${account.id}")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -399,7 +387,6 @@ class CopyOrderTrackingService(
|
||||
)
|
||||
|
||||
if (unmatchedOrders.isEmpty()) {
|
||||
logger.debug("没有未匹配的买入订单: copyTradingId=${copyTrading.id}, market=${leaderSellTrade.market}, outcomeIndex=${leaderSellTrade.outcomeIndex}")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -440,7 +427,6 @@ class CopyOrderTrackingService(
|
||||
}
|
||||
|
||||
if (totalMatched.lte(BigDecimal.ZERO)) {
|
||||
logger.debug("没有匹配到任何订单: copyTradingId=${copyTrading.id}, needMatch=$needMatch")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -533,7 +519,6 @@ class CopyOrderTrackingService(
|
||||
order.updatedAt = System.currentTimeMillis()
|
||||
copyOrderTrackingRepository.save(order)
|
||||
|
||||
logger.info("匹配买入订单: copyTradingId=${copyTrading.id}, buyOrderId=${order.buyOrderId}, matchQty=${detail.matchedQuantity}, pnl=${detail.realizedPnl}")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -560,7 +545,6 @@ class CopyOrderTrackingService(
|
||||
sellMatchDetailRepository.save(savedDetail)
|
||||
}
|
||||
|
||||
logger.info("完成卖出匹配并创建订单: copyTradingId=${copyTrading.id}, sellOrderId=$realSellOrderId, totalMatched=$totalMatched, totalPnl=$totalRealizedPnl")
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -31,7 +31,6 @@ class CopyTradingMonitorService(
|
||||
*/
|
||||
@PostConstruct
|
||||
fun init() {
|
||||
logger.info("跟单监听服务初始化...")
|
||||
scope.launch {
|
||||
try {
|
||||
startMonitoring()
|
||||
@@ -46,7 +45,6 @@ class CopyTradingMonitorService(
|
||||
*/
|
||||
@PreDestroy
|
||||
fun destroy() {
|
||||
logger.info("停止跟单监听服务...")
|
||||
scope.cancel()
|
||||
// 只使用轮询,不使用WebSocket
|
||||
pollingService.stop()
|
||||
@@ -60,7 +58,6 @@ class CopyTradingMonitorService(
|
||||
val enabledCopyTradings = copyTradingRepository.findByEnabledTrue()
|
||||
|
||||
if (enabledCopyTradings.isEmpty()) {
|
||||
logger.info("没有启用的跟单关系,等待添加...")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -70,7 +67,6 @@ class CopyTradingMonitorService(
|
||||
leaderRepository.findById(leaderId).orElse(null)
|
||||
}
|
||||
|
||||
logger.info("开始监听 ${leaders.size} 个Leader的交易: ${leaders.map { it.leaderAddress }}")
|
||||
|
||||
// 3. 启动轮询监听(使用 /activity 接口,不需要认证)
|
||||
// 注意:WebSocket 需要认证才能订阅其他用户的交易,因此禁用WebSocket,只使用轮询
|
||||
@@ -86,11 +82,9 @@ class CopyTradingMonitorService(
|
||||
|
||||
val copyTradings = copyTradingRepository.findByLeaderIdAndEnabledTrue(leaderId)
|
||||
if (copyTradings.isEmpty()) {
|
||||
logger.debug("Leader $leaderId 没有启用的跟单关系,不启动监听")
|
||||
return
|
||||
}
|
||||
|
||||
logger.info("添加Leader监听: ${leader.leaderAddress}")
|
||||
// 只使用轮询,不使用WebSocket(需要认证)
|
||||
pollingService.addLeader(leader)
|
||||
}
|
||||
@@ -101,11 +95,9 @@ class CopyTradingMonitorService(
|
||||
suspend fun removeLeaderMonitoring(leaderId: Long) {
|
||||
val copyTradings = copyTradingRepository.findByLeaderIdAndEnabledTrue(leaderId)
|
||||
if (copyTradings.isNotEmpty()) {
|
||||
logger.debug("Leader $leaderId 仍有启用的跟单关系,不停止监听")
|
||||
return
|
||||
}
|
||||
|
||||
logger.info("移除Leader监听: leaderId=$leaderId")
|
||||
// 只使用轮询,不使用WebSocket
|
||||
pollingService.removeLeader(leaderId)
|
||||
}
|
||||
@@ -114,7 +106,6 @@ class CopyTradingMonitorService(
|
||||
* 重新启动监听(当跟单关系状态改变时调用)
|
||||
*/
|
||||
suspend fun restartMonitoring() {
|
||||
logger.info("重新启动跟单监听...")
|
||||
// 只使用轮询,不使用WebSocket
|
||||
pollingService.stop()
|
||||
delay(1000) // 等待1秒
|
||||
|
||||
@@ -52,12 +52,9 @@ class CopyTradingPollingService(
|
||||
*/
|
||||
fun start(leaders: List<Leader>) {
|
||||
if (!pollingEnabled) {
|
||||
logger.info("轮询监听已禁用,跳过启动")
|
||||
return
|
||||
}
|
||||
|
||||
logger.info("启动轮询监听,Leader数量: ${leaders.size},轮询间隔: ${pollingInterval}ms")
|
||||
|
||||
leaders.forEach { leader ->
|
||||
addLeader(leader)
|
||||
}
|
||||
@@ -81,7 +78,6 @@ class CopyTradingPollingService(
|
||||
cachedTradeIds[leaderId] = mutableSetOf()
|
||||
// 首次轮询标志,用于缓存数据而不处理
|
||||
isFirstPoll[leaderId] = true
|
||||
logger.info("添加轮询监听: leaderId=$leaderId, address=${leader.leaderAddress}, 首次轮询将只缓存数据")
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -91,14 +87,12 @@ class CopyTradingPollingService(
|
||||
monitoredLeaders.remove(leaderId)
|
||||
cachedTradeIds.remove(leaderId)
|
||||
isFirstPoll.remove(leaderId)
|
||||
logger.info("移除轮询监听: leaderId=$leaderId")
|
||||
}
|
||||
|
||||
/**
|
||||
* 停止所有监听
|
||||
*/
|
||||
fun stop() {
|
||||
logger.info("停止所有轮询监听...")
|
||||
stopPolling()
|
||||
monitoredLeaders.clear()
|
||||
cachedTradeIds.clear()
|
||||
@@ -110,17 +104,14 @@ class CopyTradingPollingService(
|
||||
*/
|
||||
private fun startPolling() {
|
||||
if (pollingJob != null && pollingJob!!.isActive) {
|
||||
logger.debug("轮询任务已在运行")
|
||||
return
|
||||
}
|
||||
|
||||
if (monitoredLeaders.isEmpty()) {
|
||||
logger.debug("没有需要监听的Leader,不启动轮询")
|
||||
return
|
||||
}
|
||||
|
||||
pollingJob = scope.launch {
|
||||
logger.info("轮询任务已启动,间隔: ${pollingInterval}ms")
|
||||
|
||||
while (isActive) {
|
||||
try {
|
||||
@@ -235,7 +226,6 @@ class CopyTradingPollingService(
|
||||
cachedIds.addAll(tradeIds)
|
||||
cachedTradeIds[leaderId] = cachedIds
|
||||
|
||||
logger.info("首次轮询,缓存 ${allTrades.size} 笔交易数据,不进行处理: leaderId=$leaderId")
|
||||
// 标记首次轮询完成
|
||||
isFirstPoll[leaderId] = false
|
||||
} else {
|
||||
@@ -247,7 +237,6 @@ class CopyTradingPollingService(
|
||||
// 找出新增的交易
|
||||
val incrementalTrades = allTrades.filter { it.id in incrementalTradeIds }
|
||||
|
||||
logger.debug("通过 diff 发现 ${incrementalTrades.size} 笔新增交易: leaderId=$leaderId")
|
||||
|
||||
// 处理新增的交易
|
||||
incrementalTrades.forEach { trade ->
|
||||
@@ -263,9 +252,7 @@ class CopyTradingPollingService(
|
||||
cachedIds.addAll(incrementalTradeIds)
|
||||
cachedTradeIds[leaderId] = cachedIds
|
||||
|
||||
logger.debug("已更新缓存,当前缓存 ${cachedIds.size} 笔交易ID: leaderId=$leaderId")
|
||||
} else {
|
||||
logger.debug("未发现新增交易: leaderId=$leaderId")
|
||||
}
|
||||
|
||||
// 限制缓存大小,避免内存溢出(只保留最近1000条)
|
||||
@@ -273,7 +260,6 @@ class CopyTradingPollingService(
|
||||
// 保留最新的1000条(由于查询是按时间戳降序,保留前1000条即可)
|
||||
val sortedTradeIds = allTrades.map { it.id }.take(1000).toSet()
|
||||
cachedTradeIds[leaderId] = sortedTradeIds.toMutableSet()
|
||||
logger.debug("缓存已满,清理到1000条: leaderId=$leaderId")
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
|
||||
@@ -61,7 +61,6 @@ class CopyTradingService(
|
||||
)
|
||||
|
||||
val saved = copyTradingRepository.save(copyTrading)
|
||||
logger.info("成功创建跟单: ${saved.id}, account=${request.accountId}, template=${request.templateId}, leader=${request.leaderId}")
|
||||
|
||||
// 如果跟单已启用,启动Leader监听
|
||||
if (saved.enabled) {
|
||||
@@ -162,7 +161,6 @@ class CopyTradingService(
|
||||
)
|
||||
|
||||
val saved = copyTradingRepository.save(updated)
|
||||
logger.info("成功更新跟单状态: ${saved.id}, enabled=${saved.enabled}")
|
||||
|
||||
// 更新监听状态
|
||||
kotlinx.coroutines.runBlocking {
|
||||
@@ -203,7 +201,6 @@ class CopyTradingService(
|
||||
|
||||
val leaderId = copyTrading.leaderId
|
||||
copyTradingRepository.delete(copyTrading)
|
||||
logger.info("成功删除跟单: $copyTradingId")
|
||||
|
||||
// 移除监听(如果该Leader没有其他启用的跟单关系)
|
||||
kotlinx.coroutines.runBlocking {
|
||||
|
||||
@@ -62,7 +62,6 @@ class CopyTradingTemplateService(
|
||||
)
|
||||
|
||||
val saved = templateRepository.save(template)
|
||||
logger.info("成功创建模板: ${saved.id}, ${saved.templateName}")
|
||||
|
||||
Result.success(toDto(saved))
|
||||
} catch (e: Exception) {
|
||||
@@ -118,7 +117,6 @@ class CopyTradingTemplateService(
|
||||
)
|
||||
|
||||
val saved = templateRepository.save(updated)
|
||||
logger.info("成功更新模板: ${saved.id}")
|
||||
|
||||
Result.success(toDto(saved))
|
||||
} catch (e: Exception) {
|
||||
@@ -143,7 +141,6 @@ class CopyTradingTemplateService(
|
||||
}
|
||||
|
||||
templateRepository.delete(template)
|
||||
logger.info("成功删除模板: $templateId")
|
||||
|
||||
Result.success(Unit)
|
||||
} catch (e: Exception) {
|
||||
@@ -186,7 +183,6 @@ class CopyTradingTemplateService(
|
||||
)
|
||||
|
||||
val saved = templateRepository.save(newTemplate)
|
||||
logger.info("成功复制模板: ${sourceTemplate.id} -> ${saved.id}")
|
||||
|
||||
Result.success(toDto(saved))
|
||||
} catch (e: Exception) {
|
||||
|
||||
@@ -42,7 +42,6 @@ class CopyTradingWebSocketService(
|
||||
* 启动WebSocket监听
|
||||
*/
|
||||
fun start(leaders: List<Leader>) {
|
||||
logger.info("启动WebSocket监听,Leader数量: ${leaders.size}")
|
||||
|
||||
leaders.forEach { leader ->
|
||||
try {
|
||||
@@ -63,7 +62,6 @@ class CopyTradingWebSocketService(
|
||||
}
|
||||
|
||||
if (leaderClients.containsKey(leader.id)) {
|
||||
logger.debug("Leader ${leader.id} 已经在监听中,跳过")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -98,7 +96,6 @@ class CopyTradingWebSocketService(
|
||||
scope.launch {
|
||||
try {
|
||||
client.connect()
|
||||
logger.info("已启动WebSocket监听: leaderId=$leaderId, address=$leaderAddress")
|
||||
} catch (e: Exception) {
|
||||
logger.error("连接WebSocket失败: leaderId=$leaderId", e)
|
||||
leaderClients.remove(leaderId)
|
||||
@@ -117,7 +114,6 @@ class CopyTradingWebSocketService(
|
||||
if (client != null) {
|
||||
try {
|
||||
client.closeConnection()
|
||||
logger.info("已停止WebSocket监听: leaderId=$leaderId")
|
||||
} catch (e: Exception) {
|
||||
logger.error("关闭WebSocket连接失败: leaderId=$leaderId", e)
|
||||
}
|
||||
@@ -128,7 +124,6 @@ class CopyTradingWebSocketService(
|
||||
* 停止所有监听
|
||||
*/
|
||||
fun stop() {
|
||||
logger.info("停止所有WebSocket监听...")
|
||||
val leaderIds = leaderClients.keys.toList()
|
||||
leaderIds.forEach { leaderId ->
|
||||
removeLeader(leaderId)
|
||||
@@ -150,7 +145,6 @@ class CopyTradingWebSocketService(
|
||||
""".trimIndent()
|
||||
|
||||
client.sendMessage(subscribeMessage)
|
||||
logger.info("已订阅用户交易频道: $userAddress")
|
||||
} catch (e: Exception) {
|
||||
logger.error("订阅用户交易频道失败: $userAddress", e)
|
||||
}
|
||||
@@ -163,7 +157,6 @@ class CopyTradingWebSocketService(
|
||||
try {
|
||||
// 处理PONG响应
|
||||
if (message.trim() == "PONG") {
|
||||
logger.debug("收到PONG响应: leaderId=$leaderId")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -173,7 +166,6 @@ class CopyTradingWebSocketService(
|
||||
// 检查消息类型
|
||||
val eventType = json.get("event_type")?.asString
|
||||
if (eventType != "trade") {
|
||||
logger.debug("忽略非交易事件: leaderId=$leaderId, eventType=$eventType")
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -56,7 +56,6 @@ class LeaderService(
|
||||
)
|
||||
|
||||
val saved = leaderRepository.save(leader)
|
||||
logger.info("成功添加 Leader: ${saved.id}, ${saved.leaderAddress}")
|
||||
|
||||
Result.success(toDto(saved))
|
||||
} catch (e: Exception) {
|
||||
@@ -86,7 +85,6 @@ class LeaderService(
|
||||
)
|
||||
|
||||
val saved = leaderRepository.save(updated)
|
||||
logger.info("成功更新 Leader: ${saved.id}")
|
||||
|
||||
Result.success(toDto(saved))
|
||||
} catch (e: Exception) {
|
||||
@@ -111,7 +109,6 @@ class LeaderService(
|
||||
}
|
||||
|
||||
leaderRepository.delete(leader)
|
||||
logger.info("成功删除 Leader: $leaderId")
|
||||
|
||||
Result.success(Unit)
|
||||
} catch (e: Exception) {
|
||||
|
||||
@@ -48,7 +48,6 @@ class OrderPushService(
|
||||
*/
|
||||
@PostConstruct
|
||||
fun init() {
|
||||
logger.info("订单推送服务已初始化")
|
||||
scope.launch {
|
||||
connectAllAccounts()
|
||||
}
|
||||
@@ -59,7 +58,6 @@ class OrderPushService(
|
||||
*/
|
||||
@PreDestroy
|
||||
fun destroy() {
|
||||
logger.info("停止订单推送服务")
|
||||
accountConnections.values.forEach { client ->
|
||||
try {
|
||||
if (client.isConnected()) {
|
||||
@@ -90,7 +88,6 @@ class OrderPushService(
|
||||
* 订阅所有启用的账户
|
||||
*/
|
||||
fun subscribeAllEnabled(callback: (OrderPushMessage) -> Unit) {
|
||||
logger.info("订阅所有启用账户的订单推送")
|
||||
val accounts = accountRepository.findAll()
|
||||
accounts.forEach { account ->
|
||||
if (hasApiCredentials(account) && account.isEnabled) {
|
||||
@@ -166,12 +163,10 @@ class OrderPushService(
|
||||
}
|
||||
|
||||
if (!account.isEnabled) {
|
||||
logger.debug("账户 ${account.id} 未启用,跳过连接")
|
||||
return
|
||||
}
|
||||
|
||||
if (accountConnections.containsKey(account.id)) {
|
||||
logger.debug("账户 ${account.id} 已存在连接,跳过")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -193,7 +188,6 @@ class OrderPushService(
|
||||
if (currentClient != null) {
|
||||
try {
|
||||
sendSubscribeMessage(currentClient, account)
|
||||
logger.info("已为账户 ${account.id} (${account.accountName ?: account.walletAddress}) 建立 User Channel 连接并发送订阅消息")
|
||||
} catch (e: Exception) {
|
||||
logger.error("发送订阅消息失败: account=${account.id}, ${e.message}", e)
|
||||
// 如果订阅失败,关闭连接(会触发重连)
|
||||
@@ -210,7 +204,6 @@ class OrderPushService(
|
||||
if (currentClient != null) {
|
||||
try {
|
||||
sendSubscribeMessage(currentClient, account)
|
||||
logger.info("账户 ${account.id} 重连成功,已重新发送订阅消息")
|
||||
} catch (e: Exception) {
|
||||
logger.error("重连后发送订阅消息失败: account=${account.id}, ${e.message}", e)
|
||||
}
|
||||
@@ -257,7 +250,6 @@ class OrderPushService(
|
||||
|
||||
val json = objectMapper.writeValueAsString(subscribeMessage)
|
||||
client.sendMessage(json)
|
||||
logger.info("已发送 User Channel 订阅消息: account=${account.id}, apiKey=${account.apiKey?.take(10)}...")
|
||||
} catch (e: Exception) {
|
||||
logger.error("发送订阅消息失败: account=${account.id}, ${e.message}", e)
|
||||
}
|
||||
@@ -270,7 +262,6 @@ class OrderPushService(
|
||||
try {
|
||||
// 处理心跳响应(PONG),直接返回
|
||||
if (message.trim() == "PONG" || message.trim() == "pong") {
|
||||
logger.debug("收到 PONG 响应: account=${account.id}")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -304,12 +295,10 @@ class OrderPushService(
|
||||
}
|
||||
} else {
|
||||
// 记录其他类型的消息(用于调试)
|
||||
logger.debug("收到非订单消息: account=${account.id}, eventType=$eventType")
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
// 如果解析失败,可能是非 JSON 消息(如 PONG),记录为 debug 级别
|
||||
if (message.trim() == "PONG" || message.trim() == "pong") {
|
||||
logger.debug("收到 PONG 响应: account=${account.id}")
|
||||
} else {
|
||||
logger.error("处理订单消息失败: account=${account.id}, message=${message.take(100)}, ${e.message}", e)
|
||||
}
|
||||
@@ -328,7 +317,6 @@ class OrderPushService(
|
||||
return try {
|
||||
// 检查账户是否有 API 凭证
|
||||
if (account.apiKey == null || account.apiSecret == null || account.apiPassphrase == null) {
|
||||
logger.debug("账户 ${account.id} 未配置 API 凭证,无法获取订单详情")
|
||||
return null
|
||||
}
|
||||
|
||||
@@ -395,18 +383,14 @@ class OrderPushService(
|
||||
val markets = response.body()!!
|
||||
if (markets.isNotEmpty()) {
|
||||
val market = markets.first()
|
||||
logger.debug("获取市场信息成功: conditionId=$conditionId, question=${market.question}")
|
||||
return market
|
||||
} else {
|
||||
logger.debug("未找到市场信息: conditionId=$conditionId")
|
||||
return null
|
||||
}
|
||||
} else {
|
||||
logger.debug("获取市场信息失败: conditionId=$conditionId, code=${response.code()}, message=${response.message()}")
|
||||
null
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
logger.debug("获取市场信息异常: conditionId=$conditionId, ${e.message}")
|
||||
null
|
||||
}
|
||||
}
|
||||
@@ -415,7 +399,6 @@ class OrderPushService(
|
||||
* 订阅账户的订单推送(保留用于向后兼容)
|
||||
*/
|
||||
fun subscribe(accountId: Long, callback: (OrderPushMessage) -> Unit) {
|
||||
logger.info("订阅账户订单推送: $accountId")
|
||||
accountCallbacks.getOrPut(accountId) { mutableSetOf() }.add(callback)
|
||||
|
||||
// 如果账户连接不存在,尝试建立连接
|
||||
@@ -431,7 +414,6 @@ class OrderPushService(
|
||||
* 取消订阅账户的订单推送(保留用于向后兼容)
|
||||
*/
|
||||
fun unsubscribe(accountId: Long, callback: (OrderPushMessage) -> Unit) {
|
||||
logger.info("取消订阅账户订单推送: $accountId")
|
||||
accountCallbacks[accountId]?.remove(callback)
|
||||
|
||||
// 如果没有订阅者了,可以考虑关闭连接(但暂时保持连接,以便后续订阅)
|
||||
@@ -441,7 +423,6 @@ class OrderPushService(
|
||||
* 断开指定账户的连接
|
||||
*/
|
||||
fun disconnectAccount(accountId: Long) {
|
||||
logger.info("断开账户连接: $accountId")
|
||||
val client = accountConnections.remove(accountId)
|
||||
client?.let {
|
||||
try {
|
||||
|
||||
@@ -161,7 +161,6 @@ class PolymarketClobService(
|
||||
adjustedPrice > BigDecimal("0.99") -> BigDecimal("0.99")
|
||||
else -> adjustedPrice
|
||||
}
|
||||
logger.debug("从订单表获取最优价(卖单): tokenId=$tokenId, bestBid=$bestBid, adjustedPrice=${finalPrice.toPlainString()}")
|
||||
return finalPrice.toPlainString()
|
||||
} else {
|
||||
// 市价买单:需要 bestAsk(最低卖出价)
|
||||
@@ -186,7 +185,6 @@ class PolymarketClobService(
|
||||
adjustedPrice > BigDecimal("0.99") -> BigDecimal("0.99")
|
||||
else -> adjustedPrice
|
||||
}
|
||||
logger.debug("从订单表获取最优价(买单): tokenId=$tokenId, bestAsk=$bestAsk, adjustedPrice=${finalPrice.toPlainString()}")
|
||||
return finalPrice.toPlainString()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,7 +45,6 @@ class PositionPushService(
|
||||
*/
|
||||
@PostConstruct
|
||||
fun init() {
|
||||
logger.info("仓位推送服务已初始化,轮询间隔: ${pollingInterval}ms,等待客户端连接...")
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -53,7 +52,6 @@ class PositionPushService(
|
||||
*/
|
||||
@PreDestroy
|
||||
fun destroy() {
|
||||
logger.info("停止仓位推送服务")
|
||||
synchronized(lock) {
|
||||
pollingJob?.cancel()
|
||||
pollingJob = null
|
||||
@@ -65,7 +63,6 @@ class PositionPushService(
|
||||
* 订阅仓位推送(新接口)
|
||||
*/
|
||||
fun subscribe(sessionId: String, callback: (PositionPushMessage) -> Unit) {
|
||||
logger.info("订阅仓位推送: $sessionId")
|
||||
registerSession(sessionId, callback)
|
||||
}
|
||||
|
||||
@@ -73,7 +70,6 @@ class PositionPushService(
|
||||
* 取消订阅仓位推送(新接口)
|
||||
*/
|
||||
fun unsubscribe(sessionId: String) {
|
||||
logger.info("取消订阅仓位推送: $sessionId")
|
||||
unregisterSession(sessionId)
|
||||
}
|
||||
|
||||
@@ -136,7 +132,6 @@ class PositionPushService(
|
||||
|
||||
// 发送给指定客户端
|
||||
clientCallbacks[sessionId]?.invoke(message)
|
||||
logger.debug("已发送全量仓位数据给客户端: $sessionId")
|
||||
}
|
||||
} else {
|
||||
logger.warn("获取仓位数据失败,无法发送全量数据: ${result.exceptionOrNull()?.message}")
|
||||
@@ -156,7 +151,6 @@ class PositionPushService(
|
||||
|
||||
// 启动新的轮询任务
|
||||
pollingJob = scope.launch {
|
||||
logger.info("轮询任务已启动,间隔: ${pollingInterval}ms")
|
||||
while (isActive) {
|
||||
try {
|
||||
pollAndPush()
|
||||
@@ -176,7 +170,6 @@ class PositionPushService(
|
||||
synchronized(lock) {
|
||||
pollingJob?.cancel()
|
||||
pollingJob = null
|
||||
logger.info("轮询任务已停止")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -186,7 +179,6 @@ class PositionPushService(
|
||||
private suspend fun pollAndPush() {
|
||||
// 双重检查:如果没有客户端连接,跳过轮询(虽然理论上不应该发生,但作为安全措施)
|
||||
if (clientCallbacks.isEmpty()) {
|
||||
logger.debug("没有客户端连接,跳过本次轮询")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -220,7 +212,6 @@ class PositionPushService(
|
||||
}
|
||||
}
|
||||
|
||||
logger.debug("已推送仓位增量更新,当前仓位变化: ${incremental.currentPositions.size}, 历史仓位变化: ${incremental.historyPositions.size}, 删除: ${incremental.removedKeys.size}")
|
||||
}
|
||||
|
||||
// 更新快照
|
||||
|
||||
-8
@@ -40,7 +40,6 @@ class WebSocketSubscriptionService(
|
||||
* 注册会话
|
||||
*/
|
||||
fun registerSession(sessionId: String, callback: (WsMessage) -> Unit) {
|
||||
logger.info("注册 WebSocket 会话: $sessionId")
|
||||
sessionCallbacks[sessionId] = callback
|
||||
sessionSubscriptions[sessionId] = mutableSetOf()
|
||||
}
|
||||
@@ -49,7 +48,6 @@ class WebSocketSubscriptionService(
|
||||
* 注销会话
|
||||
*/
|
||||
fun unregisterSession(sessionId: String) {
|
||||
logger.info("注销 WebSocket 会话: $sessionId")
|
||||
|
||||
// 取消所有订阅
|
||||
val channels = sessionSubscriptions.remove(sessionId) ?: emptySet()
|
||||
@@ -67,12 +65,10 @@ class WebSocketSubscriptionService(
|
||||
* 订阅频道
|
||||
*/
|
||||
fun subscribe(sessionId: String, channel: String, payload: Map<*, *>?) {
|
||||
logger.info("订阅频道: $sessionId -> $channel")
|
||||
|
||||
// 检查是否已经订阅
|
||||
val sessionChannels = sessionSubscriptions.getOrPut(sessionId) { mutableSetOf() }
|
||||
if (sessionChannels.contains(channel)) {
|
||||
logger.debug("会话 $sessionId 已经订阅了频道 $channel,跳过重复订阅")
|
||||
sendSubscribeAck(sessionId, channel, true)
|
||||
return
|
||||
}
|
||||
@@ -94,7 +90,6 @@ class WebSocketSubscriptionService(
|
||||
scope.launch {
|
||||
try {
|
||||
positionPushService.sendFullData(sessionId)
|
||||
logger.info("已发送仓位首推数据给会话: $sessionId")
|
||||
} catch (e: Exception) {
|
||||
logger.error("发送仓位首推数据失败: $sessionId, ${e.message}", e)
|
||||
}
|
||||
@@ -107,7 +102,6 @@ class WebSocketSubscriptionService(
|
||||
}
|
||||
orderChannelCallbacks[sessionId] = callback
|
||||
orderPushService.subscribeAllEnabled(callback)
|
||||
logger.info("已订阅所有启用账户的订单推送: $sessionId")
|
||||
}
|
||||
else -> {
|
||||
logger.warn("未知的频道: $channel")
|
||||
@@ -120,7 +114,6 @@ class WebSocketSubscriptionService(
|
||||
* 取消订阅
|
||||
*/
|
||||
fun unsubscribe(sessionId: String, channel: String) {
|
||||
logger.info("取消订阅频道: $sessionId -> $channel")
|
||||
|
||||
// 移除订阅关系
|
||||
sessionSubscriptions[sessionId]?.remove(channel)
|
||||
@@ -134,7 +127,6 @@ class WebSocketSubscriptionService(
|
||||
val callback = orderChannelCallbacks.remove(sessionId)
|
||||
if (callback != null) {
|
||||
orderPushService.unsubscribeAll(callback)
|
||||
logger.debug("已取消订阅订单推送: $sessionId -> $channel")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -74,9 +74,6 @@ class PolymarketAuthInterceptor(
|
||||
val signature = generateSignature(signString, apiSecret)
|
||||
|
||||
// 调试日志(仅在 DEBUG 级别输出)
|
||||
logger.debug("L2 认证签名生成: method=$method, path=$requestPath, bodyLength=${bodyString?.length ?: 0}, timestamp=$timestamp")
|
||||
logger.debug("签名字符串: $signString")
|
||||
logger.debug("签名结果: ${signature.take(20)}...")
|
||||
|
||||
// 重新创建请求体(如果原始请求有请求体)
|
||||
val newRequestBody = originalRequest.body?.let { requestBody ->
|
||||
|
||||
@@ -196,7 +196,6 @@ class ResponseLoggingInterceptor : Interceptor {
|
||||
)
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
logger.debug("读取响应体失败: ${e.message}")
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
-13
@@ -39,7 +39,6 @@ class PolymarketWebSocketClient(
|
||||
// 如果启用了代理,配置代理
|
||||
if (proxy != null) {
|
||||
builder.proxy(proxy)
|
||||
logger.info("已配置 WebSocket 代理: ${proxy.address()}")
|
||||
}
|
||||
|
||||
builder.build()
|
||||
@@ -50,7 +49,6 @@ class PolymarketWebSocketClient(
|
||||
*/
|
||||
fun connect() {
|
||||
if (webSocket != null && isConnected) {
|
||||
logger.debug("WebSocket 已连接: $sessionId")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -61,7 +59,6 @@ class PolymarketWebSocketClient(
|
||||
|
||||
webSocket = okHttpClient.newWebSocket(request, object : WebSocketListener() {
|
||||
override fun onOpen(webSocket: WebSocket, response: okhttp3.Response) {
|
||||
logger.info("已成功连接到 Polymarket RTDS: $sessionId")
|
||||
isConnected = true
|
||||
|
||||
// 重置重连延迟(连接成功后重置为初始值)
|
||||
@@ -85,17 +82,14 @@ class PolymarketWebSocketClient(
|
||||
}
|
||||
|
||||
override fun onMessage(webSocket: WebSocket, text: String) {
|
||||
logger.debug("收到 Polymarket 消息: $sessionId, $text")
|
||||
onMessage(text)
|
||||
}
|
||||
|
||||
override fun onMessage(webSocket: WebSocket, bytes: ByteString) {
|
||||
logger.debug("收到 Polymarket 二进制消息: $sessionId")
|
||||
onMessage(bytes.utf8())
|
||||
}
|
||||
|
||||
override fun onClosing(webSocket: WebSocket, code: Int, reason: String) {
|
||||
logger.info("Polymarket 连接正在关闭: $sessionId, code: $code, reason: $reason")
|
||||
isConnected = false
|
||||
stopPing()
|
||||
// 如果不是正常关闭(code != 1000),尝试重连
|
||||
@@ -105,7 +99,6 @@ class PolymarketWebSocketClient(
|
||||
}
|
||||
|
||||
override fun onClosed(webSocket: WebSocket, code: Int, reason: String) {
|
||||
logger.info("Polymarket 连接已关闭: $sessionId, code: $code, reason: $reason")
|
||||
isConnected = false
|
||||
stopPing()
|
||||
// 如果不是正常关闭(code != 1000),尝试重连
|
||||
@@ -138,7 +131,6 @@ class PolymarketWebSocketClient(
|
||||
}
|
||||
})
|
||||
|
||||
logger.info("正在连接到 Polymarket RTDS: $sessionId, URL: $url")
|
||||
} catch (e: Exception) {
|
||||
logger.error("创建 WebSocket 连接失败: $sessionId, ${e.message}", e)
|
||||
throw e
|
||||
@@ -158,7 +150,6 @@ class PolymarketWebSocketClient(
|
||||
if (isConnected) {
|
||||
try {
|
||||
sendMessage("PING")
|
||||
logger.debug("已发送 PING: $sessionId")
|
||||
} catch (e: Exception) {
|
||||
logger.warn("发送 PING 失败: $sessionId, ${e.message}")
|
||||
break
|
||||
@@ -192,17 +183,14 @@ class PolymarketWebSocketClient(
|
||||
|
||||
// 检查是否应该重连
|
||||
if (!shouldReconnect) {
|
||||
logger.info("重连已禁用,停止重连: $sessionId")
|
||||
return@launch
|
||||
}
|
||||
|
||||
// 如果已经连接,不需要重连
|
||||
if (isConnected) {
|
||||
logger.debug("连接已恢复,取消重连: $sessionId")
|
||||
return@launch
|
||||
}
|
||||
|
||||
logger.info("尝试重连 Polymarket WebSocket: $sessionId, 延迟: ${reconnectDelay}ms")
|
||||
|
||||
// 清理旧的连接
|
||||
webSocket = null
|
||||
@@ -239,7 +227,6 @@ class PolymarketWebSocketClient(
|
||||
webSocket?.close(1000, "正常关闭")
|
||||
webSocket = null
|
||||
isConnected = false
|
||||
logger.info("已关闭 WebSocket 连接: $sessionId")
|
||||
} catch (e: Exception) {
|
||||
logger.error("关闭连接失败: $sessionId, ${e.message}", e)
|
||||
}
|
||||
|
||||
-5
@@ -23,7 +23,6 @@ class PolymarketWebSocketHandler : WebSocketHandler {
|
||||
private val polymarketConnections = ConcurrentHashMap<String, PolymarketWebSocketClient>()
|
||||
|
||||
override fun afterConnectionEstablished(session: WebSocketSession) {
|
||||
logger.info("客户端连接建立: ${session.id}")
|
||||
clientSessions[session.id] = session
|
||||
|
||||
try {
|
||||
@@ -42,7 +41,6 @@ class PolymarketWebSocketHandler : WebSocketHandler {
|
||||
// 异步连接,不阻塞
|
||||
try {
|
||||
polymarketClient.connect()
|
||||
logger.info("正在连接到 Polymarket RTDS: ${session.id}")
|
||||
} catch (e: Exception) {
|
||||
logger.error("启动 Polymarket 连接失败: ${e.message}", e)
|
||||
// 连接失败时清理资源
|
||||
@@ -64,7 +62,6 @@ class PolymarketWebSocketHandler : WebSocketHandler {
|
||||
}
|
||||
|
||||
override fun handleMessage(session: WebSocketSession, message: WebSocketMessage<*>) {
|
||||
logger.debug("收到客户端消息: ${session.id}, ${message.payload}")
|
||||
|
||||
val polymarketClient = polymarketConnections[session.id]
|
||||
if (polymarketClient != null) {
|
||||
@@ -89,7 +86,6 @@ class PolymarketWebSocketHandler : WebSocketHandler {
|
||||
}
|
||||
|
||||
override fun afterConnectionClosed(session: WebSocketSession, closeStatus: CloseStatus) {
|
||||
logger.info("客户端连接关闭: ${session.id}, 状态: $closeStatus")
|
||||
cleanup(session.id)
|
||||
}
|
||||
|
||||
@@ -132,7 +128,6 @@ class PolymarketWebSocketHandler : WebSocketHandler {
|
||||
|
||||
// 移除客户端会话
|
||||
clientSessions.remove(sessionId)
|
||||
logger.debug("已清理资源: $sessionId")
|
||||
} catch (e: Exception) {
|
||||
logger.error("清理资源时发生错误: ${sessionId}, ${e.message}", e)
|
||||
}
|
||||
|
||||
@@ -40,19 +40,16 @@ class UnifiedWebSocketHandler(
|
||||
|
||||
@PostConstruct
|
||||
fun init() {
|
||||
logger.info("统一 WebSocket 处理器已初始化,心跳超时: ${heartbeatTimeout}ms")
|
||||
startCleanupTask()
|
||||
}
|
||||
|
||||
@PreDestroy
|
||||
fun destroy() {
|
||||
logger.info("停止统一 WebSocket 处理器")
|
||||
cleanupJob?.cancel()
|
||||
scope.cancel()
|
||||
}
|
||||
|
||||
override fun afterConnectionEstablished(session: WebSocketSession) {
|
||||
logger.info("WebSocket 客户端连接建立: ${session.id}")
|
||||
clientSessions[session.id] = session
|
||||
lastActivityTime[session.id] = System.currentTimeMillis()
|
||||
|
||||
@@ -70,7 +67,6 @@ class UnifiedWebSocketHandler(
|
||||
lastActivityTime[session.id] = System.currentTimeMillis()
|
||||
try {
|
||||
session.sendMessage(TextMessage("PONG"))
|
||||
logger.debug("收到心跳并响应: ${session.id}")
|
||||
} catch (e: Exception) {
|
||||
logger.error("发送心跳响应失败: ${session.id}, ${e.message}", e)
|
||||
}
|
||||
@@ -127,7 +123,6 @@ class UnifiedWebSocketHandler(
|
||||
}
|
||||
|
||||
override fun afterConnectionClosed(session: WebSocketSession, closeStatus: CloseStatus) {
|
||||
logger.info("WebSocket 客户端连接关闭: ${session.id}, 状态: $closeStatus")
|
||||
cleanup(session.id)
|
||||
}
|
||||
|
||||
@@ -166,11 +161,9 @@ class UnifiedWebSocketHandler(
|
||||
try {
|
||||
session.close(CloseStatus.NORMAL)
|
||||
} catch (e: Exception) {
|
||||
logger.debug("关闭会话失败: $sessionId, ${e.message}")
|
||||
}
|
||||
}
|
||||
|
||||
logger.info("已清理 WebSocket 资源: $sessionId")
|
||||
} catch (e: Exception) {
|
||||
logger.error("清理 WebSocket 资源时发生错误: $sessionId, ${e.message}", e)
|
||||
}
|
||||
@@ -212,7 +205,6 @@ class UnifiedWebSocketHandler(
|
||||
}
|
||||
|
||||
if (inactiveSessions.isNotEmpty()) {
|
||||
logger.info("已清理 ${inactiveSessions.size} 个不活跃连接")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user