This commit is contained in:
Adam Hathcock
2025-10-27 09:37:00 +00:00
parent b3975b7bbd
commit f8cc4ade8a
11 changed files with 108 additions and 47 deletions

View File

@@ -14,8 +14,9 @@ public class TarArchiveEntry : TarEntry, IArchiveEntry
public virtual Stream OpenEntryStream() => Parts.Single().GetCompressedStream().NotNull();
public virtual Task<Stream> OpenEntryStreamAsync(CancellationToken cancellationToken = default) =>
Task.FromResult(OpenEntryStream());
public virtual Task<Stream> OpenEntryStreamAsync(
CancellationToken cancellationToken = default
) => Task.FromResult(OpenEntryStream());
#region IArchiveEntry Members

View File

@@ -13,8 +13,9 @@ public class ZipArchiveEntry : ZipEntry, IArchiveEntry
public virtual Stream OpenEntryStream() => Parts.Single().GetCompressedStream().NotNull();
public virtual Task<Stream> OpenEntryStreamAsync(CancellationToken cancellationToken = default) =>
Task.FromResult(OpenEntryStream());
public virtual Task<Stream> OpenEntryStreamAsync(
CancellationToken cancellationToken = default
) => Task.FromResult(OpenEntryStream());
#region IArchiveEntry Members

View File

@@ -163,7 +163,9 @@ public class EntryStream : Stream, IStreamStack
CancellationToken cancellationToken
)
{
var read = await _stream.ReadAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
var read = await _stream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
if (read <= 0)
{
_completed = true;

View File

@@ -365,7 +365,9 @@ public class DeflateStream : Stream, IStreamStack
{
throw new ObjectDisposedException("DeflateStream");
}
return await _baseStream.ReadAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
return await _baseStream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
}
#if !NETFRAMEWORK && !NETSTANDARD2_0
@@ -454,7 +456,9 @@ public class DeflateStream : Stream, IStreamStack
{
throw new ObjectDisposedException("DeflateStream");
}
await _baseStream.WriteAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
await _baseStream
.WriteAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
}
#if !NETFRAMEWORK && !NETSTANDARD2_0

View File

@@ -331,7 +331,9 @@ public class GZipStream : Stream, IStreamStack
{
throw new ObjectDisposedException("GZipStream");
}
var n = await BaseStream.ReadAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
var n = await BaseStream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
if (!_firstReadDone)
{

View File

@@ -199,7 +199,12 @@ internal class ZlibBaseStream : Stream, IStreamStack
} while (!done);
}
public override async Task WriteAsync(byte[] buffer, int offset, int count, CancellationToken cancellationToken)
public override async Task WriteAsync(
byte[] buffer,
int offset,
int count,
CancellationToken cancellationToken
)
{
// workitem 7159
// calculate the CRC on the unccompressed data (before writing)
@@ -238,7 +243,14 @@ internal class ZlibBaseStream : Stream, IStreamStack
throw new ZlibException((_wantCompress ? "de" : "in") + "flating: " + _z.Message);
}
await _stream.WriteAsync(_workingBuffer, 0, _workingBuffer.Length - _z.AvailableBytesOut, cancellationToken).ConfigureAwait(false);
await _stream
.WriteAsync(
_workingBuffer,
0,
_workingBuffer.Length - _z.AvailableBytesOut,
cancellationToken
)
.ConfigureAwait(false);
done = _z.AvailableBytesIn == 0 && _z.AvailableBytesOut != 0;
@@ -418,7 +430,14 @@ internal class ZlibBaseStream : Stream, IStreamStack
if (_workingBuffer.Length - _z.AvailableBytesOut > 0)
{
await _stream.WriteAsync(_workingBuffer, 0, _workingBuffer.Length - _z.AvailableBytesOut, cancellationToken).ConfigureAwait(false);
await _stream
.WriteAsync(
_workingBuffer,
0,
_workingBuffer.Length - _z.AvailableBytesOut,
cancellationToken
)
.ConfigureAwait(false);
}
done = _z.AvailableBytesIn == 0 && _z.AvailableBytesOut != 0;
@@ -473,9 +492,9 @@ internal class ZlibBaseStream : Stream, IStreamStack
// 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);
var bytesRead = await _stream
.ReadAsync(trailer, _z.AvailableBytesIn, bytesNeeded, cancellationToken)
.ConfigureAwait(false);
}
}
else
@@ -888,20 +907,25 @@ internal class ZlibBaseStream : Stream, IStreamStack
)
{
// For now, delegate to synchronous Read wrapped in ValueTask
return new ValueTask<int>(Task.Run(() =>
{
byte[] array = System.Buffers.ArrayPool<byte>.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<byte>.Shared.Return(array);
}
}, cancellationToken));
return new ValueTask<int>(
Task.Run(
() =>
{
byte[] array = System.Buffers.ArrayPool<byte>.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<byte>.Shared.Return(array);
}
},
cancellationToken
)
);
}
#endif

View File

@@ -342,7 +342,9 @@ public class ZlibStream : Stream, IStreamStack
{
throw new ObjectDisposedException("ZlibStream");
}
return await _baseStream.ReadAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
return await _baseStream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
}
#if !NETFRAMEWORK && !NETSTANDARD2_0
@@ -424,7 +426,9 @@ public class ZlibStream : Stream, IStreamStack
{
throw new ObjectDisposedException("ZlibStream");
}
await _baseStream.WriteAsync(buffer, offset, count, cancellationToken).ConfigureAwait(false);
await _baseStream
.WriteAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
}
#if !NETFRAMEWORK && !NETSTANDARD2_0

View File

@@ -337,7 +337,6 @@ public class SourceStream : Stream, IStreamStack
return total - count;
}
#endif
public override void Close()

View File

@@ -211,8 +211,7 @@ public abstract class AbstractReader<TEntry, TVolume> : IReader, IReaderExtracti
{
var streamListener = this as IReaderExtractionListener;
using Stream s = OpenEntryStream();
await s
.TransferToAsync(writeStream, Entry, streamListener, cancellationToken)
await s.TransferToAsync(writeStream, Entry, streamListener, cancellationToken)
.ConfigureAwait(false);
}

View File

@@ -91,7 +91,10 @@ internal static class Utility
while (source.Read(buffer.Memory.Span) > 0) { }
}
public static async Task SkipAsync(this Stream source, CancellationToken cancellationToken = default)
public static async Task SkipAsync(
this Stream source,
CancellationToken cancellationToken = default
)
{
var array = ArrayPool<byte>.Shared.Rent(TEMP_BUFFER_SIZE);
try
@@ -261,7 +264,7 @@ internal static class Utility
while (
await ReadTransferBlockAsync(source, array, maxReadSize, cancellationToken)
.ConfigureAwait(false)
is var (success, count)
is var (success, count)
&& success
)
{

View File

@@ -28,7 +28,11 @@ public class AsyncTests : TestBase
);
// Just verify some files were extracted
var extractedFiles = Directory.GetFiles(SCRATCH_FILES_PATH, "*", SearchOption.AllDirectories);
var extractedFiles = Directory.GetFiles(
SCRATCH_FILES_PATH,
"*",
SearchOption.AllDirectories
);
Assert.True(extractedFiles.Length > 0, "No files were extracted");
}
@@ -107,7 +111,11 @@ public class AsyncTests : TestBase
);
// Just verify some files were extracted
var extractedFiles = Directory.GetFiles(SCRATCH_FILES_PATH, "*", SearchOption.AllDirectories);
var extractedFiles = Directory.GetFiles(
SCRATCH_FILES_PATH,
"*",
SearchOption.AllDirectories
);
Assert.True(extractedFiles.Length > 0, "No files were extracted");
}
@@ -146,7 +154,7 @@ public class AsyncTests : TestBase
var buffer = new byte[4096];
var totalRead = 0;
int bytesRead;
// Test ReadAsync on EntryStream
while ((bytesRead = await entryStream.ReadAsync(buffer, 0, buffer.Length)) > 0)
{
@@ -166,12 +174,15 @@ public class AsyncTests : TestBase
new Random(42).NextBytes(testData);
var compressedPath = Path.Combine(SCRATCH_FILES_PATH, "async_compressed.gz");
// Test async write with GZipStream
using (var fileStream = File.Create(compressedPath))
using (var gzipStream = new Compressors.Deflate.GZipStream(
fileStream,
Compressors.CompressionMode.Compress))
using (
var gzipStream = new Compressors.Deflate.GZipStream(
fileStream,
Compressors.CompressionMode.Compress
)
)
{
await gzipStream.WriteAsync(testData, 0, testData.Length);
await gzipStream.FlushAsync();
@@ -182,15 +193,26 @@ public class AsyncTests : TestBase
// Test async read with GZipStream
using (var fileStream = File.OpenRead(compressedPath))
using (var gzipStream = new Compressors.Deflate.GZipStream(
fileStream,
Compressors.CompressionMode.Decompress))
using (
var gzipStream = new Compressors.Deflate.GZipStream(
fileStream,
Compressors.CompressionMode.Decompress
)
)
{
var decompressed = new byte[testData.Length];
var totalRead = 0;
int bytesRead;
while (totalRead < decompressed.Length &&
(bytesRead = await gzipStream.ReadAsync(decompressed, totalRead, decompressed.Length - totalRead)) > 0)
while (
totalRead < decompressed.Length
&& (
bytesRead = await gzipStream.ReadAsync(
decompressed,
totalRead,
decompressed.Length - totalRead
)
) > 0
)
{
totalRead += bytesRead;
}