16 static readonly UnboundedChannelOptions unboundedChannelOptions =
new UnboundedChannelOptions()
20 AllowSynchronousContinuations =
true
23 readonly IBufferWriter<byte> _writer;
24 readonly Channel<(IMemoryOwner<byte> Owner,
int Length)> _channel = Channel.CreateUnbounded<(IMemoryOwner<byte> Owner,
int Length)>(unboundedChannelOptions);
27 CancellationTokenSource? _stop;
36 _writer = writer ??
throw new ArgumentNullException(nameof(writer));
43 public void Write(IMemoryOwner<byte> owner,
int length)
46 throw new ArgumentNullException(nameof(owner));
48 if (_stop ==
null || _task ==
null)
52 if (_stop ==
null || _task ==
null)
54 _stop =
new CancellationTokenSource();
55 _task = DequeueLoop(_stop.Token);
60 _channel.Writer.TryWrite((owner, length));
66 async Task DequeueLoop(CancellationToken cancellationToken)
68 while (cancellationToken.IsCancellationRequested ==
false)
72 while (await _channel.Reader.WaitToReadAsync(cancellationToken))
73 while (_channel.Reader.TryRead(out var item))
74 WriteData(item.Owner, item.Length);
76 catch (OperationCanceledException)
88 void WriteData(IMemoryOwner<byte> owner,
int length)
91 var buffer = _writer.GetMemory(length);
92 owner.Memory.Slice(0, length).CopyTo(buffer);
93 _writer.Advance(length);
108 _task?.GetAwaiter().GetResult();
113 while (_channel.Reader.TryRead(out var item))
114 item.Owner.Dispose();
117 _channel.Writer.TryComplete();