using System.Buffers;
using NewLife.Data;
using NewLife.Http;
namespace NewLife.Messaging;
/// <summary>WebSocket 消æ¯ç¼–è§£ç 器(RFC 6455 å¸§æ ¼å¼ï¼‰ã€‚å¸§æ ¼å¼ï¼šFIN/OPCODE + 长度(1/2/8 å—节大端)+ [掩ç 4 å—节] + è´Ÿè½½</summary>
/// <remarks>
/// <para>æ— çŠ¶æ€ã€å¯è·¨è¿žæŽ¥å…±äº«ï¼›<see cref="IsServer"/> åªå†³å®š<b>å‘é€</b>æ–¹å‘的掩ç (æœåŠ¡ç«¯ä¸åŠ æŽ©ç ã€å®¢æˆ·ç«¯è‡ªåŠ¨åŠ éšæœºæŽ©ç ),解æžä¾§å¯¹å¸¦æŽ©ç 与ä¸å¸¦æŽ©ç 的帧都接å—ï¼Œç”±æ¶ˆè´¹æ–¹ä¾æ® <see cref="WsMessage.MaskKey"/> 决定是å¦è§£ç 。</para>
/// <para><b>掩ç è§£ç </b>:解æžäº§å‡ºçš„ <see cref="WsMessage.MaskKey"/> éžç©ºæ—¶ï¼Œæ¶ˆè´¹æ–¹é¡»å¯¹è´Ÿè½½æŒ‰æŽ©ç è§£ç (æ¯å—节 XOR 密钥,<c>data[i] ^= key[i % 4]</c>,链å¼è´Ÿè½½è·¨æ®µè¿žç»ï¼‰ã€‚整帧路径å¯åŽŸåœ°è§£ç ï¼›æµå¼è·¯å¾„建议物化åŽè§£ç 。<b>æ¤æ¡åªçº¦æŸç›´æŽ¥æ¶ˆè´¹æœ¬ codec 的调用方,且解ç åªå¯æ‰§è¡Œä¸€æ¬¡</b>(<see cref="WsMessage.Demask"/> éžå¹‚ç‰ï¼‰ï¼›ç» <c>Http/WebSocket</c>(æœåŠ¡ç«¯ï¼‰ä¸Ž <c>WebSocketClient</c>(客户端)交付的消æ¯ï¼Œæ¡†æž¶å·²åœ¨äº¤ä»˜å‰å®Œæˆè§£ç ï¼Œæ¤æ—¶ <see cref="WsMessage.MaskKey"/> 仅作诊æ–观察,业务侧ä¸å¾—冿¬¡è§£ç 。</para>
/// <para><b>ä¸è¦ç›´æŽ¥ç”¨ä½œ <c>SessionBase.Protocol</c></b>:本 codec ä¸åšæŽ©ç è§£ç ,会è¯å±‚ä¹Ÿæ²¡æœ‰è§£ç æ—¶æœºï¼ˆè´Ÿè½½å¯èƒ½è¿˜æ˜¯æµå¼çš„),
/// 直接装é…会把未解ç 的掩ç è´Ÿè½½é™é»˜äº¤ç»™ä¸šåŠ¡ã€‚WebSocket 通é“请走 <c>Http/WebSocket</c>(æœåŠ¡ç«¯ï¼‰ä¸Ž <c>WebSocketClient</c>(客户端),
/// 两者在交付å‰å®Œæˆè§£ç 与分片é‡ç»„。</para>
/// <para><b>分片</b>:FIN=0 的帧照常解æžï¼ˆFin=false),分片é‡ç»„由消费侧(WebSocket æœåŠ¡ç«¯ / WebSocketClient)累积完æˆã€‚</para>
/// </remarks>
/// <example>
/// <code>
/// // æœåŠ¡ç«¯ï¼šæŽ¥æ”¶å¸¦æŽ©ç 的客户端帧ã€å‘逿— 掩ç 帧
/// var codec = new WebSocketCodec { IsServer = true };
///
/// if (codec.TryParse(buffer) is { } rs) { var msg = rs.Message; }
/// </code>
/// </example>
public class WebSocketCodec : IMessageCodec
{
#region 属性
/// <summary>æ˜¯å¦æœåŠ¡ç«¯ã€‚é»˜è®¤ true:接收期望掩ç 帧ã€å‘é€ä¸åŠ æŽ©ç ;客户端置 false:å‘é€è‡ªåŠ¨åŠ æŽ©ç </summary>
public Boolean IsServer { get; set; } = true;
#endregion
#region 方法
/// <summary>å®šç•Œå¹¶æž„é€ æ¶ˆæ¯ã€‚è§£æžå¸§å¤´ï¼ˆFIN/OPCODE/长度/掩ç é”®ï¼‰ï¼Œæž„é€ <see cref="WsMessage"/></summary>
/// <param name="buffer">帧首窗å£ï¼ˆåªè¯»åºåˆ—,å¯è·¨æ®µï¼‰</param>
/// <returns>头部ä¸è¶³è¿”回 nullï¼ˆä¸æ¶ˆè´¹ã€ä¸äº§ç”Ÿå¯¹è±¡ï¼‰ï¼›é•¿åº¦éžæ³•返回 <see cref="ParseResult.Invalid"/>(帧æŸåé¡»ç«‹å³æŠ¥é”™ï¼‰ã€‚FIN=0 的分片帧照常产出,Fin å—æ®µæ ‡è¯†</returns>
/// <remarks>在åªè¯»åºåˆ—上顺åºè¯»å–ï¼Œä¸æ‹¼è¯»ã€ä¸ç‰©åŒ–;掩ç 键挂在消æ¯ä¸Šï¼Œè´Ÿè½½ç”±æ¶ˆè´¹æ–¹è§£ç 。</remarks>
public ParseResult? TryParse(ReadOnlySequence<Byte> buffer)
{
// 叧嗿®µç”±æ¶ˆæ¯ç±»è‡ªè¡Œè§£æžï¼ˆæ¶ˆæ¯å®šä¹‰å³å议),帧层åªè´Ÿè´£è£…é…。
// å…ˆåšæ— å‰¯ä½œç”¨çš„å¤´éƒ¨æŽ¢æµ‹ï¼šå¸§å¤´æœªåˆ°é½æ—¶ç›´æŽ¥è¿”å›žï¼Œä¸æž„é€ éšå³è¢«ä¸¢å¼ƒçš„æ¶ˆæ¯å¯¹è±¡
if (!WsMessage.TryReadHeader(buffer, out var fin, out var type, out var maskKey, out var bodyLength, out var headerSize, out var invalid))
return invalid ? new ParseResult { Invalid = true } : null;
var message = new WsMessage { Fin = fin, Type = type, MaskKey = maskKey };
return new ParseResult { Message = message, HeaderSize = headerSize, BodyLength = bodyLength };
}
/// <summary>整帧构建(帧头 + è´Ÿè½½ï¼‰ã€‚æž„å»ºä¸æ”¹åŠ¨æ¶ˆæ¯æ•°æ®</summary>
/// <param name="message">消æ¯ï¼ˆ<see cref="WsMessage"/> æä¾›ç±»åž‹ä¸ŽæŽ©ç é”®ï¼›å…¶ä»–æ¶ˆæ¯æŒ‰äºŒè¿›åˆ¶å¸§å¤„ç†ï¼‰</param>
/// <returns>æ•´å¸§æ‹¥æœ‰å¥æŸ„,调用方负责 Dispose</returns>
/// <remarks>
/// <para><b>æœåŠ¡ç«¯æ–¹å‘</b>:已预留的负载零拷è´å€Ÿä½å…±äº«ï¼ˆå¸§å¤´è½åœ¨é¢„留区);其余以新头节点挂接负载链。两ç§ç–ç•¥éƒ½ä¸æ”¹åŠ¨æ¶ˆæ¯è´Ÿè½½ã€‚</para>
/// <para><b>客户端方å‘</b>ï¼šæŽ©ç æ˜¯åŽŸåœ° XORï¼ˆç ´åæ€§ï¼‰ï¼Œå¿…须独å ——新分é…头部+负载,拷è´çš„åŒæ—¶æŽ©ç ,æºè´Ÿè½½ä¸å—å½±å“。</para>
/// <para><b>时效</b>:æœåŠ¡ç«¯ç»“æžœå¯èƒ½å¼•用消æ¯è´Ÿè½½ç¼“冲,帧å‘é€å®Œæˆå‰ä¸å¾—å¤ç”¨æˆ–改写。</para>
/// </remarks>
/// <exception cref="InvalidOperationException">消æ¯ä½“为æµå¼æ¨¡å¼ï¼ˆè¯·ä½¿ç”¨ BuildHeader + æµå¼å‘é€ï¼‰</exception>
public IOwnerPacket? Build(IMessage message)
{
if (message == null) throw new ArgumentNullException(nameof(message));
if (message.Body is { IsStreaming: true }) throw new InvalidOperationException("æµå¼æ¶ˆæ¯ä½“æ— æ³•æ•´å¸§æž„å»ºï¼Œè¯·ä½¿ç”¨ BuildHeader + æµå¼å‘é€");
var ws = message as WsMessage;
var body = message.Payload;
var len = body?.Total ?? 0;
// 掩ç :仅客户端方å‘;未指定时生æˆéšæœºå¯†é’¥ï¼ˆå¸§å†…ä¸´æ—¶ï¼Œä¸å†™å›žæ¶ˆæ¯ï¼‰
Byte[]? masks = null;
if (!IsServer)
{
masks = ws?.MaskKey;
if (masks == null || masks.Length < 4)
{
masks = new Byte[4];
#if NET6_0_OR_GREATER
Random.Shared.NextBytes(masks);
#else
new Random().NextBytes(masks);
#endif
}
}
// 头部大å°ï¼šFIN+OPCODE(1) + 长度(1/3/9) + 掩ç (0/4)
var size = len switch
{
< 126 => 1 + 1,
<= 0xFFFF => 1 + 1 + 2,
_ => 1 + 1 + 8,
};
// 线上掩ç 键固定 4 å—节:按数组实际长度算头长,会与 WriteHeader 写出的å—节数ä¸ä¸€è‡´ï¼Œå¯¼è‡´å¸§é”™ä½
if (masks != null) size += 4;
IOwnerPacket pk;
if (masks == null)
{
// æ— æŽ©ç ï¼šå·²é¢„ç•™çš„æ‹¥æœ‰å¥æŸ„é›¶æ‹·è´å€Ÿä½å…±äº«ï¼Œå…¶ä½™æ–°å¤´èŠ‚ç‚¹æŒ‚æŽ¥è´Ÿè½½é“¾
pk = body.PrepareHeader(size);
}
else
{
// æŽ©ç æ˜¯åŽŸåœ° XORï¼ˆç ´åæ€§ï¼‰ï¼Œå¿…须独å :新分é…头部+负载,拷è´çš„åŒæ—¶æŽ©ç ,æºè´Ÿè½½ä¸å—å½±å“
pk = new OwnerPacket(size + len);
var span = pk.GetSpan();
if (len > 0)
{
body!.ReadBytes(span[size..]);
ApplyMask(span[size..], masks, 0);
}
}
// 帧头由消æ¯ç±»å†™å…¥ï¼ˆFIN/OPCODE/长度/掩ç é”®ï¼›éž WsMessage æ¶ˆæ¯æŒ‰äºŒè¿›åˆ¶å¸§å¤„ç†ï¼‰
var header = ws ?? new WsMessage { Type = WebSocketMessageType.Binary };
header.WriteHeader(pk.GetSpan(), len, masks);
return pk;
}
/// <summary>仅构建头部数æ®åŒ…,声明负载长度(头 + æµå¼ä½“å‘é€ï¼‰</summary>
/// <param name="message">消æ¯ï¼ˆ<see cref="WsMessage"/> æä¾›ç±»åž‹ï¼›å…¶ä»–æ¶ˆæ¯æŒ‰äºŒè¿›åˆ¶å¸§å¤„ç†ï¼‰</param>
/// <param name="bodyLength">è´Ÿè½½å—节数(0 ~ 2147483647)</param>
/// <returns>头部数æ®åŒ…,调用方负责 Dispose</returns>
/// <exception cref="ArgumentOutOfRangeException">长度为负或超过 32 ä½å议上é™</exception>
/// <exception cref="NotSupportedException">客户端方å‘(带掩ç çš„å¸§ä¸æ”¯æŒæµå¼å‘é€ï¼Œè¯·ä½¿ç”¨æ•´å¸§æž„建)</exception>
public IOwnerPacket BuildHeader(IMessage message, Int64 bodyLength)
{
if (message == null) throw new ArgumentNullException(nameof(message));
if (bodyLength < 0) throw new ArgumentOutOfRangeException(nameof(bodyLength), "Payload length must be non-negative.");
if (bodyLength > Int32.MaxValue) throw new ArgumentOutOfRangeException(nameof(bodyLength), "Payload length exceeds the 32-bit protocol limit.");
if (!IsServer) throw new NotSupportedException("带掩ç çš„å¸§ä¸æ”¯æŒæµå¼å‘é€ï¼Œè¯·ä½¿ç”¨æ•´å¸§æž„建");
var size = bodyLength switch
{
< 126 => 1 + 1,
<= 0xFFFF => 1 + 1 + 2,
_ => 1 + 1 + 8,
};
var pk = new OwnerPacket(size);
// 帧头由消æ¯ç±»å†™å…¥ï¼ˆæœåŠ¡ç«¯æ–¹å‘æ— 掩ç ï¼›éž WsMessage æ¶ˆæ¯æŒ‰äºŒè¿›åˆ¶å¸§å¤„ç†ï¼‰
var header = message as WsMessage ?? new WsMessage { Type = WebSocketMessageType.Binary };
header.WriteHeader(pk.GetSpan(), bodyLength, null);
return pk;
}
#endregion
#region 辅助
/// <summary>æ ¡éªŒæ¶ˆæ¯è´Ÿè½½æ˜¯å¦ä¸ºåˆæ³• UTF-8(跨段连ç»ï¼‰</summary>
/// <param name="payload">负载;null è§†ä¸ºåˆæ³•</param>
/// <returns>是å¦åˆæ³•</returns>
/// <remarks>
/// <para>RFC 6455 §8.1:文本消æ¯çš„è´Ÿè½½å¿…é¡»æ˜¯åˆæ³• UTF-8ï¼Œç•¸å˜æ•°æ®æŒ‰ 1007(Invalid frame payload data)失败连接——
/// 直接交给业务会让对端拿到乱ç 而ä¸è‡ªçŸ¥ã€‚åˆæ³• UTF-8(RFC 3629ï¼‰éœ€åŒæ—¶æ»¡è¶³ï¼šç»å—èŠ‚å½¢æ€æ£ç¡®ã€
/// æ— è¿‡é•¿ç¼–ç (overlong)ã€ä¸è½åœ¨ä»£ç†åŒº U+D800~U+DFFFã€ä¸è¶…过 U+10FFFF。</para>
/// <para>分片消æ¯çš„负载是å„片å—节的连ç»ä½“,故按链å¼è´Ÿè½½è·¨æ®µè¿žç»æ ¡éªŒï¼ˆå¤šå—节åºåˆ—å¯è·¨ç‰‡/跨段)。</para>
/// </remarks>
public static Boolean IsValidUtf8(IPacket? payload)
{
if (payload == null) return true;
// 状æ€ï¼šneed 为待ç»å—节数(>0 表示上一åºåˆ—未收尾,跨段时出现),min 为本åºåˆ—å…许的最å°ç 点(overlong 判æ®ï¼‰
var need = 0;
var cp = 0u;
var min = 0u;
for (var pk = payload; pk != null; pk = pk.Next)
{
var span = pk.GetSpan();
for (var i = 0; i < span.Length; i++)
{
var b = span[i];
// ç»å—节
if (need > 0)
{
if ((b & 0xC0) != 0x80) return false;
cp = (cp << 6) | (UInt32)(b & 0x3F);
if (--need > 0) continue;
// åºåˆ—æ”¶å°¾ï¼šæ ¡éªŒè¿‡çŸç¼–ç (overlong)ã€ä»£ç†åŒºä¸Žç 点上é™
if (cp < min || (cp >= 0xD800 && cp <= 0xDFFF) || cp > 0x10FFFF) return false;
continue;
}
if (b < 0x80) continue;
if ((b & 0xE0) == 0xC0) { need = 1; cp = (UInt32)(b & 0x1F); min = 0x80; }
else if ((b & 0xF0) == 0xE0) { need = 2; cp = (UInt32)(b & 0x0F); min = 0x800; }
else if ((b & 0xF8) == 0xF0) { need = 3; cp = (UInt32)(b & 0x07); min = 0x10000; }
else return false;
}
}
// åºåˆ—未收尾(负载截æ–在多å—节å—符ä¸é—´ï¼‰
return need == 0;
}
/// <summary>判æ–å…³é—状æ€ç 是å¦å…许出现在线路上(RFC 6455 §7.4.1)</summary>
/// <param name="closeStatus">å…³é—状æ€ç </param>
/// <returns>å¯å‘é€è¿”回 true</returns>
/// <remarks>1000~1003ã€1007~1014 为å议与 IANA 已分é…值,3000~4999 供应用与注册使用;
/// 1004/1005/1006/1015 属ä¿ç•™å€¼ï¼Œ0~999 与 1016~2999 未分é…,端点å‡ç¦æ¢å‘é€â€”—
/// å‘é€è¿™ç±»ç 会被对端判定为å议错误,把优雅关é—å˜æˆå¼‚常关é—</remarks>
internal static Boolean IsSendableCloseStatus(Int32 closeStatus) => closeStatus switch
{
>= 1000 and <= 1003 => true,
>= 1007 and <= 1014 => true,
>= 3000 and <= 4999 => true,
_ => false,
};
/// <summary>æž„é€ å…³é—å¸§æ£æ–‡ï¼š2 å—节网络åºçжæ€ç åŠ UTF-8 åŽŸå› ï¼ˆRFC 6455 §5.5.1)</summary>
/// <param name="closeStatus">å…³é—状æ€ç </param>
/// <param name="statusDescription">åŽŸå› æè¿°ï¼Œå¯ä¸ºç©º</param>
/// <returns>å…³é—å¸§æ£æ–‡æ•°æ®åŒ…,调用方负责 Dispose</returns>
/// <remarks>ä¸å¯å‘é€çš„状æ€ç é€€åŒ–ä¸ºæ— è´Ÿè½½ï¼ˆå³â€œæ— 状æ€ç å…³é—帧â€è¯ä¹‰ï¼‰ï¼›åŽŸå› æŒ‰æŽ§åˆ¶å¸§ä¸Šé™æˆªæ–到 123 å—节且ä¸åˆ‡æ–多å—节å—符</remarks>
internal static IPacket BuildClosePayload(Int32 closeStatus, String? statusDescription)
{
if (!IsSendableCloseStatus(closeStatus))
{
NewLife.Log.XTrace.WriteLine("WebSocket å…³é—状æ€ç {0} ä¸å…许出现在线路上,改为å‘逿— 状æ€ç å…³é—帧", closeStatus);
return new ArrayPacket(new Byte[0]);
}
var desc = (statusDescription ?? String.Empty).GetBytes();
// RFC 6455 §5.5:控制帧负载ä¸å¾—超过 125 å—节——状æ€ç 2 å—节 + åŽŸå› æœ€å¤š 123 å—节。
// 䏿ˆªæ–会å‘å‡ºéžæ³•æŽ§åˆ¶å¸§ï¼Œå¯¹ç«¯ï¼ˆå«æœ¬åº“æœåŠ¡ç«¯ï¼‰ä¼šæŒ‰å议错误 1002 å…³é—,优雅关é—å˜æˆå¼‚常关é—
if (desc.Length > 123) desc = TruncateUtf8(desc, 123);
var buf = new Byte[2 + desc.Length];
buf[0] = (Byte)(closeStatus >> 8);
buf[1] = (Byte)closeStatus;
desc.CopyTo(buf, 2);
return new ArrayPacket(buf);
}
/// <summary>按 UTF-8 å—符边界截æ–å—节,ä¸åˆ‡æ–多å—节å—符</summary>
/// <param name="bytes">原始å—节</param>
/// <param name="max">最大å—节数</param>
/// <returns>截æ–结果</returns>
private static Byte[] TruncateUtf8(Byte[] bytes, Int32 max)
{
var len = max;
// ç»å—节(10xxxxxx)ä¸èƒ½ä½œä¸ºèµ·ç‚¹ï¼šå›žé€€åˆ°è¯¥å—符首å—节,å¦åˆ™æˆªå‡ºåŠä¸ªå—ç¬¦ï¼ˆå¯¹ç«¯è¯»åˆ°éžæ³• UTF-8)
while (len > 0 && (bytes[len] & 0xC0) == 0x80) len--;
var buf = new Byte[len];
Array.Copy(bytes, buf, len);
return buf;
}
/// <summary>按掩ç 键对数æ®åŽŸåœ° XOR è§£ç /ç¼–ç (跨段连ç»ï¼šå移按负载起点累计)</summary>
/// <param name="data">ç›®æ ‡æ•°æ®</param>
/// <param name="masks">掩ç 键(至少 4 å—节)</param>
/// <param name="offset">当剿•°æ®åœ¨è´Ÿè½½ä¸çš„åç§»</param>
private static void ApplyMask(Span<Byte> data, Byte[] masks, Int64 offset)
{
for (var i = 0; i < data.Length; i++)
{
data[i] = (Byte)(data[i] ^ masks[(Int32)((offset + i) & 3)]);
}
}
#endregion
}
|