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;