using NewLife.Data;
namespace NewLife.Messaging;
/// <summary>消æ¯å¸§æ³µã€‚在数æ®åŒ…管é“上按å议(<see cref="IMessageCodec"/>)定界消æ¯å¸§ï¼šå¤´éƒ¨åˆ°é½å³ç»‘定体并交付</summary>
/// <remarks>
/// <para><b>头部到é½å³äº¤ä»˜</b>:帧头一旦完整å³å¯äº§å‡ºæ¶ˆæ¯ï¼Œä¸å¿…ç‰æ•´å¸§åˆ°é½ã€‚体按到达情况绑定:</para>
/// <list type="bullet">
/// <item><description><b>帧已完整</b>:<see cref="PipeReader.TakeFrame(Int64)"/> é›¶æ‹·è´åˆ‡å‡ºæ•´å¸§ï¼Œä½“为内å˜è§†å›¾ï¼ˆ<see cref="LimitedReader.IsStreaming"/> 为 false)</description></item>
/// <item><description><b>帧未完整</b>:体为æµå¼è¯»å–器(<see cref="PipeReader.Limit(Int64)"/>),数æ®éšç®¡é“到达,大帧ä¸å¿…整载入内å˜</description></item>
/// </list>
/// <para>两ç§ç»‘å®šå¯¹æ¶ˆè´¹æ–¹é€æ˜Žï¼ˆç»Ÿä¸€ç» <see cref="Message.Body"/> 读å–ï¼‰ã€‚æœªè¯»ä½“åœ¨äº¤ä»˜æ”¶å°¾æ—¶ç» <see cref="DiscardAsync"/> ä¸¢å¼ƒï¼Œä¿æŒä¸‹ä¸€å¸§å®šç•Œå¯¹é½ã€‚</para>
/// <para><b>æ— çŠ¶æ€</b>:åè®®å®žä¾‹æ— çŠ¶æ€å¯è·¨è¿žæŽ¥å…±äº«ï¼›å¸§æ³µå®žä¾‹å¯å¤ç”¨äºŽå¤šè¿žæŽ¥ï¼ˆè·Ÿéš <see cref="PipeReader"/> å•读约æŸï¼‰ã€‚</para>
/// </remarks>
/// <example>
/// <code>
/// var pump = new MessagePump(new SrmpCodec());
/// while (true)
/// {
/// var msg = await pump.ReadAsync(pipe.Reader, cancellationToken);
/// if (msg == null) break; // æµç»“æŸ
/// try
/// {
/// // å¤´éƒ¨å—æ®µç«‹å³å¯ç”¨ï¼›è´Ÿè½½æŒ‰éœ€è¯»å–:await msg.Body.ReadAllAsync()
/// }
/// finally
/// {
/// await MessagePump.DiscardAsync(msg); // 丢弃未读体对é½ä¸‹ä¸€å¸§
/// msg.TryDispose();
/// }
/// }
/// </code>
/// </example>
public class MessagePump
{
#region 属性
/// <summary>消æ¯ç¼–è§£ç 器(帧å议)</summary>
public IMessageCodec? Codec { get; set; }
/// <summary>最大缓å˜å—èŠ‚æ•°ï¼ˆæ— æ³•å®šç•Œçš„æ®‹ä½™ä¸Šé™ï¼‰ï¼Œé»˜è®¤ 1M。0 表示ä¸é™åˆ¶</summary>
/// <remarks>
/// <para>æ£å¸¸å¸§å¤´éƒ¨åˆ°é½å³äº¤ä»˜ï¼Œä¸äº§ç”Ÿç´¯ç§¯ï¼›æ®‹ä½™æŒç»å¢žé•¿è¯´æ˜Žå¯¹ç«¯æ•°æ®ä¸Žåè®®ä¸åŒ¹é…或已æŸå。</para>
/// <para><see cref="TryRead"/> 䏿‰§è¡Œæœ¬æ£€æŸ¥ï¼ˆçº¯ä¸ç‰å¾…è¯ä¹‰ï¼‰ï¼Œç”± <see cref="ReadAsync"/> 在ç‰å¾…è¿‡ç¨‹ä¸æ‰§è¡Œï¼šè¾¾åˆ°ä¸Šé™æ—¶æŠ›å‡ºå¼‚常,é¿å…连接僵æ»ã€‚</para>
/// </remarks>
public Int32 MaxCache { get; set; } = 1024 * 1024;
/// <summary>è¦æ±‚整帧完整æ‰äº§å‡ºï¼ˆæ•´å¸§æ¨¡å¼ï¼‰ã€‚默认 false:头部到é½å³äº¤ä»˜ï¼ˆä½“å¯ä¸ºæµå¼ï¼‰</summary>
/// <remarks>
/// åŒæ¥æ³µåœºæ™¯ï¼ˆå¦‚ WebSocket æœåŠ¡ç«¯çš„åŒæ¥å¸§å¾ªçŽ¯ï¼‰æ— æ³•å¼‚æ¥æ¶ˆè´¹æµå¼ä½“:帧未完整时ä¸äº§å‡ºã€ä¸æ¶ˆè´¹ï¼Œç•™å¾…æ•°æ®åˆ°é½åŽé‡æ–°è§£æžã€‚
/// 该模å¼ä¸‹çš„å•帧长度上é™ç”± <see cref="MaxFrameSize"/> 约æŸã€‚
/// </remarks>
public Boolean RequireFullFrame { get; set; }
/// <summary>å•帧长度上é™ï¼ˆæ•´å¸§æ¨¡å¼ä¸Žå议预绑定体),默认 16M。0 表示ä¸é™åˆ¶</summary>
/// <remarks>
/// <para>整帧模å¼ä¸Žåè®®é¢„ç»‘å®šä½“éƒ½ä¸æ¶ˆè´¹æœªå®Œæ•´å¸§ï¼Œä»…å‡å¯¹ç«¯å£°æ˜Žä¸€ä¸ªè¶…大帧长度å³å¯è®©ç®¡é“æ— é™å 用内å˜ï¼š<see cref="MaxCache"/> åªç®¡â€œæ— 法定界的残余â€ï¼Œ
/// 管ä¸åˆ°â€œå·²å®šç•Œä½†æ°¸è¿œåˆ°ä¸é½â€çš„å¸§ã€‚è¶…è¿‡ä¸Šé™æ—¶æŠ›å‡ºå¼‚常,由调用方按å议错误关é—连接。</para>
/// </remarks>
public Int32 MaxFrameSize { get; set; } = 16 * 1024 * 1024;
#endregion
#region æž„é€
/// <summary>实例化帧泵</summary>
public MessagePump() { }
/// <summary>实例化帧泵</summary>
/// <param name="codec">消æ¯ç¼–è§£ç 器</param>
public MessagePump(IMessageCodec codec) => Codec = codec;
#endregion
#region 读å–
/// <summary>å°è¯•读å–一帧(ä¸ç‰å¾…)。头部到é½å³è¿”回消æ¯ï¼›æ— 消æ¯å¸§ï¼ˆç©ºè¡Œ/心跳ç‰ï¼‰è‡ªåŠ¨è·³è¿‡</summary>
/// <param name="reader">æ•°æ®åŒ…读å–器</param>
/// <param name="message">è§£æžå‡ºçš„æ¶ˆæ¯ï¼ˆæˆåŠŸæ—¶æœ‰æ•ˆï¼›å¤´éƒ¨å—æ®µå°±ä½ã€ä½“已绑定)</param>
/// <returns>是å¦è¯»å–到消æ¯ï¼›å¤´éƒ¨ä¸è¶³æˆ–窗å£å†…ä»…æœ‰æ— æ¶ˆæ¯å¸§æ—¶è¿”回 false(窗å£ä¸åŠ¨ï¼Œç‰è¿½åŠ ï¼‰</returns>
/// <exception cref="InvalidOperationException"><see cref="Codec"/> 未设置;或å议帧æŸåï¼ˆå¤´éƒ¨å·²å®Œæ•´ä½†é•¿åº¦å—æ®µéžæ³•)</exception>
/// <exception cref="FrameTooLargeException">整帧模å¼ä¸‹å•帧超过 <see cref="MaxFrameSize"/></exception>
/// <remarks>åè®®è¿”å›žæ— æ¶ˆæ¯å¸§ï¼ˆ<see cref="IMessageCodec.TryParse"/> æˆåŠŸä½†æ¶ˆæ¯ä¸º null)时消费该帧并继ç»è§£æžä¸‹ä¸€å¸§ï¼Œç›´åˆ°äº§å‡ºæ¶ˆæ¯æˆ–æ•°æ®ä¸è¶³ã€‚</remarks>
public Boolean TryRead(PipeReader reader, out IMessage? message) => TryReadCore(reader, out message, out _);
/// <summary>å°è¯•读å–ä¸€å¸§ï¼ˆæ ¸å¿ƒï¼‰ã€‚frameLength 返回“已定界但整帧未到é½â€çš„帧长,0 è¡¨ç¤ºæ— è¿›å±•ï¼ˆå¤´éƒ¨ä¸è¶³ï¼‰</summary>
private Boolean TryReadCore(PipeReader reader, out IMessage? message, out Int64 frameLength)
{
if (reader == null) throw new ArgumentNullException(nameof(reader));
message = null;
frameLength = 0;
var codec = Codec ?? throw new InvalidOperationException("MessagePump.Codec not set.");
while (true)
{
var buffer = reader.Buffer;
if (buffer.IsEmpty) return false;
// 定界:头部ä¸è¶³è¿”回 nullï¼Œä¸æ¶ˆè´¹ã€ä¸äº§ç”Ÿå¯¹è±¡
var rs = codec.TryParse(buffer);
if (rs == null) return false;
// æŸåå¸§ï¼šç«‹å³æŠ¥é”™ï¼ˆç”±ä¼šè¯å±‚å…³é—连接),ä¸è¿›å…¥ç‰å¾…——å¦åˆ™ä¼šåƒµæ»åˆ°æ®‹ä½™ä¸Šé™æ‰æ–å¼€
if (rs.Value.Invalid) throw new InvalidOperationException($"å议帧æŸåï¼šå¤´éƒ¨å·²å®Œæ•´ä½†å¸§é•¿åº¦å—æ®µéžæ³•(åè®® {codec.GetType().Name})");
var headerSize = rs.Value.HeaderSize;
var bodyLength = rs.Value.BodyLength;
// 已定界但需整帧(装饰å议声明,或本泵整帧模å¼ï¼‰ï¼šä¸äº§å‡ºã€ä¸æ¶ˆè´¹ï¼Œç‰æ•°æ®åˆ°é½åŽé‡æ–°è§£æžã€‚
// 必须与“头部ä¸è¶³â€åŒºåˆ†â€”—å¦åˆ™ä¸Šå±‚会把æ£åœ¨æ£å¸¸ç´¯ç§¯çš„å¤§å¸§å½“æˆæ— 法定界的残余而误æ–连接
if ((rs.Value.NeedFullFrame || (RequireFullFrame && rs.Value.Message != null)) && headerSize + bodyLength > buffer.Length)
{
rs.Value.Message?.Dispose();
frameLength = headerSize + bodyLength;
// å•帧长度超上é™ç«‹å³æŠ¥é”™ï¼šå¦åˆ™å¯¹ç«¯åªéœ€å£°æ˜Žä¸€ä¸ªè¶…å¤§å¸§ï¼Œå°±èƒ½è®©æœ¬è¿žæŽ¥çš„å†…å˜æ— é™å¢žé•¿
if (MaxFrameSize > 0 && frameLength > MaxFrameSize)
throw new FrameTooLargeException(frameLength, MaxFrameSize);
return false;
}
// æ— æ¶ˆæ¯å¸§ï¼ˆç©ºè¡Œ/心跳ç‰ï¼‰ï¼šæ¶ˆè´¹è¯¥å¸§åŽç»§ç»è§£æžä¸‹ä¸€å¸§
if (rs.Value.Message == null)
{
var skip = headerSize;
// å议错误防护:既ä¸äº§å‡ºæ¶ˆæ¯åˆä¸æ¶ˆè´¹å—节会æ»å¾ªçŽ¯ï¼Œä¹Ÿä¸å¾—消费超出窗å£
if (skip <= 0 || skip > buffer.Length) return false;
reader.AdvanceTo(skip, skip);
continue;
}
var msg = rs.Value.Message;
// å议已预绑定体(如压缩å议解压åŽé‡ç»‘):直接消费整帧,ä¸å†äºŒæ¬¡ç»‘定。
// 预绑定体的å议自己ä¿è¯æ•´å¸§åˆ°é½ï¼Œä½†è¿™é‡Œä»è¦æ ¡éªŒçª—å£ï¼šè‡ªå®šä¹‰å议若æå‰ç»‘定体,
// å¸§æœªåˆ°é½æ—¶æŽ¨è¿›çª—å£ä¼šè¶Šç•ŒæŠ›å¼‚å¸¸ï¼Œæ‰“æ–æ•´æ¡æŽ¥æ”¶é“¾
if (msg.Payload != null)
{
if (headerSize + bodyLength > buffer.Length)
{
msg.Dispose();
// 已定界但整帧未到é½ï¼Œä¸Ž NeedFullFrame 分支åŒå£å¾„:
// ä¸æŠ¥å‘Šå¸§é•¿ä¼šè¢« ReadAsync å½“ä½œâ€œæ— æ³•å®šç•Œçš„æ®‹ä½™â€ï¼Œå¸§é•¿è¶…过 MaxCache 时把æ£å¸¸ç´¯ç§¯çš„å¤§å¸§è¯¯åˆ¤ä¸ºåæ•°æ®è€Œæ–连
frameLength = headerSize + bodyLength;
if (MaxFrameSize > 0 && frameLength > MaxFrameSize)
throw new FrameTooLargeException(frameLength, MaxFrameSize);
return false;
}
reader.AdvanceTo(headerSize + bodyLength, headerSize + bodyLength);
message = msg;
return true;
}
// 消费头部;体起点å³å½“å‰è¯»å–ä½ç½®
reader.AdvanceTo(headerSize, headerSize);
// 帧已完整:零拷è´åˆ‡å‡ºæ•´å¸§ï¼Œä½“为内å˜è§†å›¾ï¼›å¦åˆ™ç»‘定æµå¼ä½“,数æ®éšç®¡é“到达
if (headerSize + bodyLength <= buffer.Length)
{
// 整帧模å¼ä¸‹çš„å•帧上é™å¯¹â€œå·²å®Œæ•´å¸§â€åŒæ ·ç”Ÿæ•ˆï¼šå¦åˆ™å¯¹ç«¯æŠŠè¶…é™å¸§åŽ‹åœ¨ä¸€è½®é‡Œå³å¯ç»•过上é™
if (RequireFullFrame && MaxFrameSize > 0 && headerSize + bodyLength > MaxFrameSize)
{
msg.Dispose();
throw new FrameTooLargeException(headerSize + bodyLength, MaxFrameSize);
}
var frame = reader.TakeFrame(bodyLength);
msg.SetBody(frame);
}
else
{
msg.BindBody(reader.Limit(bodyLength));
}
message = msg;
return true;
}
}
/// <summary>读å–一帧;数æ®ä¸è¶³æ—¶ç‰å¾…è¿½åŠ </summary>
/// <param name="reader">æ•°æ®åŒ…读å–器</param>
/// <param name="cancellationToken">å–æ¶ˆé€šçŸ¥</param>
/// <returns>消æ¯ï¼ˆå¤´éƒ¨å—段就ä½ã€ä½“已绑定);æµç»“æŸæˆ–æœ‰æ®‹ä½™æ— æ³•æˆå¸§æ—¶è¿”回 null</returns>
/// <remarks>æ•°æ®ä¸è¶³æ—¶æŒ‰ examined è¯ä¹‰æ ‡è®°å·²æ£€æŸ¥å¹¶ç»§ç»ç‰å¾…ï¼›æµç»“æŸæ—¶çª—å£å†…æ— æ³•æˆå¸§çš„æ®‹ä½™æ•°æ®è¢«ä¸¢å¼ƒã€‚</remarks>
public async ValueTask<IMessage?> ReadAsync(PipeReader reader, CancellationToken cancellationToken = default)
{
if (reader == null) throw new ArgumentNullException(nameof(reader));
while (true)
{
var rr = await reader.ReadAsync(cancellationToken).ConfigureAwait(false);
if (rr.IsCanceled) return null;
// 头部到é½ï¼šç›´æŽ¥äº§å‡ºæ¶ˆæ¯ï¼ˆå¸§å®Œæ•´/æœªå®Œæ•´åœ¨æ ¸å¿ƒå†…éƒ¨åˆ†æµï¼‰
if (TryReadCore(reader, out var message, out var frameLength)) return message;
// æ— æ³•å®šç•Œï¼šæµå·²ç»“æŸåˆ™ä¸¢å¼ƒæ®‹ä½™ï¼Œå¦åˆ™æ ‡è®°å·²æ£€æŸ¥åˆ°çª—壿œ«å°¾ç‰å¾…è¿½åŠ ã€‚
// å¸¦å¼‚å¸¸ç»“æŸæ—¶ä¸ä¼ªè£…æˆä¼˜é›…å…³é—(BCL çš„ IsCompletedOrThrow åŒä¹‰ï¼‰ï¼šæŠŠæ•…障抛给调用方按错误处ç†
if (rr.IsCompleted)
{
if (reader.Error is { } error) throw error;
return null;
}
var buffer = reader.Buffer;
// åè®®é”™è¯¯é˜²æŠ¤ï¼šæ®‹ä½™æ— æ³•å®šç•Œä¸”æŒç»å¢žé•¿ï¼Œè¾¾åˆ°ä¸Šé™å³å¿«é€Ÿå¤±è´¥ï¼Œé¿å…连接僵æ»ã€‚
// 已定界但整帧未到é½ï¼ˆframeLength > 0)ä¸å±žäºŽâ€œæ— 法定界的残余â€ï¼Œå…¶ä¸Šé™ç”± MaxFrameSize 把关
if (frameLength <= 0 && MaxCache > 0 && buffer.Length >= MaxCache)
throw new InvalidOperationException($"æ— æ³•å®šç•Œçš„æ®‹ä½™æ•°æ® {buffer.Length} å—èŠ‚è¾¾åˆ°ä¸Šé™ {MaxCache},对端数æ®ä¸Žåè®®ä¸åŒ¹é…或已æŸå");
reader.AdvanceTo(0, buffer.Length);
}
}
#endregion
#region 交付收尾
/// <summary>ä¸¢å¼ƒæ¶ˆæ¯æœªè¯»ä½“,使主读å–器对é½å¸§å°¾ï¼ˆäº¤ä»˜æ”¶å°¾ï¼›éšåŽåº”释放消æ¯ï¼‰</summary>
/// <param name="message">消æ¯ï¼ˆå¯ä¸º null,安全)</param>
/// <param name="cancellationToken">å–æ¶ˆé€šçŸ¥</param>
/// <returns>æœªè¯»ä½™é‡æ¸…é›¶ï¼›æµç»“æŸæ—¶å¯èƒ½ä¿ç•™ä½™é‡ï¼ˆæ— 法对é½ï¼‰</returns>
/// <remarks>å¤„ç†æ–¹ä¸è¯»æ¶ˆæ¯ä½“时(仅头部è¯ä¹‰å³å¯å®Œæˆå¤„ç†ï¼‰ï¼Œç”±äº¤ä»˜è·¯å¾„调用本方法跳过余é‡ï¼Œä¿æŒä¸‹ä¸€å¸§å®šç•Œå¯¹é½ï¼›å·²è¯»æ»¡æˆ–已释放的消æ¯è°ƒç”¨æ— æ“作。</remarks>
public static async ValueTask DiscardAsync(IMessage? message, CancellationToken cancellationToken = default)
{
if (message?.Body is { Remaining: > 0 } body) await body.DrainAsync(cancellationToken).ConfigureAwait(false);
}
#endregion
}
|