using System.Buffers;
namespace NewLife.Data;
/// <summary>é™é•¿è¯»å–器。在管é“主读å–器上é™å®šå—节预算(æµå¼æ¨¡å¼ï¼‰ï¼Œæˆ–å¯¹å†…å˜æ•°æ®åŒ…é™å®šçª—å£ï¼ˆå†…å˜æ¨¡å¼ï¼‰ï¼Œç”¨äºŽæŒ‰å¸§é•¿è¯»å–负载(body)</summary>
/// <remarks>
/// <para><b>ä¸¤ç§æ¨¡å¼</b>:æµå¼æ¨¡å¼ç”± <see cref="PipeReader.Limit(Int64)"/> 创建,数æ®éšç®¡é“åˆ°è¾¾ï¼›å†…å˜æ¨¡å¼ç”±æ¶ˆæ¯æ•´å¸§è§£æžåˆ›å»ºï¼Œæ•°æ®å·²åœ¨å†…å˜ï¼ˆè¯»å–ç«‹å³å®Œæˆï¼‰ã€‚</para>
/// <para><b>视图与消费</b>:<see cref="AsPacket"/> å–å‰©ä½™ä½“è§†å›¾ï¼ˆä¸æ¶ˆè´¹ï¼‰ï¼›<see cref="ReadAllAsync"/> 读满并推进预算(消费);需è¦é•¿é©»å†…容时请使用 ReadAllAsync。</para>
/// <para><b>çŸè¯»ä¸é™é»˜</b>:æµå¼æ¨¡å¼ <see cref="ReadAllAsync"/> 未读满就é‡åˆ°å–消或æµç»“æŸï¼ˆå¯¹ç«¯å…³é—ã€ç®¡é“故障)时抛异常,ä¸è¿”å›žåŠæˆªæ¶ˆæ¯ä½“â€”â€”åŠæˆªä½“除了长度对ä¸ä¸Šä¹‹å¤–没有任何å¯åˆ¤å®šçš„ä¿¡å·ï¼Œé™é»˜è¿”回ç‰äºŽæŠŠæˆªæ–ä¼ªè£…æˆæ£å¸¸æ¶ˆæ¯ã€‚</para>
/// <para>读å–窗å£è£å‰ªåˆ°é¢„算内;<see cref="AdvanceTo(Int64, Int64)"/> é€ä¼ 主读å–器并扣å‡é¢„算。</para>
/// <para><see cref="DrainAsync"/> 丢弃未读余é‡ï¼ˆç‰å¾…æ•°æ®åˆ°è¾¾é€æ¥è·³è¿‡ï¼‰ï¼Œä½¿ä¸»è¯»å–器对é½åˆ°å¸§å°¾ï¼›æœªè¯»æ»¡å°±è¿›å…¥ä¸‹ä¸€å¸§å‰å¿…é¡» Drain,å¦åˆ™çª—å£é”™ä½ã€‚</para>
/// <para>预算耗尽åŽçš„读å–返回 <c>IsCompleted=true</c> 的空结果(体结æŸè¯ä¹‰ï¼‰ã€‚</para>
/// <para>啿¶ˆè´¹è€…ä½¿ç”¨ï¼›å†…å˜æ¨¡å¼çš„底层数æ®åŒ…所有æƒå½’æ¶ˆæ¯æŒæœ‰ï¼Œæœ¬è¯»å–å™¨åªæä¾›çª—å£è§†å›¾ã€‚</para>
/// </remarks>
public sealed class LimitedReader
{
#region 属性
private readonly PipeReader? _reader;
private readonly IPacket? _packet;
private readonly ReadOnlySequence<Byte> _sequence;
/// <summary>å†…å˜æ¨¡å¼èµ·ç‚¹å移(用于å¤ä½ï¼‰</summary>
private readonly Int64 _start;
/// <summary>å†…å˜æ¨¡å¼åˆå§‹é¢„算(用于å¤ä½ï¼‰</summary>
private readonly Int64 _initial;
private Int64 _offset;
private Int64 _remaining;
/// <summary>æ˜¯å¦æµå¼æ¨¡å¼ã€‚false ä¸ºå†…å˜æ¨¡å¼ï¼ˆæ•´å¸§è§£æžï¼Œè¯»å–ç«‹å³å®Œæˆï¼‰</summary>
public Boolean IsStreaming => _reader != null;
/// <summary>剩余预算(å—节)</summary>
public Int64 Remaining => _remaining;
/// <summary>当å‰çª—å£ï¼ˆè£å‰ªåˆ°å‰©ä½™é¢„ç®—å†…ï¼‰ã€‚æ— æ•°æ®æ—¶ä¸ºç©ºåºåˆ—</summary>
public ReadOnlySequence<Byte> Buffer
{
get
{
// å†…å˜æ¨¡å¼ï¼šç›´æŽ¥å¯¹é¢„建åºåˆ—切片
if (_reader == null) return _sequence.Slice(_offset, _remaining);
var buffer = _reader.Buffer;
return _remaining >= buffer.Length ? buffer : buffer.Slice(0, _remaining);
}
}
#endregion
#region æž„é€
/// <summary>在主读å–器上é™å®šå—节预算(æµå¼æ¨¡å¼ï¼‰</summary>
/// <param name="reader">主读å–器</param>
/// <param name="length">预算å—节数</param>
internal LimitedReader(PipeReader reader, Int64 length)
{
if (reader == null) throw new ArgumentNullException(nameof(reader));
if (length < 0) throw new ArgumentOutOfRangeException(nameof(length));
_reader = reader;
_remaining = length;
}
/// <summary>åœ¨å†…å˜æ•°æ®åŒ…上é™å®šçª—å£ï¼ˆå†…å˜æ¨¡å¼ï¼‰ã€‚æ•°æ®åŒ…所有æƒå½’调用方(消æ¯ï¼‰ï¼Œæœ¬è¯»å–å™¨åªæä¾›è§†å›¾</summary>
/// <param name="packet">底层数æ®åŒ…</param>
/// <param name="offset">èµ·å§‹åç§»</param>
/// <param name="length">窗å£å—节数</param>
internal LimitedReader(IPacket packet, Int32 offset, Int64 length)
{
if (packet == null) throw new ArgumentNullException(nameof(packet));
if (offset < 0) throw new ArgumentOutOfRangeException(nameof(offset));
if (length < 0) throw new ArgumentOutOfRangeException(nameof(length));
_packet = packet;
_sequence = packet.AsReadOnlySequence();
_start = offset;
_initial = length;
_offset = offset;
_remaining = length;
}
#endregion
#region 方法
/// <summary>å¤ä½åˆ°èµ·ç‚¹ï¼Œä½¿å†…容å¯è¢«é‡æ–°è¯»å–ã€‚ä»…å†…å˜æ¨¡å¼å¯ç”¨</summary>
/// <remarks>用于事件链(å¯è§‚测)已消费消æ¯ä½“åŽã€æŠŠæ¶ˆæ¯äº¤ä»˜ç»™ç‰å¾…æ–¹å‰æ¢å¤å…¶å¯è¯»æ€§ã€‚
/// æµå¼ä½“的数æ®ç”±ç®¡é“承载ã€ä¸å¯é‡æ”¾ï¼Œè°ƒç”¨å°†æŠ›å¼‚常。</remarks>
/// <exception cref="InvalidOperationException">æµå¼æ¨¡å¼ä¸å¯å¤ä½</exception>
internal void Reset()
{
if (_reader != null) throw new InvalidOperationException("æµå¼æ¨¡å¼ä¸æ”¯æŒå¤ä½");
_offset = _start;
_remaining = _initial;
}
/// <summary>è¯»å–æ•°æ®ã€‚预算耗尽时返回已结æŸçš„空结果;其余è¯ä¹‰ä¸Žä¸»è¯»å–器一致</summary>
/// <param name="cancellationToken">å–æ¶ˆé€šçŸ¥</param>
/// <returns>读å–结果(窗å£å·²è£å‰ªåˆ°é¢„算内)</returns>
public async ValueTask<ReadResult> ReadAsync(CancellationToken cancellationToken = default)
{
if (_remaining <= 0) return new ReadResult(ReadOnlySequence<Byte>.Empty, true, false);
// å†…å˜æ¨¡å¼ï¼šæ•°æ®å·²åœ¨å†…å˜ï¼Œè¯»å–ç«‹å³å®Œæˆ
if (_reader == null) return new ReadResult(Buffer, true, false);
// 主读å–器已éšç®¡é“结æŸå®Œæˆï¼šæŒ‰â€œä½“结æŸâ€è¯ä¹‰è¿”回,ä¸å†è¯»å–(结æŸåŽå†è¯»ä¼šæŠ›å¼‚常)
if (_reader.IsReaderCompleted) return new ReadResult(ReadOnlySequence<Byte>.Empty, true, false);
var rr = await _reader.ReadAsync(cancellationToken).ConfigureAwait(false);
if (rr.IsCanceled) return rr;
var buffer = rr.Buffer;
if (buffer.Length > _remaining) buffer = buffer.Slice(0, _remaining);
return new ReadResult(buffer, rr.IsCompleted, false);
}
/// <summary>å°è¯•åŒæ¥è¯»å–(ä¸ç‰å¾…)。有数æ®ã€å·²å–消ã€é¢„算耗尽或æµç»“æŸæ—¶è¿”回 true</summary>
/// <param name="result">读å–结果(窗å£å·²è£å‰ªåˆ°é¢„算内)</param>
/// <returns>是å¦ç«‹å³å¯è¯»ï¼›æ— æ•°æ®ä¸”æœªç»“æŸæ—¶è¿”回 false</returns>
public Boolean TryRead(out ReadResult result)
{
// 预算耗尽:体结æŸè¯ä¹‰
if (_remaining <= 0)
{
result = new ReadResult(ReadOnlySequence<Byte>.Empty, true, false);
return true;
}
// å†…å˜æ¨¡å¼ï¼šæ•°æ®å·²åœ¨å†…å˜ï¼Œç«‹å³è¿”回
if (_reader == null)
{
result = new ReadResult(Buffer, true, false);
return true;
}
// 主读å–器已éšç®¡é“结æŸå®Œæˆï¼šæŒ‰â€œä½“结æŸâ€è¯ä¹‰è¿”回,ä¸å†è¯»å–(结æŸåŽå†è¯»ä¼šæŠ›å¼‚常)
if (_reader.IsReaderCompleted)
{
result = new ReadResult(ReadOnlySequence<Byte>.Empty, true, false);
return true;
}
if (!_reader.TryRead(out var rr))
{
result = default;
return false;
}
if (rr.IsCanceled)
{
result = rr;
return true;
}
var buffer = rr.Buffer;
if (buffer.Length > _remaining) buffer = buffer.Slice(0, _remaining);
result = new ReadResult(buffer, rr.IsCompleted, false);
return true;
}
/// <summary>消费推进,ä¸å¾—超过剩余预算</summary>
/// <param name="consumedBytes">已消费å—节数</param>
public void AdvanceTo(Int64 consumedBytes) => AdvanceTo(consumedBytes, consumedBytes);
/// <summary>消费推进(å«å·²æ£€æŸ¥é•¿åº¦ï¼‰ï¼Œä¸å¾—超过剩余预算</summary>
/// <param name="consumedBytes">已消费å—节数</param>
/// <param name="examinedBytes">已检查å—节数</param>
public void AdvanceTo(Int64 consumedBytes, Int64 examinedBytes)
{
if (consumedBytes < 0 || consumedBytes > _remaining) throw new ArgumentOutOfRangeException(nameof(consumedBytes));
if (examinedBytes < consumedBytes) examinedBytes = consumedBytes;
if (examinedBytes > _remaining) examinedBytes = _remaining;
if (_reader == null)
_offset += consumedBytes;
else
_reader.AdvanceTo(consumedBytes, examinedBytes);
_remaining -= consumedBytes;
}
/// <summary>丢弃未读余é‡ï¼Œä½¿ä¸»è¯»å–器对é½åˆ°å¸§å°¾</summary>
/// <param name="cancellationToken">å–æ¶ˆé€šçŸ¥</param>
/// <returns>æœªè¯»ä½™é‡æ¸…é›¶ï¼›æµç»“æŸæ—¶å¯èƒ½ä¿ç•™ä½™é‡ï¼ˆæ— 法对é½ï¼‰</returns>
/// <remarks>ç‰å¾…æ•°æ®åˆ°è¾¾å¹¶é€æ¥è·³è¿‡ï¼›å–消时ä¸ä½œå¤„ç†ï¼Œç”±è°ƒç”¨æ–¹å†³å®š</remarks>
public async ValueTask DrainAsync(CancellationToken cancellationToken = default)
{
// å†…å˜æ¨¡å¼ï¼šç›´æŽ¥ä¸¢å¼ƒä½™é‡
if (_reader == null)
{
_offset += _remaining;
_remaining = 0;
return;
}
while (_remaining > 0)
{
var rr = await _reader.ReadAsync(cancellationToken).ConfigureAwait(false);
if (rr.IsCanceled) return;
var take = (Int64)Math.Min(rr.Buffer.Length, _remaining);
if (take > 0)
{
_reader.AdvanceTo(take, take);
_remaining -= take;
}
else if (rr.IsCompleted)
{
// æ•°æ®æœªåˆ°é½ä½†æµå·²ç»“æŸï¼šæ— 法对é½ï¼Œä¿ç•™ä½™é‡ç”±è°ƒç”¨æ–¹æ„ŸçŸ¥
return;
}
}
}
/// <summary>å–剩余体的数æ®åŒ…è§†å›¾ï¼ˆä¸æ¶ˆè´¹ï¼‰ã€‚ä»…å†…å˜æ¨¡å¼å¯ç”¨ï¼›æµå¼æ¨¡å¼è¯·ç”¨ <see cref="ReadAllAsync"/></summary>
/// <returns>剩余体视图;预算耗尽时返回 nullã€‚æ‹¥æœ‰å¥æŸ„为共享切片(å„自 Disposeï¼‰ï¼Œå€Ÿé˜…è§†å›¾ä»…åœ¨åº•å±‚æ•°æ®æœ‰æ•ˆæœŸå†…å¯ç”¨</returns>
/// <exception cref="InvalidOperationException">æµå¼æ¨¡å¼ä¸å…许直接å–包</exception>
public IPacket? AsPacket()
{
if (_reader != null) throw new InvalidOperationException("æµå¼æ¨¡å¼ä¸å…许直接å–包,请使用 ReadAllAsync è¯»å–æˆ– DrainAsync 丢弃");
if (_remaining <= 0) return null;
return _packet!.Slice((Int32)_offset, (Int32)_remaining);
}
/// <summary>读满剩余数æ®ï¼ˆè¯»å–完æˆè¯ä¹‰ï¼‰</summary>
/// <param name="cancellationToken">å–æ¶ˆé€šçŸ¥</param>
/// <returns>剩余体的数æ®åŒ…ï¼›å†…å˜æ¨¡å¼ä¸ºç«‹å³å®Œæˆçš„视图,æµå¼æ¨¡å¼ä¸ºè¯»æ»¡åŽçš„æ± 化包</returns>
/// <remarks>æµå¼æ¨¡å¼æœªè¯»æ»¡å°±é‡å–消或æµç»“æŸæ—¶æŠ›å¼‚常,ä¸è¿”å›žåŠæˆªæ¶ˆæ¯ä½“ï¼›å†…å˜æ¨¡å¼æ•°æ®å·²åœ¨å†…å˜ï¼Œä¸ä¼šçŸè¯»ã€‚</remarks>
/// <exception cref="OperationCanceledException">æµå¼æ¨¡å¼è¯»å–è¢«å–æ¶ˆï¼Œæ¶ˆæ¯ä½“未读满</exception>
/// <exception cref="EndOfStreamException">æµå¼æ¨¡å¼æ•°æ®æµæå‰ç»“æŸï¼Œæ¶ˆæ¯ä½“未读满</exception>
public async ValueTask<IPacket> ReadAllAsync(CancellationToken cancellationToken = default)
{
if (_remaining <= 0) return new OwnerPacket(0);
// å†…å˜æ¨¡å¼ï¼šå–视图并推进至末尾
if (_reader == null)
{
var view = _packet!.Slice((Int32)_offset, (Int32)_remaining);
_offset += _remaining;
_remaining = 0;
return view;
}
// æµå¼æ¨¡å¼ï¼šé¢„算已知,一次分é…读满
if (_remaining > Int32.MaxValue) throw new NotSupportedException($"消æ¯ä½“过大({_remaining} å—èŠ‚ï¼‰ï¼Œæ— æ³•ä¸€æ¬¡æ€§ç‰©åŒ–");
var buffer = new OwnerPacket((Int32)_remaining);
try
{
var total = _remaining;
var memory = buffer.GetMemory();
var offset = 0;
var canceled = false;
while (_remaining > 0)
{
var rr = await _reader.ReadAsync(cancellationToken).ConfigureAwait(false);
if (rr.IsCanceled)
{
canceled = true;
break;
}
var data = rr.Buffer;
if (!data.IsEmpty)
{
var take = (Int32)Math.Min(data.Length, _remaining);
data.Slice(0, take).CopyTo(memory.Span[offset..]);
offset += take;
AdvanceTo(take);
}
else if (rr.IsCompleted) break;
}
// æœªè¯»æ»¡å°±é€€å‡ºï¼šåŠæˆªæ¶ˆæ¯ä½“比异常å±é™©å¾—多(上层åªèƒ½ä»Žé•¿åº¦å¯¹ä¸ä¸ŠçŒœå‡ºé—®é¢˜ï¼‰ï¼ŒæŒ‰åŽŸå› æŠ›ç»™è°ƒç”¨æ–¹
if (_remaining > 0)
{
if (canceled) throw new OperationCanceledException($"消æ¯ä½“读å–è¢«å–æ¶ˆï¼ŒæœŸæœ› {total} å—节,实读 {offset} å—节");
// 管é“å¸¦å¼‚å¸¸ç»“æŸæ—¶é€ä¼ 原始故障(与帧泵的å£å¾„一致),å¦åˆ™è¯´æ˜Žå¯¹ç«¯æå‰å…³é—
if (_reader.Error is { } error) throw error;
throw new EndOfStreamException($"消æ¯ä½“未读满:期望 {total} å—节,实读 {offset} å—节");
}
return buffer.Resize(offset);
}
catch
{
// å¤±è´¥è·¯å¾„å¿…é¡»å½’è¿˜æ± ç¼“å†²ï¼šç•™ç€åªèƒ½é æžæž„兜底(打æ¼é‡Šæ”¾è¦å‘Šï¼Œä¸”归还时机ä¸å®šï¼‰
buffer.TryDispose();
throw;
}
}
#endregion
}
|