using System.Net.WebSockets; using CoreAgent.WebSocketTransport.Interfaces; using Microsoft.Extensions.Logging; namespace CoreAgent.WebSocketTransport.Services; /// /// WebSocket 连接实现 /// 单一职责:负责网络连接管理 /// public class WebSocketConnection : IWebSocketConnection { private ClientWebSocket _webSocket; private readonly ILogger _logger; private readonly object _lock = new object(); public WebSocketState State => _webSocket.State; public bool IsConnected => _webSocket.State == WebSocketState.Open; public WebSocketConnection(ILogger logger) { _webSocket = new ClientWebSocket(); _logger = logger; } public async Task ConnectAsync(Uri uri, CancellationToken cancellationToken = default) { lock (_lock) { // 如果连接已关闭或处于CloseReceived状态,需要重新创建WebSocket实例 if (_webSocket.State == WebSocketState.Closed || _webSocket.State == WebSocketState.CloseReceived || _webSocket.State == WebSocketState.CloseSent) { _logger.LogInformation("WebSocket 处于关闭状态,重新创建连接实例"); _webSocket?.Dispose(); _webSocket = new ClientWebSocket(); } if (_webSocket.State == WebSocketState.Open) { _logger.LogInformation("WebSocket 已连接"); return; } } _logger.LogInformation("正在连接 WebSocket 服务器: {Uri}", uri); // 👇 这行:跳过所有证书验证(只用于开发测试) _webSocket.Options.RemoteCertificateValidationCallback += (sender, cert, chain, sslPolicyErrors) => true; await _webSocket.ConnectAsync(uri, cancellationToken); _logger.LogInformation("WebSocket 连接成功"); } public async Task SendAsync(ArraySegment buffer, WebSocketMessageType messageType, bool endOfMessage, CancellationToken cancellationToken = default) { if (_webSocket.State != WebSocketState.Open) { throw new InvalidOperationException($"WebSocket 未连接,当前状态: {_webSocket.State}"); } await _webSocket.SendAsync(buffer, messageType, endOfMessage, cancellationToken); } public async Task ReceiveAsync(ArraySegment buffer, CancellationToken cancellationToken = default) { if (_webSocket.State != WebSocketState.Open) { throw new InvalidOperationException($"WebSocket 未连接,当前状态: {_webSocket.State}"); } return await _webSocket.ReceiveAsync(buffer, cancellationToken); } public async Task CloseAsync(WebSocketCloseStatus closeStatus, string statusDescription, CancellationToken cancellationToken = default) { if (_webSocket.State == WebSocketState.Open) { await _webSocket.CloseAsync(closeStatus, statusDescription, cancellationToken); } } /// /// 强制关闭连接并重新创建WebSocket实例 /// public void ForceClose() { lock (_lock) { _logger.LogInformation("强制关闭 WebSocket 连接并重新创建实例"); _webSocket?.Dispose(); _webSocket = new ClientWebSocket(); } } public void Dispose() { lock (_lock) { _webSocket?.Dispose(); } GC.SuppressFinalize(this); } }