feat: 优化WebSocket连接和心跳机制
- 在main.py和standalone.py中添加ws_ping_interval和ws_ping_timeout配置 - 调整ws.py中的心跳发送逻辑,先发送ping再等待 - 在host_client中优化消息处理,使用任务队列处理转发请求 - 更新WebTool以适配新的API格式并增加搜索结果限制 - 在agent.py中添加日期显示和web调用次数限制 - 修复bot/handler.py中的事件循环问题
This commit is contained in:
+26
-16
@@ -43,6 +43,7 @@ class NodeClient:
|
||||
self._running = False
|
||||
self._last_heartbeat = time.time()
|
||||
self._reconnect_delay = 1.0
|
||||
self._forward_tasks: set[asyncio.Task] = set()
|
||||
|
||||
async def connect(self) -> bool:
|
||||
"""Connect to the router WebSocket."""
|
||||
@@ -53,9 +54,9 @@ class NodeClient:
|
||||
try:
|
||||
self.ws = await websockets.connect(
|
||||
self.config.router_url,
|
||||
extra_headers=headers,
|
||||
ping_interval=30,
|
||||
ping_timeout=10,
|
||||
additional_headers=headers,
|
||||
ping_interval=20,
|
||||
ping_timeout=60,
|
||||
)
|
||||
logger.info("Connected to router: %s", self.config.router_url)
|
||||
self._reconnect_delay = 1.0
|
||||
@@ -145,17 +146,9 @@ class NodeClient:
|
||||
except Exception as e:
|
||||
logger.error("Failed to send status: %s", e)
|
||||
|
||||
async def handle_message(self, data: str) -> None:
|
||||
"""Handle an incoming message from the router."""
|
||||
try:
|
||||
msg = decode(data)
|
||||
except Exception as e:
|
||||
logger.error("Failed to decode message: %s", e)
|
||||
return
|
||||
|
||||
if isinstance(msg, ForwardRequest):
|
||||
await self.handle_forward(msg)
|
||||
elif isinstance(msg, Heartbeat):
|
||||
async def handle_message_decoded(self, msg: Any) -> None:
|
||||
"""Handle an already-decoded message from the router."""
|
||||
if isinstance(msg, Heartbeat):
|
||||
if msg.type == "ping":
|
||||
if self.ws:
|
||||
try:
|
||||
@@ -165,7 +158,7 @@ class NodeClient:
|
||||
elif msg.type == "pong":
|
||||
self._last_heartbeat = time.time()
|
||||
else:
|
||||
logger.debug("Received message type: %s", msg.type)
|
||||
logger.debug("Received message type: %s", type(msg).__name__)
|
||||
|
||||
async def receive_loop(self) -> None:
|
||||
"""Main receive loop for incoming messages."""
|
||||
@@ -174,7 +167,20 @@ class NodeClient:
|
||||
|
||||
try:
|
||||
async for data in self.ws:
|
||||
await self.handle_message(data)
|
||||
try:
|
||||
msg = decode(data)
|
||||
except Exception as e:
|
||||
logger.error("Failed to decode message: %s", e)
|
||||
continue
|
||||
|
||||
if isinstance(msg, ForwardRequest):
|
||||
# Dispatch as a task so pings are handled without waiting
|
||||
# for the full agent run to complete.
|
||||
task = asyncio.create_task(self.handle_forward(msg))
|
||||
self._forward_tasks.add(task)
|
||||
task.add_done_callback(self._forward_tasks.discard)
|
||||
else:
|
||||
await self.handle_message_decoded(msg)
|
||||
except websockets.ConnectionClosed as e:
|
||||
logger.warning("Connection closed: %s", e)
|
||||
except Exception as e:
|
||||
@@ -243,6 +249,10 @@ class NodeClient:
|
||||
async def stop(self) -> None:
|
||||
"""Stop the client."""
|
||||
self._running = False
|
||||
for task in list(self._forward_tasks):
|
||||
task.cancel()
|
||||
if self._forward_tasks:
|
||||
await asyncio.gather(*self._forward_tasks, return_exceptions=True)
|
||||
if self.ws:
|
||||
await self.ws.close()
|
||||
await manager.stop()
|
||||
|
||||
Reference in New Issue
Block a user