ASP.NET Core 6 WebSocket代理中间件实现方案咨询
ASP.NET Core 6 WebSocket 代理实现方案
核心思路
基于ASP.NET Core中间件实现WebSocket代理,拦截客户端的WebSocket握手请求,完成与后端目标服务器的握手,之后双向转发并处理消息,全程无需修改客户端代码,仅需让客户端连接到代理服务的URL。
避免直接操作Socket的坑
直接捕获Socket转发崩溃,大概率是因为没正确处理WebSocket的握手协议、帧格式,以及异步流的同步问题。推荐使用ASP.NET Core内置的WebSocket对象处理,而非直接操作底层Socket。
完整实现示例
1. 代理中间件代码
using System.Net.WebSockets; using System.Text; namespace WebSocketProxy.Middleware; public class WebSocketProxyMiddleware { private readonly RequestDelegate _next; private readonly ILogger<WebSocketProxyMiddleware> _logger; private readonly string _targetServerUrl; public WebSocketProxyMiddleware(RequestDelegate next, ILogger<WebSocketProxyMiddleware> logger, IConfiguration config) { _next = next; _logger = logger; _targetServerUrl = config.GetValue<string>("WebSocket:TargetServerUrl") ?? throw new ArgumentNullException("TargetServerUrl未配置"); } public async Task InvokeAsync(HttpContext context) { // 仅处理WebSocket握手请求 if (!context.WebSockets.IsWebSocketRequest) { await _next(context); return; } // 1. 接受客户端WebSocket连接 var clientWebSocket = await context.WebSockets.AcceptWebSocketAsync(); _logger.LogInformation("客户端WebSocket连接已建立"); // 2. 连接到目标服务器的WebSocket端点 using var targetWebSocket = new ClientWebSocket(); try { // 复制客户端的WebSocket子协议(如果有) if (!string.IsNullOrEmpty(context.WebSockets.WebSocketRequestedProtocols.FirstOrDefault())) { targetWebSocket.Options.AddSubProtocol(context.WebSockets.WebSocketRequestedProtocols.First()); } await targetWebSocket.ConnectAsync(new Uri(_targetServerUrl), context.RequestAborted); _logger.LogInformation("已连接到目标WebSocket服务器"); // 3. 启动双向消息转发任务 var receiveFromClientTask = ForwardClientToTarget(clientWebSocket, targetWebSocket); var receiveFromTargetTask = ForwardTargetToClient(clientWebSocket, targetWebSocket); // 等待任意一方连接关闭 var completedTask = await Task.WhenAny(receiveFromClientTask, receiveFromTargetTask); if (completedTask.IsFaulted) { _logger.LogError(completedTask.Exception, "消息转发过程中出现错误"); } } catch (Exception ex) { _logger.LogError(ex, "代理连接建立失败"); // 关闭客户端连接 if (clientWebSocket.State == WebSocketState.Open) { await clientWebSocket.CloseAsync(WebSocketCloseStatus.InternalServerError, "代理连接失败", context.RequestAborted); } } finally { // 确保资源释放 if (clientWebSocket.State == WebSocketState.Open) { await clientWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "连接关闭", CancellationToken.None); } clientWebSocket.Dispose(); } } /// <summary> /// 转发客户端消息到目标服务器,同时处理请求数据 /// </summary> private async Task ForwardClientToTarget(WebSocket clientSocket, WebSocket targetSocket) { var buffer = new byte[4096]; while (clientSocket.State == WebSocketState.Open && targetSocket.State == WebSocketState.Open) { var result = await clientSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None); if (result.MessageType == WebSocketMessageType.Close) { await targetSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None); break; } // 处理客户端发送的消息(示例:修改文本消息) var processedBuffer = ProcessClientMessage(buffer, result.Count); await targetSocket.SendAsync(new ArraySegment<byte>(processedBuffer, 0, processedBuffer.Length), result.MessageType, result.EndOfMessage, CancellationToken.None); } } /// <summary> /// 转发目标服务器消息到客户端,同时处理响应数据 /// </summary> private async Task ForwardTargetToClient(WebSocket clientSocket, WebSocket targetSocket) { var buffer = new byte[4096]; while (clientSocket.State == WebSocketState.Open && targetSocket.State == WebSocketState.Open) { var result = await targetSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None); if (result.MessageType == WebSocketMessageType.Close) { await clientSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None); break; } // 处理目标服务器返回的消息(示例:修改文本消息) var processedBuffer = ProcessTargetMessage(buffer, result.Count); await clientSocket.SendAsync(new ArraySegment<byte>(processedBuffer, 0, processedBuffer.Length), result.MessageType, result.EndOfMessage, CancellationToken.None); } } /// <summary> /// 自定义客户端请求消息处理逻辑 /// </summary> private byte[] ProcessClientMessage(byte[] buffer, int count) { // 示例:如果是文本消息,添加代理标识 var message = Encoding.UTF8.GetString(buffer, 0, count); var processedMessage = $"[Proxy Processed] {message}"; return Encoding.UTF8.GetBytes(processedMessage); // 如果是二进制消息,直接返回或自定义二进制处理逻辑 // return buffer.Take(count).ToArray(); } /// <summary> /// 自定义目标服务器响应消息处理逻辑 /// </summary> private byte[] ProcessTargetMessage(byte[] buffer, int count) { // 示例:如果是文本消息,添加代理标识 var message = Encoding.UTF8.GetString(buffer, 0, count); var processedMessage = $"[Proxy Response] {message}"; return Encoding.UTF8.GetBytes(processedMessage); // 二进制消息处理同理 // return buffer.Take(count).ToArray(); } }
2. 注册中间件
在Program.cs中注册代理中间件:
var builder = WebApplication.CreateBuilder(args); // 添加配置(可以在appsettings.json中配置目标服务器地址) builder.Configuration.AddJsonFile("appsettings.json", optional: false, reloadOnChange: true); // 添加日志服务 builder.Services.AddLogging(); var app = builder.Build(); // 注册WebSocket代理中间件,放在其他中间件之前 app.UseMiddleware<WebSocketProxy.Middleware.WebSocketProxyMiddleware>(); app.Run();
3. 配置目标服务器地址
在appsettings.json中添加:
{ "WebSocket": { "TargetServerUrl": "ws://localhost:5000/your-websocket-endpoint" } }
关键注意事项
- 必须处理WebSocket的关闭帧,否则会导致连接泄漏或异常
- 消息转发要使用异步任务并行处理,避免阻塞
- 子协议要从客户端请求中复制并传递给目标服务器,确保兼容性
- 异常处理要覆盖连接建立、消息转发的全流程,避免崩溃
- 二进制消息和文本消息要分别处理,不要混用编码
内容的提问来源于stack exchange,提问作者DaCh
相关产品推荐
相关产品推荐

