From e072d0c894713c660a97744080e49f23bc6d1469 Mon Sep 17 00:00:00 2001 From: WrBug Date: Tue, 13 Jan 2026 16:04:20 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=B7=BB=E5=8A=A0=20Activity=20WebSock?= =?UTF-8?q?et=20=E6=B6=88=E6=81=AF=E8=B6=85=E6=97=B6=E6=A3=80=E6=B5=8B?= =?UTF-8?q?=E5=92=8C=E8=87=AA=E5=8A=A8=E9=87=8D=E8=BF=9E=E6=9C=BA=E5=88=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 添加 lastActivityTime 变量记录最后一次收到 activity 消息的时间 - 实现 startActivityTimeoutCheck() 方法,每30秒检查一次消息接收情况 - 如果超过30秒未收到 activity 消息,自动触发 WebSocket 重连 - 在订阅成功后自动启动超时检测任务 - 在 stop() 方法中正确清理检测任务资源 --- .../monitor/PolymarketActivityWsService.kt | 69 ++++++++++++++++++- 1 file changed, 66 insertions(+), 3 deletions(-) diff --git a/backend/src/main/kotlin/com/wrbug/polymarketbot/service/copytrading/monitor/PolymarketActivityWsService.kt b/backend/src/main/kotlin/com/wrbug/polymarketbot/service/copytrading/monitor/PolymarketActivityWsService.kt index cb094d4..b2d8989 100644 --- a/backend/src/main/kotlin/com/wrbug/polymarketbot/service/copytrading/monitor/PolymarketActivityWsService.kt +++ b/backend/src/main/kotlin/com/wrbug/polymarketbot/service/copytrading/monitor/PolymarketActivityWsService.kt @@ -7,11 +7,11 @@ import com.wrbug.polymarketbot.entity.Leader import com.wrbug.polymarketbot.repository.LeaderRepository import com.wrbug.polymarketbot.service.copytrading.statistics.CopyOrderTrackingService import com.wrbug.polymarketbot.util.fromJson +import com.wrbug.polymarketbot.constants.PolymarketConstants import com.wrbug.polymarketbot.websocket.PolymarketWebSocketClient import jakarta.annotation.PreDestroy import kotlinx.coroutines.* import org.slf4j.LoggerFactory -import org.springframework.beans.factory.annotation.Value import org.springframework.stereotype.Service import java.math.BigDecimal import java.util.concurrent.ConcurrentHashMap @@ -29,8 +29,7 @@ class PolymarketActivityWsService( private val logger = LoggerFactory.getLogger(PolymarketActivityWsService::class.java) - @Value("\${polymarket.websocket.activity.url:wss://ws-live-data.polymarket.com}") - private var websocketUrl: String = "wss://ws-live-data.polymarket.com" + private val websocketUrl: String = PolymarketConstants.ACTIVITY_WS_URL private val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) @@ -44,6 +43,13 @@ class PolymarketActivityWsService( @Volatile private var isSubscribed = false + // 最后一次收到 activity 消息的时间(毫秒时间戳) + @Volatile + private var lastActivityTime: Long = 0 + + // Activity 消息超时检测任务 + private var activityTimeoutJob: Job? = null + /** * 启动监听 */ @@ -184,6 +190,10 @@ class PolymarketActivityWsService( client.sendMessage(subscribeMessage) isSubscribed = true + // 重置最后一次收到 activity 消息的时间 + lastActivityTime = System.currentTimeMillis() + // 启动 Activity 消息超时检测 + startActivityTimeoutCheck() logger.info("Activity WebSocket 订阅成功(全局交易流)") } catch (e: Exception) { logger.error("订阅 Activity WebSocket 失败", e) @@ -191,6 +201,54 @@ class PolymarketActivityWsService( } } + /** + * 启动 Activity 消息超时检测 + * 每30秒检查一次,如果超过30秒没有收到activity消息,则重连 + */ + private fun startActivityTimeoutCheck() { + // 先停止之前的检测任务 + stopActivityTimeoutCheck() + + activityTimeoutJob = scope.launch { + while (isActive && isSubscribed) { + delay(30000) // 每30秒检查一次 + + // 如果已经取消订阅,停止检测 + if (!isSubscribed) { + break + } + + // 如果 lastActivityTime 为 0,说明还没有收到过消息,跳过本次检测 + if (lastActivityTime == 0L) { + continue + } + + val currentTime = System.currentTimeMillis() + val timeSinceLastActivity = currentTime - lastActivityTime + + // 如果超过30秒没有收到activity消息,触发重连 + if (timeSinceLastActivity >= 30000) { + logger.warn("超过30秒未收到 Activity 消息,触发重连。距离上次消息: ${timeSinceLastActivity}ms") + // 关闭当前连接并重连 + wsClient?.closeConnection() + wsClient = null + isSubscribed = false + // 重新连接 + connectAndSubscribe() + break // 重连后会重新启动检测任务 + } + } + } + } + + /** + * 停止 Activity 消息超时检测 + */ + private fun stopActivityTimeoutCheck() { + activityTimeoutJob?.cancel() + activityTimeoutJob = null + } + /** * 处理消息 */ @@ -214,6 +272,9 @@ class PolymarketActivityWsService( return } + // 更新最后一次收到 activity 消息的时间(即使不是我们监听的 Leader 的交易) + lastActivityTime = System.currentTimeMillis() + val payload = tradeMessage.payload // 提取交易者地址 @@ -383,10 +444,12 @@ class PolymarketActivityWsService( */ fun stop() { logger.info("停止 Activity WebSocket 监听") + stopActivityTimeoutCheck() wsClient?.closeConnection() wsClient = null isSubscribed = false monitoredAddresses.clear() + lastActivityTime = 0 } /**