namespace NewLife.Data;
/// <summary>æ•°æ®åŒ…管é“。连接接收方与消费方的å—节æµç¼“å†²ï¼Œå†™ä¾§è¿½åŠ æ•°æ®ã€è¯»ä¾§æŒ‰åºåˆ—消费</summary>
/// <remarks>
/// <para><b>å¯¹é½ BCL</b>:命å与形æ€å¯¹é½ System.IO.Pipelines çš„ <c>Pipe</c>(åŒåä¸åŒå‘½åç©ºé—´ï¼›ä¸¤è€…éœ€åŒæ—¶å¼•用时用别å,如 <c>using NlPipe = NewLife.Data.Pipe;</c>)。本类承载“管é“级â€å…³æ³¨ç‚¹â€”—åŒç«¯å¥æŸ„装é…ã€å…±äº«åŒæ¥é”与关é—状æ€ã€æ°´ä½èƒŒåŽ‹ä¸Žæ¢å¤äº‹ä»¶ï¼›è¯»å†™è¡Œä¸ºåˆ†åˆ«åœ¨ <see cref="PipeReader"/> 与 <see cref="PipeWriter"/> ä¸Šå®žçŽ°ã€‚è‡ªç ”åŠ¨æœºï¼šBCL çš„ System.IO.Pipelines 最低 netstandard2.0(net45 æ— æ³•å¼•ç”¨ï¼‰ï¼Œæœ¬åº“éœ€è¦åŒä¸€å¥—管é“å½¢æ€è¦†ç›–å« net45 åœ¨å†…çš„å…¨éƒ¨ç›®æ ‡æ¡†æž¶ã€‚</para>
/// <para><b>模型</b>:未消费数æ®ä»¥æ®µé“¾æŒæœ‰åœ¨ <see cref="Reader"/>â€”â€”æ¯æ®µåŒæ—¶æ‰¿è½½åºåˆ—内å˜ï¼ˆå¯¹å¤–暴露 <see cref="System.Buffers.ReadOnlySequence{T}"/> é›¶æ‹·è´çª—å£ï¼‰ä¸Žæ‹¥æœ‰å¥æŸ„ï¼ˆæ¶ˆè´¹æŽ¨è¿›æ—¶åŒæ¥å½’è¿˜ï¼‰ï¼›ä»»æ„æ—¶åˆ»è¿½åŠ æ•°æ®å³å¯å”¤é†’挂起的读å–。</para>
/// <para><b>背压</b>:未消费数æ®è¾¾åˆ° <see cref="PauseThreshold"/> æ—¶ <see cref="IsPaused"/> 为 true,接收方应暂åœç»§ç»æŽ¥æ”¶ï¼›å†™ä¾§æäº¤ï¼ˆ<see cref="PipeWriter.FlushAsync(CancellationToken)"/> è¾¾åˆ°æš‚åœæ°´ä½å³æŒ‚èµ·ï¼Œå¯¹é½ BCL)亦å¯åœ¨æ¤æŒ‚èµ·ç‰å¾…,形æˆåŒå‘背压。消费推进到 <see cref="ResumeThreshold"/> ä»¥ä¸‹æ—¶è§¦å‘ <see cref="Resumed"/> 并唤醒挂起æäº¤ï¼Œæ¢å¤æŽ¥æ”¶ã€‚缓冲有界,消费者ä¸å–æ•°åˆ™æŽ¥æ”¶æ–¹åœæ¢æ‹‰åŠ¨ï¼ˆTCP 窗å£è‡ªç„¶å›žåŽ‹ï¼‰ã€‚èƒŒåŽ‹è¿‘ä¹Žé›¶å¼€é”€ï¼šæœªè§¦å‘æš‚åœæ—¶çš„æäº¤å¿«è·¯å¾„基准实测 28nsã€é›¶åˆ†é…â€”â€”åªæœ‰çœŸæ£ç§¯åŽ‹åˆ°æš‚åœæ°´ä½æ‰ä»˜å‡ºæŒ‚èµ·/唤醒。</para>
/// <para><b>线程模型</b>:å•写(<see cref="Writer"/> 的方法)å•读(<see cref="Reader"/> çš„æ–¹æ³•ï¼‰ï¼Œè¯»å†™å¯æ¥è‡ªä¸åŒçº¿ç¨‹ã€‚</para>
/// <para><b>所有æƒ</b>:<see cref="PipeWriter.Append(IPacket)"/> æ— æ¡ä»¶æŽ¥ç®¡å…¥å‚奿Ÿ„(管é“已关闿—¶ç”±ç®¡é“负责释放);消费方读å–çš„åºåˆ—仅在对应数æ®è¢«æ¶ˆè´¹å‰æœ‰æ•ˆã€‚</para>
/// <para><b>与 BCL 的差异</b>:消费推进以å—节计数 <see cref="PipeReader.AdvanceTo(Int64)"/> ä¸ºä¸»â€”â€”è¿½åŠ æ•°æ®ä¼šé‡å»ºçª—å£åºåˆ—,ä½ç½®åœ¨è·¨è¿½åŠ åœºæ™¯ä¸ç¨³å®šï¼›å—节计数对跨段ã€è·¨è½®åœºæ™¯ç®€å•å¯é ,且å¯åœ¨å…¨éƒ¨ç›®æ ‡æ¡†æž¶å®žçŽ°ï¼ˆå« net45ï¼‰ï¼Œå¦æä¾› SequencePosition é‡è½½ï¼ˆå–自最近一次读å–窗å£ï¼Œå½¢æ€å¯¹é½ï¼‰ã€‚<see cref="PipeWriter.Append(IPacket)"/>ã€<see cref="PipeReader.TakeFrame(Int64)"/>ã€<see cref="PipeReader.Limit(Int64)"/> 与å—节计数推进为本库扩展,差异清å•è§ã€Šæ•°æ®ç®¡é“Pipe》。</para>
/// </remarks>
/// <example>
/// <code>
/// var pipe = new Pipe();
/// pipe.Writer.Append(received); // 接收线程投递(所有æƒè½¬ç§»ï¼‰
/// var rr = await pipe.Reader.ReadAsync(); // 消费线程读å–
/// if (!rr.Buffer.IsEmpty) { /* è§£æžçª—å£ */ pipe.Reader.AdvanceTo(consumedBytes); }
/// pipe.Writer.Complete(); // 写入结æŸï¼ˆå”¤é†’挂起读)
/// </code>
/// </example>
public sealed class Pipe : IDisposable
{
#region 属性
/// <summary>è¯»ä¾§å¥æŸ„ã€‚æ¶ˆè´¹æ–¹ä»Žè¿™é‡Œè¯»å–æ•°æ®ï¼ˆå¯¹æ ‡ BCL çš„ PipeReader)</summary>
public PipeReader Reader { get; }
/// <summary>写侧奿Ÿ„。生产方从这里写入数æ®ï¼ˆå¯¹æ ‡ BCL çš„ PipeWriter)</summary>
public PipeWriter Writer { get; }
/// <summary>æš‚åœæ°´ä½ï¼ˆå—节)。未检查数æ®è¾¾åˆ°è¯¥å€¼æ—¶ <see cref="IsPaused"/> 为 trueï¼ŒæŽ¥æ”¶æ–¹åº”æš‚åœæŽ¥æ”¶ã€å†™ä¾§æäº¤å¯æŒ‚èµ·ç‰å¾…ï¼›0或负数ä¸å¯ç”¨èƒŒåŽ‹ã€‚é»˜è®¤1M</summary>
/// <remarks>è®°è´¦é‡ä¸ºâ€œæœªæ£€æŸ¥æ•°æ® = <see cref="UnconsumedLength"/> − 已检查å—节数â€ï¼Œå£å¾„å¯¹é½ System.IO.Pipelines:
/// 读侧 <see cref="PipeReader.AdvanceTo(Int64, Int64)"/> 声明为已检查的å—节å³è§£é™¤è®¡å…¥ï¼Œå³ä½¿å°šæœªæ¶ˆè´¹â€”—背压约æŸçš„æ˜¯æ¶ˆè´¹è€…还没看过的积压。
/// å•å‚ <see cref="PipeReader.AdvanceTo(Int64)"/> ç‰ä»· examined=consumed,故常规消费推进的暂åœ/æ¢å¤è¡Œä¸ºä¸Žâ€œæŒ‰æœªæ¶ˆè´¹è®°è´¦â€å®Œå…¨ä¸€è‡´ã€‚</remarks>
public Int64 PauseThreshold { get; set; } = 1024 * 1024;
/// <summary>æ¢å¤æ°´ä½ï¼ˆå—节)。未检查数æ®é™åˆ°è¯¥å€¼ä»¥ä¸‹æ—¶è§£é™¤æš‚åœå¹¶è§¦å‘ <see cref="Resumed"/>。默认512K</summary>
/// <remarks>è®°è´¦å£å¾„è§ <see cref="PauseThreshold"/>。</remarks>
public Int64 ResumeThreshold { get; set; } = 512 * 1024;
/// <summary>未消费数æ®é•¿åº¦ã€‚æ¥è‡ªè¯»ä¾§æ®µé“¾</summary>
public Int64 UnconsumedLength => Reader.UnconsumedLength;
/// <summary>是å¦å·²è¾¾æš‚åœæ°´ä½ã€‚达到åŽä¿æŒï¼Œç›´åˆ°æ¶ˆè´¹é™åˆ°æ¢å¤æ°´ä½ä»¥ä¸‹æ‰è§£é™¤</summary>
public Boolean IsPaused => PauseThreshold > 0 && _paused;
/// <summary>写侧是å¦å·²å®Œæˆ</summary>
public Boolean IsCompleted => WriterCompleted;
/// <summary>å®Œæˆæ—¶çš„异常。读写两侧的 Complete(error) å‡å¯æºå¸¦ï¼Œå…ˆåˆ°å…ˆå¾—(已有异常时ä¸è¢«ç©ºå€¼è¦†ç›–)</summary>
public Exception? Error { get; internal set; }
/// <summary>消费推进使未消费数æ®é™åˆ°æ¢å¤æ°´ä½ä»¥ä¸‹ã€æˆ–读侧挂起ç‰å¾…触å‘é¥¥é¥¿è®©ä½æ—¶è§¦å‘ï¼ŒæŽ¥æ”¶æ–¹å¯æ¢å¤æŽ¥æ”¶</summary>
public event EventHandler? Resumed;
#endregion
#region 共享状æ€ï¼ˆè¯»å†™ä¸¤ä¾§å…±ç”¨ï¼›çжæ€å†™å…¥ä¸€å¾‹åœ¨ SyncRoot ä¿æŠ¤ä¸‹ï¼‰
/// <summary>å…±äº«åŒæ¥é”。读写两侧的状æ€å†™å…¥å…±ç”¨è¿™ä¸€æŠŠé”</summary>
internal readonly Object SyncRoot = new();
/// <summary>写侧是å¦å·²å®Œæˆã€‚写入在 SyncRoot ä¿æŠ¤ä¸‹ï¼›volatile ä¿è¯æ— é”快照读(IsCompleted)的跨线程å¯è§æ€§</summary>
internal volatile Boolean WriterCompleted;
/// <summary>读侧是å¦å·²ç»“æŸã€‚写入在 SyncRoot ä¿æŠ¤ä¸‹</summary>
internal volatile Boolean ReaderCompleted;
/// <summary>æš‚åœæ€ã€‚è¾¾åˆ°æš‚åœæ°´ä½åŽç½®ä½ï¼Œæ¶ˆè´¹é™åˆ°æ¢å¤æ°´ä½ä»¥ä¸‹æ‰è§£é™¤ï¼ˆè¿Ÿæ»žï¼‰ï¼›å†™å…¥åœ¨ SyncRoot ä¿æŠ¤ä¸‹ï¼Œvolatile ä¿è¯æ— é”快照读(IsPaused)的跨线程å¯è§æ€§</summary>
private volatile Boolean _paused;
#endregion
#region æž„é€
/// <summary>创建数æ®åŒ…管é“</summary>
public Pipe()
{
Reader = new PipeReader(this);
Writer = new PipeWriter(this);
}
#endregion
#region 方法
/// <summary>释放管é“。ç‰ä»·äºŽå®Œæˆå†™å…¥å¹¶ç»“æŸè¯»å–,全部未消费数æ®å½’è¿˜å†…å˜æ± </summary>
public void Dispose()
{
Writer.Complete();
Reader.Complete();
}
/// <summary>å¤ä½ç®¡é“,供对象å¤ç”¨ã€‚è¦æ±‚读写两端å‡å·² <see cref="PipeWriter.Complete(Exception?)"/></summary>
/// <remarks>å¯¹é½ System.IO.Pipelines çš„ Pipe.Resetï¼›å¤ä½åŽè¯»å†™å¥æŸ„ä¿æŒæœ‰æ•ˆã€å®Œæˆçжæ€ä¸Žæ°´ä½æ¸…é›¶ã€é”™è¯¯æ¸…除。</remarks>
/// <exception cref="InvalidOperationException">读侧或写侧尚未完æˆ</exception>
public void Reset()
{
lock (SyncRoot)
{
if (!WriterCompleted || !ReaderCompleted) throw new InvalidOperationException("Pipe must be completed on both sides before reset.");
WriterCompleted = false;
ReaderCompleted = false;
_paused = false;
Error = null;
Reader.ResetForReuse();
}
}
/// <summary>刷新暂åœçжæ€ï¼ˆè°ƒç”¨æ–¹æŒé”)。返回本次是å¦éœ€è¦è§¦å‘æ¢å¤äº‹ä»¶</summary>
/// <param name="pending">当剿œªæ£€æŸ¥æ•°æ®é•¿åº¦ï¼ˆæœªæ¶ˆè´¹ − 已检查)</param>
/// <remarks>è¿Ÿæ»žï¼šè¾¾åˆ°æš‚åœæ°´ä½è½¬å…¥æš‚åœæ€åŽä¿æŒï¼Œç›´åˆ°é™å›žæ¢å¤æ°´ä½ä»¥ä¸‹æ‰è§£é™¤ã€‚
/// ä¸èƒ½æŒ‰çž¬æ—¶é•¿åº¦åˆ¤æ–â€”â€”å¤šæ¬¡å°æ¥æ¶ˆè´¹æ—¶ï¼Œè·¨è¿‡æ¢å¤æ°´ä½çš„é‚£ä¸€æ¬¡æŽ¨è¿›çš„èµ·å§‹é•¿åº¦å·²ä½ŽäºŽæš‚åœæ°´ä½ï¼Œçž¬æ—¶åˆ¤æ–ä¼šæ¼æŠ¥æ¢å¤ã€‚</remarks>
internal Boolean UpdatePauseLocked(Int64 pending)
{
if (PauseThreshold <= 0)
{
_paused = false;
return false;
}
if (_paused)
{
// 已暂åœï¼šé™åˆ°æ¢å¤æ°´ä½ä»¥ä¸‹æ—¶è§£é™¤å¹¶æŠ¥å‘Š
if (pending < ResumeThreshold)
{
_paused = false;
return true;
}
}
else if (pending >= PauseThreshold)
{
_paused = true;
}
return false;
}
/// <summary>é‡ç½®æš‚åœæ€ï¼ˆè°ƒç”¨æ–¹æŒé”ï¼‰ã€‚è¯»ä¾§ç»“æŸæ—¶è°ƒç”¨ï¼Œä¸å†è§¦å‘æ¢å¤äº‹ä»¶</summary>
internal void ResetPauseLocked() => _paused = false;
/// <summary>读饥饿让ä½ï¼ˆè°ƒç”¨æ–¹æŒé”)。返回本次是å¦éœ€è¦è§¦å‘æ¢å¤äº‹ä»¶</summary>
/// <remarks>è¯»ä¾§åœ¨â€œæ— æ–°æ•°æ®å¯äº¤ä»˜â€æ—¶æ‰ä¼šæŒ‚èµ·ç‰å¾…ï¼šæ¤æ—¶ä¸ä¼šå†æœ‰æ¶ˆè´¹ã€æš‚åœå·²ä¸å¯èƒ½æŒ‰å¸¸è§„路径(消费é™åŽ‹ï¼‰è§£é™¤ï¼Œ
/// è‹¥ç»§ç»æŒæœ‰ï¼ŒæŽ¥æ”¶æ–¹å°†åœæ‘†ï¼Œè¯»è€…永远ç‰ä¸åˆ°åŽç»å—节(整帧/最å°é•¿åº¦è¯»å–æ»é”)。故挂起å‰è§£é™¤æš‚åœï¼Œç”±è°ƒç”¨æ–¹è§¦å‘ <see cref="Resumed"/> 放行接收。
/// 按“未检查â€è®°è´¦åŽï¼Œè¯»è€…进入ç‰å¾…时未检查é‡å¿…ç„¶ä¸å¤§äºŽ 0ã€æš‚åœå·²è‡ªè¡Œè§£é™¤ï¼Œè¿™é‡Œä½œä¸ºå…œåº•ä¿ç•™ï¼š
/// è¿è¡ŒæœŸè°ƒä½Ž <see cref="PauseThreshold"/> ç‰é…置使暂åœåœ¨è¯»è€…ç‰å¾…æœŸé—´é‡æ–°æˆç«‹æ—¶ï¼Œä»æ®æ¤æ”¾è¡ŒæŽ¥æ”¶ã€‚</remarks>
internal Boolean ReleasePauseForReaderLocked()
{
if (!_paused) return false;
_paused = false;
return true;
}
/// <summary>è§¦å‘æ¢å¤ï¼ˆé”外调用):唤醒挂起的写侧æäº¤å¹¶è§¦å‘æ¢å¤äº‹ä»¶ï¼Œé¿å…用户代ç 进入é”内</summary>
internal void RaiseResumed()
{
TaskCompletionSource<FlushResult>? waiter;
CancellationTokenRegistration reg;
lock (SyncRoot)
{
waiter = Writer.TakeFlushWaiterLocked(out reg);
}
PipeWriter.NotifyFlushWaiter(waiter, reg, new FlushResult(IsCompleted, false));
Resumed?.Invoke(this, EventArgs.Empty);
}
#endregion
}
|