From beeb37b4fd3e64eacb0ea2627e234ff4993adaa9 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Mon, 27 Oct 2025 09:11:29 +0000
Subject: [PATCH] Add async support to EntryStream, ZlibStream, and
ZlibBaseStream
Co-authored-by: adamhathcock <527620+adamhathcock@users.noreply.github.com>
---
src/SharpCompress/Common/EntryStream.cs | 77 ++++++
.../Compressors/Deflate/ZlibBaseStream.cs | 227 ++++++++++++++++++
.../Compressors/Deflate/ZlibStream.cs | 86 +++++++
src/SharpCompress/Utility.cs | 22 ++
4 files changed, 412 insertions(+)
diff --git a/src/SharpCompress/Common/EntryStream.cs b/src/SharpCompress/Common/EntryStream.cs
index a0fe736a..cfaef0aa 100644
--- a/src/SharpCompress/Common/EntryStream.cs
+++ b/src/SharpCompress/Common/EntryStream.cs
@@ -1,6 +1,8 @@
using System;
using System.IO;
using System.IO.Compression;
+using System.Threading;
+using System.Threading.Tasks;
using SharpCompress.IO;
using SharpCompress.Readers;
@@ -51,6 +53,15 @@ public class EntryStream : Stream, IStreamStack
_completed = true;
}
+ ///
+ /// Asynchronously skip the rest of the entry stream.
+ ///
+ public async Task SkipEntryAsync(CancellationToken cancellationToken = default)
+ {
+ await this.SkipAsync(cancellationToken).ConfigureAwait(false);
+ _completed = true;
+ }
+
protected override void Dispose(bool disposing)
{
if (!(_completed || _reader.Cancelled))
@@ -83,6 +94,40 @@ public class EntryStream : Stream, IStreamStack
_stream.Dispose();
}
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override async ValueTask DisposeAsync()
+ {
+ if (!(_completed || _reader.Cancelled))
+ {
+ await SkipEntryAsync().ConfigureAwait(false);
+ }
+
+ //Need a safe standard approach to this - it's okay for compression to overreads. Handling needs to be standardised
+ if (_stream is IStreamStack ss)
+ {
+ if (ss.BaseStream() is SharpCompress.Compressors.Deflate.DeflateStream deflateStream)
+ {
+ await deflateStream.FlushAsync().ConfigureAwait(false);
+ }
+ else if (ss.BaseStream() is SharpCompress.Compressors.LZMA.LzmaStream lzmaStream)
+ {
+ await lzmaStream.FlushAsync().ConfigureAwait(false);
+ }
+ }
+
+ if (_isDisposed)
+ {
+ return;
+ }
+ _isDisposed = true;
+#if DEBUG_STREAMS
+ this.DebugDispose(typeof(EntryStream));
+#endif
+ await base.DisposeAsync().ConfigureAwait(false);
+ await _stream.DisposeAsync().ConfigureAwait(false);
+ }
+#endif
+
public override bool CanRead => true;
public override bool CanSeek => false;
@@ -91,6 +136,8 @@ public class EntryStream : Stream, IStreamStack
public override void Flush() { }
+ public override Task FlushAsync(CancellationToken cancellationToken) => Task.CompletedTask;
+
public override long Length => _stream.Length;
public override long Position
@@ -109,6 +156,36 @@ public class EntryStream : Stream, IStreamStack
return read;
}
+ public override async Task ReadAsync(
+ byte[] buffer,
+ int offset,
+ int count,
+ CancellationToken cancellationToken
+ )
+ {
+ var read = await _stream.ReadAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
+ if (read <= 0)
+ {
+ _completed = true;
+ }
+ return read;
+ }
+
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override async ValueTask ReadAsync(
+ Memory buffer,
+ CancellationToken cancellationToken = default
+ )
+ {
+ var read = await _stream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
+ if (read <= 0)
+ {
+ _completed = true;
+ }
+ return read;
+ }
+#endif
+
public override int ReadByte()
{
var value = _stream.ReadByte();
diff --git a/src/SharpCompress/Compressors/Deflate/ZlibBaseStream.cs b/src/SharpCompress/Compressors/Deflate/ZlibBaseStream.cs
index 155a3556..52442f90 100644
--- a/src/SharpCompress/Compressors/Deflate/ZlibBaseStream.cs
+++ b/src/SharpCompress/Compressors/Deflate/ZlibBaseStream.cs
@@ -31,6 +31,8 @@ using System.Buffers.Binary;
using System.Collections.Generic;
using System.IO;
using System.Text;
+using System.Threading;
+using System.Threading.Tasks;
using SharpCompress.Common.Tar.Headers;
using SharpCompress.IO;
@@ -197,6 +199,57 @@ internal class ZlibBaseStream : Stream, IStreamStack
} while (!done);
}
+ public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
+ {
+ // workitem 7159
+ // calculate the CRC on the unccompressed data (before writing)
+ if (crc != null)
+ {
+ crc.SlurpBlock(buffer, offset, count);
+ }
+
+ if (_streamMode == StreamMode.Undefined)
+ {
+ _streamMode = StreamMode.Writer;
+ }
+ else if (_streamMode != StreamMode.Writer)
+ {
+ throw new ZlibException("Cannot Write after Reading.");
+ }
+
+ if (count == 0)
+ {
+ return;
+ }
+
+ // first reference of z property will initialize the private var _z
+ z.InputBuffer = buffer;
+ _z.NextIn = offset;
+ _z.AvailableBytesIn = count;
+ var done = false;
+ do
+ {
+ _z.OutputBuffer = workingBuffer;
+ _z.NextOut = 0;
+ _z.AvailableBytesOut = _workingBuffer.Length;
+ var rc = (_wantCompress) ? _z.Deflate(_flushMode) : _z.Inflate(_flushMode);
+ if (rc != ZlibConstants.Z_OK && rc != ZlibConstants.Z_STREAM_END)
+ {
+ throw new ZlibException((_wantCompress ? "de" : "in") + "flating: " + _z.Message);
+ }
+
+ await _stream.WriteAsync(_workingBuffer, 0, _workingBuffer.Length - _z.AvailableBytesOut, cancellationToken).ConfigureAwait(false);
+
+ done = _z.AvailableBytesIn == 0 && _z.AvailableBytesOut != 0;
+
+ // If GZIP and de-compress, we're done when 8 bytes remain.
+ if (_flavor == ZlibStreamFlavor.GZIP && !_wantCompress)
+ {
+ done = (_z.AvailableBytesIn == 8 && _z.AvailableBytesOut != 0);
+ }
+ } while (!done);
+ }
+
private void finish()
{
if (_z is null)
@@ -335,6 +388,104 @@ internal class ZlibBaseStream : Stream, IStreamStack
}
}
+ private async Task finishAsync(CancellationToken cancellationToken = default)
+ {
+ if (_z is null)
+ {
+ return;
+ }
+
+ if (_streamMode == StreamMode.Writer)
+ {
+ var done = false;
+ do
+ {
+ _z.OutputBuffer = workingBuffer;
+ _z.NextOut = 0;
+ _z.AvailableBytesOut = _workingBuffer.Length;
+ var rc =
+ (_wantCompress) ? _z.Deflate(FlushType.Finish) : _z.Inflate(FlushType.Finish);
+
+ if (rc != ZlibConstants.Z_STREAM_END && rc != ZlibConstants.Z_OK)
+ {
+ var verb = (_wantCompress ? "de" : "in") + "flating";
+ if (_z.Message is null)
+ {
+ throw new ZlibException(String.Format("{0}: (rc = {1})", verb, rc));
+ }
+ throw new ZlibException(verb + ": " + _z.Message);
+ }
+
+ if (_workingBuffer.Length - _z.AvailableBytesOut > 0)
+ {
+ await _stream.WriteAsync(_workingBuffer, 0, _workingBuffer.Length - _z.AvailableBytesOut, cancellationToken).ConfigureAwait(false);
+ }
+
+ done = _z.AvailableBytesIn == 0 && _z.AvailableBytesOut != 0;
+
+ // If GZIP and de-compress, we're done when 8 bytes remain.
+ if (_flavor == ZlibStreamFlavor.GZIP && !_wantCompress)
+ {
+ done = (_z.AvailableBytesIn == 8 && _z.AvailableBytesOut != 0);
+ }
+ } while (!done);
+
+ await FlushAsync(cancellationToken).ConfigureAwait(false);
+
+ // workitem 7159
+ if (_flavor == ZlibStreamFlavor.GZIP)
+ {
+ if (_wantCompress)
+ {
+ // Emit the GZIP trailer: CRC32 and size mod 2^32
+ byte[] intBuf = new byte[4];
+ BinaryPrimitives.WriteInt32LittleEndian(intBuf, crc.Crc32Result);
+ await _stream.WriteAsync(intBuf, 0, 4, cancellationToken).ConfigureAwait(false);
+ var c2 = (int)(crc.TotalBytesRead & 0x00000000FFFFFFFF);
+ BinaryPrimitives.WriteInt32LittleEndian(intBuf, c2);
+ await _stream.WriteAsync(intBuf, 0, 4, cancellationToken).ConfigureAwait(false);
+ }
+ else
+ {
+ throw new ZlibException("Writing with decompression is not supported.");
+ }
+ }
+ }
+ // workitem 7159
+ else if (_streamMode == StreamMode.Reader)
+ {
+ if (_flavor == ZlibStreamFlavor.GZIP)
+ {
+ if (!_wantCompress)
+ {
+ // workitem 8501: handle edge case (decompress empty stream)
+ if (_z.TotalBytesOut == 0L)
+ {
+ return;
+ }
+
+ // Read and potentially verify the GZIP trailer: CRC32 and size mod 2^32
+ byte[] trailer = new byte[8];
+
+ // workitem 8679
+ if (_z.AvailableBytesIn != 8)
+ {
+ // Make sure we have read to the end of the stream
+ _z.InputBuffer.AsSpan(_z.NextIn, _z.AvailableBytesIn).CopyTo(trailer);
+ var bytesNeeded = 8 - _z.AvailableBytesIn;
+ var bytesRead = await _stream.ReadAsync(
+ trailer, _z.AvailableBytesIn, bytesNeeded, cancellationToken
+ ).ConfigureAwait(false);
+ }
+ }
+ else
+ {
+ throw new ZlibException("Reading with compression is not supported.");
+ }
+ }
+ }
+ }
+
private void end()
{
if (z is null)
@@ -382,6 +533,38 @@ internal class ZlibBaseStream : Stream, IStreamStack
}
}
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override async ValueTask DisposeAsync()
+ {
+ if (isDisposed)
+ {
+ return;
+ }
+ isDisposed = true;
+#if DEBUG_STREAMS
+ this.DebugDispose(typeof(ZlibBaseStream));
+#endif
+ await base.DisposeAsync().ConfigureAwait(false);
+ if (_stream is null)
+ {
+ return;
+ }
+ try
+ {
+ await finishAsync().ConfigureAwait(false);
+ }
+ finally
+ {
+ end();
+ if (_stream != null)
+ {
+ await _stream.DisposeAsync().ConfigureAwait(false);
+ _stream = null;
+ }
+ }
+ }
+#endif
+
public override void Flush()
{
_stream.Flush();
@@ -390,6 +573,14 @@ internal class ZlibBaseStream : Stream, IStreamStack
z.AvailableBytesIn = 0;
}
+ public override async Task FlushAsync(CancellationToken cancellationToken)
+ {
+ await _stream.FlushAsync(cancellationToken).ConfigureAwait(false);
+ //rewind the buffer
+ ((IStreamStack)this).Rewind(z.AvailableBytesIn); //unused
+ z.AvailableBytesIn = 0;
+ }
+
public override Int64 Seek(Int64 offset, SeekOrigin origin) =>
throw new NotSupportedException();
@@ -678,6 +869,42 @@ internal class ZlibBaseStream : Stream, IStreamStack
return rc;
}
+ public override Task ReadAsync(
+ byte[] buffer,
+ int offset,
+ int count,
+ CancellationToken cancellationToken
+ )
+ {
+ // For now, delegate to synchronous Read wrapped in Task
+ // A full async implementation would require refactoring the complex zlib codec logic
+ return Task.Run(() => Read(buffer, offset, count), cancellationToken);
+ }
+
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override ValueTask ReadAsync(
+ Memory buffer,
+ CancellationToken cancellationToken = default
+ )
+ {
+ // For now, delegate to synchronous Read wrapped in ValueTask
+ return new ValueTask(Task.Run(() =>
+ {
+ byte[] array = System.Buffers.ArrayPool.Shared.Rent(buffer.Length);
+ try
+ {
+ int read = Read(array, 0, buffer.Length);
+ array.AsSpan(0, read).CopyTo(buffer.Span);
+ return read;
+ }
+ finally
+ {
+ System.Buffers.ArrayPool.Shared.Return(array);
+ }
+ }, cancellationToken));
+ }
+#endif
+
public override Boolean CanRead => _stream.CanRead;
public override Boolean CanSeek => _stream.CanSeek;
diff --git a/src/SharpCompress/Compressors/Deflate/ZlibStream.cs b/src/SharpCompress/Compressors/Deflate/ZlibStream.cs
index d2051649..8ce52d48 100644
--- a/src/SharpCompress/Compressors/Deflate/ZlibStream.cs
+++ b/src/SharpCompress/Compressors/Deflate/ZlibStream.cs
@@ -28,6 +28,8 @@
using System;
using System.IO;
using System.Text;
+using System.Threading;
+using System.Threading.Tasks;
using SharpCompress.IO;
namespace SharpCompress.Compressors.Deflate;
@@ -266,6 +268,34 @@ public class ZlibStream : Stream, IStreamStack
_baseStream.Flush();
}
+ public override async Task FlushAsync(CancellationToken cancellationToken)
+ {
+ if (_disposed)
+ {
+ throw new ObjectDisposedException("ZlibStream");
+ }
+ await _baseStream.FlushAsync(cancellationToken).ConfigureAwait(false);
+ }
+
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override async ValueTask DisposeAsync()
+ {
+ if (_disposed)
+ {
+ return;
+ }
+ _disposed = true;
+ if (_baseStream != null)
+ {
+ await _baseStream.DisposeAsync().ConfigureAwait(false);
+ }
+#if DEBUG_STREAMS
+ this.DebugDispose(typeof(ZlibStream));
+#endif
+ await base.DisposeAsync().ConfigureAwait(false);
+ }
+#endif
+
///
/// Read data from the stream.
///
@@ -301,6 +331,34 @@ public class ZlibStream : Stream, IStreamStack
return _baseStream.Read(buffer, offset, count);
}
+ public override async Task ReadAsync(
+ byte[] buffer,
+ int offset,
+ int count,
+ CancellationToken cancellationToken
+ )
+ {
+ if (_disposed)
+ {
+ throw new ObjectDisposedException("ZlibStream");
+ }
+ return await _baseStream.ReadAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
+ }
+
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override async ValueTask ReadAsync(
+ Memory buffer,
+ CancellationToken cancellationToken = default
+ )
+ {
+ if (_disposed)
+ {
+ throw new ObjectDisposedException("ZlibStream");
+ }
+ return await _baseStream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
+ }
+#endif
+
public override int ReadByte()
{
if (_disposed)
@@ -355,6 +413,34 @@ public class ZlibStream : Stream, IStreamStack
_baseStream.Write(buffer, offset, count);
}
+ public override async Task WriteAsync(
+ byte[] buffer,
+ int offset,
+ int count,
+ CancellationToken cancellationToken
+ )
+ {
+ if (_disposed)
+ {
+ throw new ObjectDisposedException("ZlibStream");
+ }
+ await _baseStream.WriteAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
+ }
+
+#if !NETFRAMEWORK && !NETSTANDARD2_0
+ public override async ValueTask WriteAsync(
+ ReadOnlyMemory buffer,
+ CancellationToken cancellationToken = default
+ )
+ {
+ if (_disposed)
+ {
+ throw new ObjectDisposedException("ZlibStream");
+ }
+ await _baseStream.WriteAsync(buffer, cancellationToken).ConfigureAwait(false);
+ }
+#endif
+
public override void WriteByte(byte value)
{
if (_disposed)
diff --git a/src/SharpCompress/Utility.cs b/src/SharpCompress/Utility.cs
index c6cb3e95..fb5021d1 100644
--- a/src/SharpCompress/Utility.cs
+++ b/src/SharpCompress/Utility.cs
@@ -91,6 +91,28 @@ internal static class Utility
while (source.Read(buffer.Memory.Span) > 0) { }
}
+ public static async Task SkipAsync(this Stream source, CancellationToken cancellationToken = default)
+ {
+ var array = ArrayPool.Shared.Rent(TEMP_BUFFER_SIZE);
+ try
+ {
+ while (true)
+ {
+ var read = await source
+ .ReadAsync(array, 0, array.Length, cancellationToken)
+ .ConfigureAwait(false);
+ if (read <= 0)
+ {
+ break;
+ }
+ }
+ }
+ finally
+ {
+ ArrayPool.Shared.Return(array);
+ }
+ }
+
public static DateTime DosDateToDateTime(ushort iDate, ushort iTime)
{
var year = (iDate / 512) + 1980;