using System.Net;
using System.Security.Cryptography;
using NewLife.Data;
using NewLife.Messaging;
using NewLife.Net;
namespace NewLife.Http;
/// <summary>WebSocket消息处理(新协议栈消息)</summary>
/// <param name="socket"></param>
/// <param name="message">消息(回调内负载完整可用,返回后收尾)</param>
public delegate void WsMessageDelegate(WebSocket socket, WsMessage message);
/// <summary>WebSocket会话管理</summary>
/// <remarks>HTTP 服务端升级后的 WS 会话(由 HttpSession 触发),自持消息泵(MessagePump + WebSocketCodec)解析帧。</remarks>
public class WebSocket : IDisposable
{
#region 属性
/// <summary>是否还在连接</summary>
public Boolean Connected { get; set; }
/// <summary>消息处理器</summary>
/// <remarks>回调内消息负载完整可用(已解码),返回后消息收尾</remarks>
public WsMessageDelegate? MessageHandler { get; set; }
/// <summary>Http上下文</summary>
public IHttpContext? Context { get; set; }
/// <summary>版本</summary>
public String? Version { get; set; }
/// <summary>协议。如mqtt</summary>
public String? Protocol { get; set; }
/// <summary>活跃时间</summary>
public DateTime ActiveTime { get; set; }
private Pipe? _pipe;
/// <summary>消息编解码器(服务端角色:接收带掩码帧、发送无掩码)。无状态,可跨会话共享</summary>
private static readonly WebSocketCodec _codec = new() { IsServer = true };
/// <summary>帧泵。整帧模式:同步泵只能消费整帧,帧未完整留待下一轮</summary>
/// <remarks>每实例一份:整帧模式的帧长上限(<see cref="MaxFrameSize"/>)随会话配置,不能跨连接共享</remarks>
private MessagePump? _pump;
/// <summary>单帧长度上限,默认 16M。0 表示不限制</summary>
/// <remarks>整帧解析要求整个帧驻留内存,本上限是单连接的内存安全阀。首次帧泵启动时随配置定格,运行中修改不再生效</remarks>
public Int32 MaxFrameSize { get; set; } = 16 * 1024 * 1024;
/// <summary>分片重组器(RFC 6455 §5.4)。数据帧 FIN=0 累积,末片合并成完整消息后交付</summary>
private readonly WebSocketFragment _fragment = new();
#endregion
#region 方法
/// <summary>WebSocket 握手</summary>
/// <param name="context"></param>
/// <returns></returns>
public static WebSocket? Handshake(IHttpContext context)
{
var request = context.Request;
if (!request.Headers.TryGetValue("Sec-WebSocket-Key", out var key) || key.IsNullOrEmpty()) return null;
var manager = new WebSocket();
// 校验不通过必须返回 null:否则调用方(HttpSession)会把普通请求长期当 WebSocket 接管,
// 既不给客户端任何握手响应,也让该连接永远脱离 HTTP 解析
if (!manager.ProcessRequest(context)) return null;
return manager;
}
/// <summary>处理 WebSocket 握手</summary>
/// <param name="context"></param>
/// <remarks>按 RFC 6455 §4.2 校验握手四要素:只看 Sec-WebSocket-Key 会让任意路径的普通请求也被升级为 WebSocket</remarks>
public Boolean ProcessRequest(IHttpContext context)
{
var request = context.Request;
if (!request.Headers.TryGetValue("Sec-WebSocket-Key", out var key) || key.IsNullOrEmpty()) return false;
var upgrade = request.Headers["Upgrade"];
if (!HasToken(upgrade, "websocket")) return false;
// Connection 可能形如 “keep-alive, Upgrade”:按逗号拆分的令牌列表判断,不是子串匹配
var connection = request.Headers["Connection"];
if (!HasToken(connection, "Upgrade")) return false;
// 仅支持 RFC 6455(版本 13)。客户端可能以列表形式声明多个版本,只要含 13 即接受
if (!HasToken(request.Headers["Sec-WebSocket-Version"], "13")) return false;
var buf = SHA1.Create().ComputeHash((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").GetBytes());
key = buf.ToBase64();
var response = context.Response;
response.StatusCode = HttpStatusCode.SwitchingProtocols;
response.Headers["Upgrade"] = "websocket";
response.Headers["Connection"] = "Upgrade";
response.Headers["Sec-WebSocket-Accept"] = key;
if (context is DefaultHttpContext dhc) dhc.WebSocket = this;
if (!Protocol.IsNullOrEmpty())
response.Headers["Sec-WebSocket-Protocol"] = Protocol;
if (!Version.IsNullOrEmpty())
response.Headers["Sec-WebSocket-Version"] = Version;
//if (request.Headers.TryGetValue("Sec-WebSocket-Version", out var ver)) Version = ver;
Context = context;
Connected = true;
ActiveTime = DateTime.Now;
return true;
}
/// <summary>头部值是否含指定令牌。头部可能是逗号分隔列表(如 “keep-alive, Upgrade”),逐项 Trim 后不区分大小写比较</summary>
/// <param name="header">头部值</param>
/// <param name="token">目标令牌</param>
/// <returns>是否含该令牌</returns>
private static Boolean HasToken(String header, String token)
{
if (header.IsNullOrEmpty()) return false;
foreach (var item in header.Split(','))
{
if (item.Trim().EqualIgnoreCase(token)) return true;
}
return false;
}
/// <summary>处理WebSocket数据包。数据进入数据管道,逐帧同步泵出完整帧交给消息处理,支持跨接收边界的粘包/分包</summary>
/// <param name="pk">已到达的原始数据包,可能包含零个或多个完整 WebSocket 帧</param>
/// <remarks>帧未完整时残片保留在管道内等下一轮;入参为借阅视图时自动转为自有拷贝(跨轮安全)</remarks>
public void Process(IPacket pk)
{
// 数据进入管道:共享切片保持调用方句柄不受影响;借阅视图不能跨轮保留,转为自有拷贝
var node = pk.Slice(0, -1);
if (node is not OwnerPacket op || op.RefCount == 0) node = node.Clone();
_pipe ??= new Pipe();
_pipe.Writer.Append(node);
// 同步泵:当前缓冲内可成的整帧全部处理;头部不足或帧未完整则留给下一轮
var pump = _pump ??= new MessagePump(_codec) { RequireFullFrame = true, MaxFrameSize = MaxFrameSize };
try
{
ReadFrames(pump);
}
catch (FrameTooLargeException ex)
{
// 单帧超限:RFC 6455 §7.4.1 用 1009(Message Too Big),与帧损坏(1002)区分
NewLife.Log.XTrace.WriteLine("WebSocket 帧超限,关闭连接:{0}", ex.Message);
FailConnection(1009, "message too big");
}
catch (InvalidOperationException ex)
{
// 协议帧损坏:残片无法恢复,留在管道会让后续每个包重复解析同一坏帧,直接关闭连接
NewLife.Log.XTrace.WriteLine("WebSocket 帧损坏,关闭连接:{0}", ex.Message);
FailConnection(1002, "protocol error");
}
}
/// <summary>同步泵出并向业务交付管道内已经可成的整帧;帧未完整则留在管道等下一轮</summary>
/// <param name="pump">帧泵</param>
private void ReadFrames(MessagePump pump)
{
var reader = _pipe!.Reader;
// 读取器已随连接关闭完成时直接收尾:结束后再读会抛异常,不应当成“帧损坏”去关连接
while (!reader.IsReaderCompleted && pump.TryRead(reader, out var message))
{
try
{
if (message is not WsMessage ws) continue;
// RFC 6455 §5.1:客户端发给服务端的每一帧都必须带掩码(防中间设备缓存投毒),未掩码帧按协议错误关闭连接
if (ws.MaskKey == null) throw new InvalidOperationException("客户端帧必须带掩码(RFC 6455 §5.1)");
// RFC 6455 §5.5:控制帧(Close/Ping/Pong)必须 FIN=1 且负载不超过 125 字节。
// 不校验会让对端用一个超大 Ping 触发等量 Pong 回显(反射放大)
if (ws.Type is WebSocketMessageType.Close or WebSocketMessageType.Ping or WebSocketMessageType.Pong
&& (!ws.Fin || ws.Payload?.Total > 125))
throw new InvalidOperationException("控制帧必须 FIN=1 且负载不超过 125 字节(RFC 6455 §5.5)");
// 客户端帧带掩码:整帧路径对负载原地解码(帧内字节独享)
ws.Demask();
// 分片重组:数据帧 FIN=0 累积,续片追加,末片合并成完整消息后交付;控制帧直通
if (ws.Type is WebSocketMessageType.Text or WebSocketMessageType.Binary && !ws.Fin)
{
_fragment.Begin(ws.Type, ws.Payload);
if (_fragment.TooBig) { FailConnection(1009, "message too big"); return; }
continue;
}
if (ws.Type == WebSocketMessageType.Data)
{
var whole = _fragment.Append(ws.Fin, ws.Payload);
// 无首片的孤立续片(分片序列未开始或已结束):RFC 6455 §5.4 属协议错误,
// 静默忽略会让对端以为已送达
if (whole == null && !_fragment.Active)
{
if (_fragment.TooBig) { FailConnection(1009, "message too big"); return; }
FailConnection(1002, "unexpected continuation frame");
return;
}
if (whole != null)
{
try
{
// RFC 6455 §8.1:文本消息必须是合法 UTF-8(分片消息在重组完成后校验,多字节序列可跨片)
if (whole.Type == WebSocketMessageType.Text && !WebSocketCodec.IsValidUtf8(whole.Payload))
{
FailConnection(1007, "invalid utf-8");
return;
}
Process(whole);
}
finally
{
// 与单帧路径同口径:回调返回后消息收尾
whole.TryDispose();
}
}
if (_fragment.TooBig) { FailConnection(1009, "message too big"); return; }
continue;
}
// RFC 6455 §8.1:文本消息的负载必须是合法 UTF-8,畸变数据按 1007 失败连接(不能让业务拿到乱码)
if (ws.Type == WebSocketMessageType.Text && !WebSocketCodec.IsValidUtf8(ws.Payload))
{
FailConnection(1007, "invalid utf-8");
return;
}
// 消息化入口:业务回调直达(负载已解码),协议帧处理内联
Process(ws);
}
finally
{
message.TryDispose();
}
}
}
/// <summary>失败连接:发关闭帧(对端可感知原因)后释放连接</summary>
/// <param name="closeStatus">关闭状态码,如 1002 协议错误、1009 消息过大</param>
/// <param name="reason">关闭描述,随关闭帧下发</param>
private void FailConnection(Int32 closeStatus, String reason)
{
NewLife.Log.XTrace.WriteLine("WebSocket 失败连接 {0}:{1}", closeStatus, reason);
try { Close(closeStatus, reason); } catch { }
Context?.Connection?.TryDispose();
Context?.Socket?.TryDispose();
Connected = false;
}
/// <summary>处理WebSocket消息</summary>
/// <param name="message">消息(负载已解码;回调内完整可用,返回后收尾)</param>
public void Process(WsMessage message)
{
ActiveTime = DateTime.Now;
// Close 帧:从负载解析状态码与描述(在业务回调前,供回调与回显使用)
if (message.Type == WebSocketMessageType.Close) message.TryReadCloseStatus();
// 业务回调:负载完整可用
MessageHandler?.Invoke(this, message);
// 协议帧处理:Close 回显关闭 / Ping 回显 Pong
var session = Context?.Connection;
var socket = Context?.Socket;
if (session == null && socket == null) return;
switch (message.Type)
{
case WebSocketMessageType.Close:
{
// RFC 6455 §5.5.1/§7.4.1:有状态码则回显;无状态码时发空负载关闭帧。
// 1005/1006/1015 是保留值,0~999 与 1016~2999 未分配,都禁止出现在线路上:
// 收到这类码本身就是协议错误,回显等于把违规码原样发回,正解是按协议错误失败连接
if (WebSocketCodec.IsSendableCloseStatus(message.CloseStatus))
Close(message.CloseStatus, message.StatusDescription ?? "Finished");
else if (message.CloseStatus > 0)
{
FailConnection(1002, "invalid close code");
return;
}
else
Close();
session?.Dispose();
socket?.Dispose();
Connected = false;
}
break;
case WebSocketMessageType.Ping:
{
// RFC 6455 §5.5.3:Pong 必须回传 Ping 的 Application Data。
// 共享切片(引用计数各自释放):Pong 帧持有独立句柄,Ping 消息与其负载均不受影响
var pong = new WsMessage { Type = WebSocketMessageType.Pong };
try
{
var payload = message.Payload;
if (payload != null) pong.SetBody(payload is IOwnerPacket owner ? owner.Slice(0, -1) : payload);
Send(pong);
}
finally
{
// 容器随发送结束释放:否则每收到一个 Ping 就多一份接收缓冲引用永不归还
pong.TryDispose();
}
}
break;
}
// 负载不在此释放:所有权随消息容器(帧泵回调 using / 调用方负责)
}
private void Send(WsMessage msg)
{
var session = Context?.Connection;
var socket = Context?.Socket;
if (session == null && socket == null) throw new ObjectDisposedException(nameof(Context));
var data = _codec.Build(msg)!;
if (session != null)
session.Send(data);
else
socket?.Send(data);
data.TryDispose();
}
/// <summary>发送消息</summary>
/// <param name="data">负载。借用语义:调用方保留句柄并自行释放</param>
/// <param name="type"></param>
public void Send(IPacket data, WebSocketMessageType type)
{
var ws = new WsMessage { Type = type };
try
{
// 借用:拥有句柄按引用计数共享给容器(零拷贝),其余视图转自有拷贝;调用方句柄始终有效
ws.SetBody(data is OwnerPacket op && op.RefCount > 0 ? op.Slice(0, -1) : data.Clone());
Send(ws);
}
finally
{
// 归还容器自己那份负载引用
ws.TryDispose();
}
}
/// <summary>发送消息</summary>
/// <param name="data"></param>
/// <param name="type"></param>
public void Send(Byte[] data, WebSocketMessageType type) => Send((ArrayPacket)data, type);
/// <summary>发送文本消息</summary>
/// <param name="message"></param>
public void Send(String message) => Send(message.GetBytes(), WebSocketMessageType.Text);
/// <summary>向所有连接发送消息</summary>
/// <param name="data"></param>
/// <param name="type"></param>
/// <param name="predicate"></param>
/// <returns>已群发客户端总数</returns>
public async Task<Int32> SendAllAsync(IPacket data, WebSocketMessageType type, Func<INetSession, Boolean>? predicate = null)
{
var session = (Context?.Connection) ?? throw new ObjectDisposedException(nameof(Context));
var ws = new WsMessage { Type = type };
try
{
// 借用:负载引用共享给容器,调用方句柄保持有效
ws.SetBody(data is OwnerPacket op && op.RefCount > 0 ? op.Slice(0, -1) : data.Clone());
var data2 = _codec.Build(ws)!;
try
{
// 经服务端对各会话并行送出,等待完成后再归还封包(封包持有负载引用,释放封包即归还整链)
return await session.Host.SendAllAsync(data2, predicate).ConfigureAwait(false);
}
finally
{
data2.TryDispose();
}
}
finally
{
ws.TryDispose();
}
}
/// <summary>向所有连接发送文本消息</summary>
/// <param name="message"></param>
/// <param name="predicate"></param>
/// <returns>已群发客户端总数</returns>
public Task<Int32> SendAllAsync(String message, Func<INetSession, Boolean>? predicate = null) => SendAllAsync((ArrayPacket)message.GetBytes(), WebSocketMessageType.Text, predicate);
/// <summary>发送关闭连接</summary>
/// <param name="closeStatus"></param>
/// <param name="statusDescription"></param>
public void Close(Int32 closeStatus, String statusDescription)
{
var ws = new WsMessage { Type = WebSocketMessageType.Close };
try
{
ws.SetBody(WebSocketCodec.BuildClosePayload(closeStatus, statusDescription));
Send(ws);
}
finally
{
ws.TryDispose();
}
}
/// <summary>发送空负载的关闭帧(对端未带状态码时用)。RFC 6455 §7.4.1 禁止把 1005/1006 等保留值发到线上</summary>
public void Close()
{
var ws = new WsMessage { Type = WebSocketMessageType.Close };
try
{
Send(ws);
}
finally
{
ws.TryDispose();
}
}
#endregion
#region 销毁
/// <summary>销毁。归还数据管道的段链缓冲(连接结束时调用)</summary>
public void Dispose()
{
_pipe?.Dispose();
_pipe = null;
}
#endregion
}
|