Provide a better way to manage async disposable streams

This commit is contained in:
Adam Hathcock
2026-08-04 09:11:26 +01:00
parent 369ab91522
commit 62d1b5dcc8
26 changed files with 141 additions and 141 deletions

View File

@@ -108,15 +108,12 @@ public static class IArchiveEntryExtensions
throw new ExtractionException("Entry is a file directory and cannot be extracted.");
}
#if LEGACY_DOTNET
using var entryStream = await archiveEntry
var entryStream = await archiveEntry
.OpenEntryStreamAsync(cancellationToken)
.ConfigureAwait(false);
#else
await using var entryStream = await archiveEntry
.OpenEntryStreamAsync(cancellationToken)
await using var entryStreamScope = entryStream
.DisposeAsyncScope()
.ConfigureAwait(false);
#endif
var checkedStream = options is null
? entryStream
: IEntryExtensions.WrapWithChecksumValidation(archiveEntry, entryStream, options);

View File

@@ -17,7 +17,6 @@ public partial class EntryStream
_completed = true;
}
#if !LEGACY_DOTNET
public override async ValueTask DisposeAsync()
{
if (_isDisposed)
@@ -43,9 +42,8 @@ public partial class EntryStream
}
}
await base.DisposeAsync().ConfigureAwait(false);
await _stream.DisposeAsync().ConfigureAwait(false);
await _stream.DisposeAsyncCompat().ConfigureAwait(false);
}
#endif
public override async Task<int> ReadAsync(
byte[] buffer,

View File

@@ -8,7 +8,7 @@ using SharpCompress.Readers;
namespace SharpCompress.Common;
public partial class EntryStream : Stream
public partial class EntryStream : AsyncDisposableStream
{
private readonly IReader _reader;
private readonly Stream _stream;

View File

@@ -0,0 +1,36 @@
using System;
using System.IO;
using System.Threading.Tasks;
namespace SharpCompress.IO;
/// <summary>
/// A <see cref="Stream"/> that is guaranteed to be asynchronously disposable on every target framework.
/// </summary>
/// <remarks>
/// <para>
/// On .NET Framework 4.8 and .NET Standard 2.0, <see cref="Stream"/> has no <c>DisposeAsync</c>.
/// <c>Microsoft.Bcl.AsyncInterfaces</c> supplies the <see cref="IAsyncDisposable"/> interface on those
/// targets but cannot retrofit it onto the BCL's <see cref="Stream"/>, and C# will not accept an
/// extension method for the pattern - <c>await using</c> requires a reachable <em>instance</em>
/// <c>DisposeAsync</c>. Deriving from this class instead of <see cref="Stream"/> therefore makes a type
/// usable with <c>await using</c> uniformly, with no conditional compilation at the call site.
/// </para>
/// <para>
/// The fallback below is the same behaviour as the BCL's own default <see cref="Stream.DisposeAsync"/>,
/// so a derived type may call <c>await base.DisposeAsync()</c> unconditionally on any target.
/// </para>
/// </remarks>
public abstract class AsyncDisposableStream : Stream
#if NO_STREAM_DISPOSEASYNC
, IAsyncDisposable
#endif
{
#if NO_STREAM_DISPOSEASYNC
public virtual ValueTask DisposeAsync()
{
Dispose();
return default;
}
#endif
}

View File

@@ -0,0 +1,36 @@
using System;
using System.Threading.Tasks;
namespace SharpCompress.IO;
/// <summary>
/// Makes any resource usable with <c>await using</c>, disposing it asynchronously when the runtime type
/// supports it and synchronously otherwise.
/// </summary>
/// <remarks>
/// Needed for locals whose <em>static</em> type is <see cref="System.IO.Stream"/> (or another type that
/// only sometimes has <c>DisposeAsync</c>), where <c>await using</c> cannot bind directly on
/// .NET Framework 4.8 / .NET Standard 2.0. Unlike a compile-time guard, this picks the asynchronous path
/// based on the runtime type, so a stream that really is asynchronously disposable is disposed that way on
/// every target framework. Prefer deriving from <see cref="AsyncDisposableStream"/> where the type is ours.
/// </remarks>
internal readonly struct AsyncDisposeScope(IDisposable? resource) : IAsyncDisposable
{
public ValueTask DisposeAsync()
{
if (resource is IAsyncDisposable asyncDisposable)
{
return asyncDisposable.DisposeAsync();
}
resource?.Dispose();
return default;
}
/// <summary>
/// Mirrors <c>ConfiguredAsyncDisposable</c> so <c>await using</c> can specify context capture without
/// boxing this struct through <see cref="IAsyncDisposable"/>.
/// </summary>
public ConfiguredAsyncDisposeScope ConfigureAwait(bool continueOnCapturedContext) =>
new(resource, continueOnCapturedContext);
}

View File

@@ -0,0 +1,19 @@
using System;
using System.Runtime.CompilerServices;
using System.Threading.Tasks;
namespace SharpCompress.IO;
internal readonly struct ConfiguredAsyncDisposeScope(IDisposable? resource, bool continueOnCapturedContext)
{
public ConfiguredValueTaskAwaitable DisposeAsync()
{
if (resource is IAsyncDisposable asyncDisposable)
{
return asyncDisposable.DisposeAsync().ConfigureAwait(continueOnCapturedContext);
}
resource?.Dispose();
return default(ValueTask).ConfigureAwait(continueOnCapturedContext);
}
}

View File

@@ -25,6 +25,20 @@ public static class StreamExtensions
public void Skip() => stream.CopyTo(Stream.Null);
/// <summary>
/// Returns a scope that disposes this stream when awaited, asynchronously where the runtime type
/// supports it. Lets <c>await using</c> be written against a <see cref="Stream"/>-typed local on
/// every target framework.
/// </summary>
internal AsyncDisposeScope DisposeAsyncScope() => new(stream);
/// <summary>
/// Disposes this stream, asynchronously where the runtime type supports it. Use where the static
/// type is <see cref="Stream"/>, which has no <c>DisposeAsync</c> on .NET Framework 4.8 /
/// .NET Standard 2.0.
/// </summary>
internal ValueTask DisposeAsyncCompat() => new AsyncDisposeScope(stream).DisposeAsync();
public async ValueTask SkipAsync(CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();

View File

@@ -104,13 +104,8 @@ public abstract partial class AbstractReader<TEntry, TVolume>
}
}
//don't know the size so we have to try to decompress to skip
#if LEGACY_DOTNET
using var s = await OpenEntryStreamAsync(cancellationToken).ConfigureAwait(false);
await s.SkipEntryAsync(cancellationToken).ConfigureAwait(false);
#else
await using var s = await OpenEntryStreamAsync(cancellationToken).ConfigureAwait(false);
await s.SkipEntryAsync(cancellationToken).ConfigureAwait(false);
#endif
}
public async ValueTask WriteEntryToAsync(
@@ -139,19 +134,11 @@ public abstract partial class AbstractReader<TEntry, TVolume>
private async ValueTask WriteAsync(Stream writeStream, CancellationToken cancellationToken)
{
#if LEGACY_DOTNET
using Stream s = await OpenEntryStreamAsync(cancellationToken).ConfigureAwait(false);
await using var s = await OpenEntryStreamAsync(cancellationToken).ConfigureAwait(false);
var sourceStream = WrapWithProgress(s, Entry);
await sourceStream
.CopyToAsync(writeStream, Options.BufferSize, cancellationToken)
.ConfigureAwait(false);
#else
await using Stream s = await OpenEntryStreamAsync(cancellationToken).ConfigureAwait(false);
var sourceStream = WrapWithProgress(s, Entry);
await sourceStream
.CopyToAsync(writeStream, Options.BufferSize, cancellationToken)
.ConfigureAwait(false);
#endif
}
public async ValueTask<EntryStream> OpenEntryStreamAsync(

View File

@@ -108,15 +108,9 @@ public static class IAsyncReaderExtensions
CancellationToken cancellationToken
)
{
#if LEGACY_DOTNET
using var entryStream = await reader
.OpenEntryStreamAsync(cancellationToken)
.ConfigureAwait(false);
#else
await using var entryStream = await reader
.OpenEntryStreamAsync(cancellationToken)
.ConfigureAwait(false);
#endif
var checkedStream = IEntryExtensions.WrapWithChecksumValidation(
reader.Entry,
entryStream,

View File

@@ -151,19 +151,8 @@ public class BZip2StreamAsyncTests
Assert.True(compressed.Length > 0);
// Decompress and verify
#if LEGACY_DOTNET
// MemoryStream has nothing to dispose asynchronously
using (var readStream = new MemoryStream(compressed))
{
using (
var bzip2Stream = await BZip2Stream.CreateAsync(
new AsyncOnlyStream(readStream),
SharpCompress.Compressors.CompressionMode.Decompress,
false
)
)
{
#else
await using (var readStream = new MemoryStream(compressed))
{
await using (
var bzip2Stream = await BZip2Stream.CreateAsync(
@@ -173,7 +162,6 @@ public class BZip2StreamAsyncTests
)
)
{
#endif
var result = new StringBuilder();
var buffer = new byte[256];
int bytesRead;

View File

@@ -68,11 +68,8 @@ public class GZipCrcExtractionTests : TestBase
[Fact]
public async Task GZipArchive_WriteToFileAsync_Throws_On_Crc_Mismatch()
{
#if LEGACY_DOTNET
// MemoryStream has nothing to dispose asynchronously
using var stream = new MemoryStream(ReadCorruptedGZipTrailer(corruptCrc: true));
#else
await using var stream = new MemoryStream(ReadCorruptedGZipTrailer(corruptCrc: true));
#endif
await using var archive = await GZipArchive.OpenAsyncArchive(stream);
var entry = await archive.EntriesAsync.SingleAsync();
var destination = Path.Combine(SCRATCH_FILES_PATH, Guid.NewGuid().ToString());

View File

@@ -40,24 +40,15 @@ public class LargeArchiveTests : TestBase
[InlineData("Large/Large.7z")]
public async Task OpenAsyncArchive_ShouldStreamLargeEntry(string fixtureName)
{
#if LEGACY_DOTNET
using var stream = new AsyncOnlyStream(
File.OpenRead(await GetMaterializedFixturePathAsync(fixtureName))
);
#else
await using var stream = new AsyncOnlyStream(
File.OpenRead(await GetMaterializedFixturePathAsync(fixtureName))
);
#endif
await using var archive = await ArchiveFactory.OpenAsyncArchive(stream);
var entry = await GetSingleEntryAsync(archive);
#if LEGACY_DOTNET
using var entryStream = await entry.OpenEntryStreamAsync();
#else
await using var entryStream = await entry.OpenEntryStreamAsync();
#endif
var entryStream = await entry.OpenEntryStreamAsync();
await using var entryStreamScope = entryStream.DisposeAsyncScope();
await VerifyContentAsync(entry.Key, entryStream);
}
@@ -83,25 +74,15 @@ public class LargeArchiveTests : TestBase
[InlineData("Large/Large.tar.gz")]
public async Task OpenAsyncReader_ShouldStreamLargeEntry(string fixtureName)
{
#if LEGACY_DOTNET
using var stream = new AsyncOnlyStream(
File.OpenRead(await GetMaterializedFixturePathAsync(fixtureName))
);
#else
await using var stream = new AsyncOnlyStream(
File.OpenRead(await GetMaterializedFixturePathAsync(fixtureName))
);
#endif
await using var reader = await ReaderFactory.OpenAsyncReader(stream);
Assert.True(await reader.MoveToNextEntryAsync());
Assert.False(reader.Entry.IsDirectory);
#if LEGACY_DOTNET
using var entryStream = await reader.OpenEntryStreamAsync();
#else
await using var entryStream = await reader.OpenEntryStreamAsync();
#endif
await VerifyContentAsync(reader.Entry.Key, entryStream);
Assert.False(await reader.MoveToNextEntryAsync());
}

View File

@@ -2,10 +2,11 @@ using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using SharpCompress.IO;
namespace SharpCompress.Test.Mocks;
public class AsyncOnlyStream(Stream stream, bool disposeStream = true) : Stream
public class AsyncOnlyStream(Stream stream, bool disposeStream = true) : AsyncDisposableStream
{
private readonly Stream _stream = stream ?? throw new ArgumentNullException(nameof(stream));

View File

@@ -1,6 +1,7 @@
using System;
using System.IO;
using System.Threading.Tasks;
using SharpCompress.IO;
namespace SharpCompress.Test.Mocks;
@@ -8,7 +9,7 @@ namespace SharpCompress.Test.Mocks;
// CryptoStream doesn't always trigger the Flush, so this class is used instead
// See https://referencesource.microsoft.com/#mscorlib/system/security/cryptography/cryptostream.cs,141
public class FlushOnDisposeStream(Stream innerStream) : Stream
public class FlushOnDisposeStream(Stream innerStream) : AsyncDisposableStream
{
public override bool CanRead => innerStream.CanRead;
@@ -48,12 +49,10 @@ public class FlushOnDisposeStream(Stream innerStream) : Stream
base.Dispose(disposing);
}
#if !LEGACY_DOTNET
public override async ValueTask DisposeAsync()
{
await innerStream.FlushAsync();
innerStream.Close();
await base.DisposeAsync();
}
#endif
}

View File

@@ -2,6 +2,7 @@ using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using SharpCompress.IO;
namespace SharpCompress.Test.Mocks;
@@ -9,7 +10,7 @@ namespace SharpCompress.Test.Mocks;
/// A forward-only stream wrapper that delegates directly to the underlying stream
/// without any buffering. Supports reading and writing but not seeking.
/// </summary>
public class ForwardOnlyStream : Stream
public class ForwardOnlyStream : AsyncDisposableStream
{
private readonly Stream _stream;
private bool _isDisposed;
@@ -142,17 +143,15 @@ public class ForwardOnlyStream : Stream
}
}
#if !LEGACY_DOTNET
public override async ValueTask DisposeAsync()
{
if (!_isDisposed)
{
await _stream.DisposeAsync();
await _stream.DisposeAsyncCompat();
_isDisposed = true;
}
await base.DisposeAsync();
}
#endif
private void ThrowIfDisposed()
{

View File

@@ -2,10 +2,11 @@
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using SharpCompress.IO;
namespace SharpCompress.Test.Mocks;
public class TestStream(Stream stream, bool read, bool write, bool seek) : Stream
public class TestStream(Stream stream, bool read, bool write, bool seek) : AsyncDisposableStream
{
public TestStream(Stream stream)
: this(stream, stream.CanRead, stream.CanWrite, stream.CanSeek) { }
@@ -50,14 +51,14 @@ public class TestStream(Stream stream, bool read, bool write, bool seek) : Strea
Memory<byte> buffer,
CancellationToken cancellationToken = default
) => stream.ReadAsync(buffer, cancellationToken);
#endif
public override async ValueTask DisposeAsync()
{
await base.DisposeAsync();
await stream.DisposeAsync();
await stream.DisposeAsyncCompat();
IsDisposed = true;
}
#endif
public override long Seek(long offset, SeekOrigin origin) => stream.Seek(offset, origin);

View File

@@ -2,6 +2,7 @@ using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
using SharpCompress.IO;
namespace SharpCompress.Test.Mocks;
@@ -9,7 +10,7 @@ namespace SharpCompress.Test.Mocks;
/// A stream wrapper that throws NotSupportedException on Flush() calls.
/// This is used to test that archive iteration handles streams that don't support flushing.
/// </summary>
public class ThrowOnFlushStream : Stream
public class ThrowOnFlushStream : AsyncDisposableStream
{
private readonly Stream inner;

View File

@@ -1,5 +1,6 @@
using System;
using System.IO;
using SharpCompress.IO;
namespace SharpCompress.Test.Mocks;
@@ -7,7 +8,7 @@ namespace SharpCompress.Test.Mocks;
/// A stream wrapper that truncates the underlying stream after reading a specified number of bytes.
/// Used for testing error handling when streams end prematurely.
/// </summary>
public class TruncatedStream : Stream
public class TruncatedStream : AsyncDisposableStream
{
private readonly Stream baseStream;
private readonly long truncateAfterBytes;

View File

@@ -159,18 +159,12 @@ public abstract class ReaderTests : TestBase
{
using var file = File.OpenRead(testArchive);
#if !LEGACY_DOTNET
await using var protectedStream = SharpCompressStream.CreateNonDisposing(
// SharpCompressStream is not yet an AsyncDisposableStream, so scope its disposal
var protectedStream = SharpCompressStream.CreateNonDisposing(
new ForwardOnlyStream(file, options.BufferSize)
);
await using var protectedStreamScope = protectedStream.DisposeAsyncScope();
await using var testStream = new TestStream(protectedStream);
#else
using var protectedStream = SharpCompressStream.CreateNonDisposing(
new ForwardOnlyStream(file, options.BufferSize)
);
using var testStream = new TestStream(protectedStream);
#endif
await using (
var reader = await ReaderFactory.OpenAsyncReader(
new AsyncOnlyStream(testStream),

View File

@@ -187,10 +187,6 @@ public class TarArchiveAsyncTests : ArchiveTests
}
}
}
#if LEGACY_DOTNET
//add a delay because old .net sucks on DisposeAsync
await Task.Delay(TimeSpan.FromSeconds(1));
#endif
}
[Fact]
@@ -338,11 +334,8 @@ public class TarArchiveAsyncTests : ArchiveTests
{
++numberOfEntries;
#if LEGACY_DOTNET
using var tarEntryStream = await entry.OpenEntryStreamAsync();
#else
await using var tarEntryStream = await entry.OpenEntryStreamAsync();
#endif
var tarEntryStream = await entry.OpenEntryStreamAsync();
await using var tarEntryStreamScope = tarEntryStream.DisposeAsyncScope();
using var testFileStream = new MemoryStream();
await tarEntryStream.CopyToAsync(testFileStream);
Assert.Equal(testBytes.Length, testFileStream.Length);

View File

@@ -337,13 +337,8 @@ public class TarReaderAsyncTests : ReaderTests
Assert.True(await reader.MoveToNextEntryAsync());
Assert.Equal("inner.tar.gz", reader.Entry.Key);
#if !LEGACY_DOTNET
await using var entryStream = await reader.OpenEntryStreamAsync();
await using var flushingStream = new FlushOnDisposeStream(entryStream);
#else
using var entryStream = await reader.OpenEntryStreamAsync();
using var flushingStream = new FlushOnDisposeStream(entryStream);
#endif
// Extract inner.tar.gz
await using var innerReader = await ReaderFactory.OpenAsyncReader(flushingStream);

View File

@@ -89,11 +89,8 @@ public class TarWriterNonSeekableTests
{
var entry = fileEntries.Single(e => e.Key == name);
using var extracted = new MemoryStream();
#if LEGACY_DOTNET
using (var entryStream = await entry.OpenEntryStreamAsync())
#else
await using (var entryStream = await entry.OpenEntryStreamAsync())
#endif
var entryStream = await entry.OpenEntryStreamAsync();
await using (entryStream.DisposeAsyncScope())
{
await entryStream.CopyToAsync(extracted);
}

View File

@@ -198,11 +198,7 @@ public class Zip64AsyncTests : WriterTests
count++;
lastKey = rd.Entry.Key;
#if LEGACY_DOTNET
using var entryStream = await rd.OpenEntryStreamAsync();
#else
await using var entryStream = await rd.OpenEntryStreamAsync();
#endif
if (rd.Entry.Key == "small")
{
using var ms = new MemoryStream();
@@ -337,17 +333,10 @@ public class Zip64AsyncTests : WriterTests
);
while (await rd.MoveToNextEntryAsync())
{
#if LEGACY_DOTNET
using (var entryStream = await rd.OpenEntryStreamAsync())
{
await entryStream.SkipEntryAsync();
}
#else
await using (var entryStream = await rd.OpenEntryStreamAsync())
{
await entryStream.SkipEntryAsync();
}
#endif
count++;
if (prev != null)
{

View File

@@ -122,11 +122,8 @@ public class ZipCrcExtractionTests : ArchiveTests
using var zipStream = CreateZipWithInvalidCrc(useDataDescriptor: false);
using var archive = ZipArchive.OpenArchive(zipStream);
var entry = archive.Entries.Single(e => !e.IsDirectory);
#if LEGACY_DOTNET
// MemoryStream has nothing to dispose asynchronously
using var destination = new MemoryStream();
#else
await using var destination = new MemoryStream();
#endif
var exception = await Assert.ThrowsAsync<InvalidFormatException>(async () =>
await entry.WriteToAsync(destination, new ExtractionOptions { CheckCrc = true })
@@ -141,11 +138,8 @@ public class ZipCrcExtractionTests : ArchiveTests
using var zipStream = CreateZipWithInvalidCrc(useDataDescriptor: false);
using var archive = ZipArchive.OpenArchive(zipStream);
var entry = archive.Entries.Single(e => !e.IsDirectory);
#if LEGACY_DOTNET
// MemoryStream has nothing to dispose asynchronously
using var destination = new MemoryStream();
#else
await using var destination = new MemoryStream();
#endif
await entry.WriteToAsync(destination, new ExtractionOptions { CheckCrc = false });

View File

@@ -315,11 +315,7 @@ public class ZipReaderAsyncTests : ReaderTests
{
if (!reader.Entry.IsDirectory)
{
#if LEGACY_DOTNET
using var entryStream = await reader.OpenEntryStreamAsync();
#else
await using var entryStream = await reader.OpenEntryStreamAsync();
#endif
// Read some data
var buffer = new byte[1024];
await entryStream.ReadAsync(buffer, 0, buffer.Length);
@@ -342,11 +338,7 @@ public class ZipReaderAsyncTests : ReaderTests
{
if (!reader.Entry.IsDirectory)
{
#if LEGACY_DOTNET
using var entryStream = await reader.OpenEntryStreamAsync();
#else
await using var entryStream = await reader.OpenEntryStreamAsync();
#endif
// Read some data
var buffer = new byte[1024];
await entryStream.ReadAsync(buffer, 0, buffer.Length);

View File

@@ -163,11 +163,8 @@ public class ZipWriterNonSeekableTests
{
var entry = archive.Entries.Single(e => e.Key == name);
using var extracted = new MemoryStream();
#if LEGACY_DOTNET
using (var entryStream = await entry.OpenEntryStreamAsync())
#else
await using (var entryStream = await entry.OpenEntryStreamAsync())
#endif
var entryStream = await entry.OpenEntryStreamAsync();
await using (entryStream.DisposeAsyncScope())
{
await entryStream.CopyToAsync(extracted);
}