Skip to content

Commit 2170911

Browse files
[release/6.0-rc1] Fixed StreamPipeReader.CopyToAsync (#57966)
* Fixed StreamPipeReader.CopyToAsync - Take the segment index into account when copying buffered data. This handles the case where ReadAsync has consumed a partial segment and then the same PipeReader instance is used to copy to a Stream and PipeWriter. - Added tests * Always slice Co-authored-by: David Fowler <davidfowl@gmail.com>
1 parent 0e4a871 commit 2170911

2 files changed

Lines changed: 44 additions & 2 deletions

File tree

src/libraries/System.IO.Pipelines/src/System/IO/Pipelines/StreamPipeReader.cs

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -383,18 +383,21 @@ public override async Task CopyToAsync(PipeWriter destination, CancellationToken
383383
try
384384
{
385385
BufferSegment? segment = _readHead;
386+
int segmentIndex = _readIndex;
387+
386388
try
387389
{
388390
while (segment != null)
389391
{
390-
FlushResult flushResult = await destination.WriteAsync(segment.Memory, tokenSource.Token).ConfigureAwait(false);
392+
FlushResult flushResult = await destination.WriteAsync(segment.Memory.Slice(segmentIndex), tokenSource.Token).ConfigureAwait(false);
391393

392394
if (flushResult.IsCanceled)
393395
{
394396
ThrowHelper.ThrowOperationCanceledException_FlushCanceled();
395397
}
396398

397399
segment = segment.NextSegment;
400+
segmentIndex = 0;
398401

399402
if (flushResult.IsCompleted)
400403
{
@@ -451,13 +454,16 @@ public override async Task CopyToAsync(Stream destination, CancellationToken can
451454
try
452455
{
453456
BufferSegment? segment = _readHead;
457+
int segmentIndex = _readIndex;
458+
454459
try
455460
{
456461
while (segment != null)
457462
{
458-
await destination.WriteAsync(segment.Memory, tokenSource.Token).ConfigureAwait(false);
463+
await destination.WriteAsync(segment.Memory.Slice(segmentIndex), tokenSource.Token).ConfigureAwait(false);
459464

460465
segment = segment.NextSegment;
466+
segmentIndex = 0;
461467
}
462468
}
463469
finally

src/libraries/System.IO.Pipelines/tests/PipeReaderCopyToAsyncTests.cs

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -286,5 +286,41 @@ public async Task ThrowingFromStreamCallsAdvanceToWithStartOfLastReadResult(int
286286
Assert.True(startPosition.Equals(wrappedPipeReader.LastConsumed));
287287
Assert.True(startPosition.Equals(wrappedPipeReader.LastExamined));
288288
}
289+
290+
[Fact]
291+
public async Task CopyToAsyncStreamCopiesRemainderAfterReadingSome()
292+
{
293+
var buffer = Encoding.UTF8.GetBytes("Hello World");
294+
await Pipe.Writer.WriteAsync(buffer);
295+
Pipe.Writer.Complete();
296+
297+
var result = await PipeReader.ReadAsync();
298+
Assert.Equal(result.Buffer.ToArray(), buffer);
299+
// Consume Hello
300+
PipeReader.AdvanceTo(result.Buffer.GetPosition(5));
301+
302+
var ms = new MemoryStream();
303+
await PipeReader.CopyToAsync(ms);
304+
305+
Assert.Equal(buffer.AsMemory(5).ToArray(), ms.ToArray());
306+
}
307+
308+
[Fact]
309+
public async Task CopyToAsyncPipeWriterCopiesRemainderAfterReadingSome()
310+
{
311+
var buffer = Encoding.UTF8.GetBytes("Hello World");
312+
await Pipe.Writer.WriteAsync(buffer);
313+
Pipe.Writer.Complete();
314+
315+
var result = await PipeReader.ReadAsync();
316+
Assert.Equal(result.Buffer.ToArray(), buffer);
317+
// Consume Hello
318+
PipeReader.AdvanceTo(result.Buffer.GetPosition(5));
319+
320+
var ms = new MemoryStream();
321+
await PipeReader.CopyToAsync(PipeWriter.Create(ms));
322+
323+
Assert.Equal(buffer.AsMemory(5).ToArray(), ms.ToArray());
324+
}
289325
}
290326
}

0 commit comments

Comments
 (0)