Files
sharpcompress/src/SharpCompress/IO/RewindableStream.Async.cs

368 lines
12 KiB
C#
Raw Normal View History

2026-01-28 16:50:35 +00:00
using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
namespace SharpCompress.IO;
internal partial class RewindableStream
{
public override async Task<int> ReadAsync(
byte[] buffer,
int offset,
int count,
CancellationToken cancellationToken
)
{
if (count == 0)
{
return 0;
}
2026-01-29 15:23:53 +00:00
2026-02-04 08:20:05 +00:00
// If recording is active or we're reading from the recording buffer, use legacy behavior
if (IsRecording || (isRewound && bufferStream.Position != bufferStream.Length))
{
return await ReadWithRecordingAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
}
// If rolling buffer is enabled (and not recording), use rolling buffer logic
if (_rollingBuffer is not null)
{
return await ReadWithRollingBufferAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
}
// No buffering - read directly from stream
int read = await stream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
streamPosition += read;
_logicalPosition = streamPosition;
return read;
}
/// <summary>
/// Async version of ReadWithRecording (legacy behavior for format detection).
/// </summary>
private async Task<int> ReadWithRecordingAsync(
byte[] buffer,
int offset,
int count,
CancellationToken cancellationToken
)
{
2026-01-28 16:50:35 +00:00
int read;
2026-01-29 15:23:53 +00:00
if (isRewound && bufferStream.Position != bufferStream.Length)
2026-01-28 16:50:35 +00:00
{
2026-01-29 15:23:53 +00:00
var readCount = Math.Min(count, (int)(bufferStream.Length - bufferStream.Position));
read = await bufferStream
.ReadAsync(buffer, offset, readCount, cancellationToken)
.ConfigureAwait(false);
if (read < count)
2026-01-28 16:50:35 +00:00
{
2026-01-29 15:23:53 +00:00
var tempRead = await stream
2026-01-28 16:50:35 +00:00
.ReadAsync(buffer, offset + read, count - read, cancellationToken)
.ConfigureAwait(false);
if (IsRecording)
{
2026-01-29 15:23:53 +00:00
await bufferStream
.WriteAsync(buffer, offset + read, tempRead, cancellationToken)
.ConfigureAwait(false);
2026-01-28 16:50:35 +00:00
}
2026-02-04 08:20:05 +00:00
else if (_rollingBuffer is not null && tempRead > 0)
{
// When transitioning out of recording mode, add to rolling buffer
// so that future rewinds will work
AddToRollingBuffer(buffer, offset + read, tempRead);
}
2026-01-29 15:23:53 +00:00
streamPosition += tempRead;
2026-02-04 08:20:05 +00:00
_logicalPosition = streamPosition;
2026-01-28 16:50:35 +00:00
read += tempRead;
}
2026-01-29 15:23:53 +00:00
if (bufferStream.Position == bufferStream.Length)
2026-01-28 16:50:35 +00:00
{
2026-01-29 15:23:53 +00:00
isRewound = false;
2026-01-28 16:50:35 +00:00
}
return read;
}
read = await stream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
if (IsRecording)
{
2026-01-29 15:23:53 +00:00
await bufferStream
.WriteAsync(buffer, offset, read, cancellationToken)
.ConfigureAwait(false);
2026-01-28 16:50:35 +00:00
}
2026-01-29 15:23:53 +00:00
streamPosition += read;
2026-02-04 08:20:05 +00:00
_logicalPosition = streamPosition;
2026-01-28 16:50:35 +00:00
return read;
}
2026-02-04 08:20:05 +00:00
/// <summary>
/// Async version of ReadWithRollingBuffer.
/// </summary>
private async Task<int> ReadWithRollingBufferAsync(
byte[] buffer,
int offset,
int count,
CancellationToken cancellationToken
)
{
int totalRead = 0;
// If logical position is behind stream position, read from rolling buffer first
while (count > 0 && _logicalPosition < streamPosition)
{
long bytesFromEnd = streamPosition - _logicalPosition;
if (bytesFromEnd > _rollingBufferLength)
{
throw new InvalidOperationException(
"Logical position is outside rolling buffer range."
);
}
int bufferIndex = (int)(
(_rollingBufferWritePos - bytesFromEnd + _rollingBufferSize) % _rollingBufferSize
);
int availableFromBuffer = (int)Math.Min(bytesFromEnd, count);
int firstPart = Math.Min(availableFromBuffer, _rollingBufferSize - bufferIndex);
Array.Copy(_rollingBuffer!, bufferIndex, buffer, offset, firstPart);
if (firstPart < availableFromBuffer)
{
Array.Copy(
_rollingBuffer!,
0,
buffer,
offset + firstPart,
availableFromBuffer - firstPart
);
}
totalRead += availableFromBuffer;
offset += availableFromBuffer;
count -= availableFromBuffer;
_logicalPosition += availableFromBuffer;
}
// If more data needed, read from underlying stream
if (count > 0)
{
int read = await stream
.ReadAsync(buffer, offset, count, cancellationToken)
.ConfigureAwait(false);
if (read > 0)
{
AddToRollingBuffer(buffer, offset, read);
streamPosition += read;
_logicalPosition += read;
totalRead += read;
}
}
return totalRead;
}
2026-01-28 16:50:35 +00:00
#if !LEGACY_DOTNET
public override async ValueTask<int> ReadAsync(
Memory<byte> buffer,
CancellationToken cancellationToken = default
)
{
if (buffer.Length == 0)
{
return 0;
}
2026-01-29 15:23:53 +00:00
2026-02-04 08:20:05 +00:00
// If recording is active or we're reading from the recording buffer, use legacy behavior
if (IsRecording || (isRewound && bufferStream.Position != bufferStream.Length))
{
return await ReadWithRecordingAsync(buffer, cancellationToken).ConfigureAwait(false);
}
// If rolling buffer is enabled (and not recording), use rolling buffer logic
if (_rollingBuffer is not null)
{
return await ReadWithRollingBufferAsync(buffer, cancellationToken)
.ConfigureAwait(false);
}
// No buffering - read directly from stream
int read = await stream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
streamPosition += read;
_logicalPosition = streamPosition;
return read;
}
/// <summary>
/// Async version of ReadWithRecording for Memory&lt;byte&gt; (legacy behavior for format detection).
/// </summary>
private async ValueTask<int> ReadWithRecordingAsync(
Memory<byte> buffer,
CancellationToken cancellationToken
)
{
2026-01-28 16:50:35 +00:00
int read;
2026-01-29 15:23:53 +00:00
if (isRewound && bufferStream.Position != bufferStream.Length)
2026-01-28 16:50:35 +00:00
{
2026-01-29 15:23:53 +00:00
var readCount = (int)
Math.Min(buffer.Length, bufferStream.Length - bufferStream.Position);
read = await bufferStream
.ReadAsync(buffer.Slice(0, readCount), cancellationToken)
.ConfigureAwait(false);
if (read < buffer.Length)
2026-01-28 16:50:35 +00:00
{
2026-01-29 15:23:53 +00:00
var tempRead = await stream
2026-01-28 16:50:35 +00:00
.ReadAsync(buffer.Slice(read), cancellationToken)
.ConfigureAwait(false);
if (IsRecording)
{
2026-01-29 15:23:53 +00:00
await bufferStream
.WriteAsync(buffer.Slice(read, tempRead), cancellationToken)
.ConfigureAwait(false);
2026-01-28 16:50:35 +00:00
}
2026-02-04 08:20:05 +00:00
else if (_rollingBuffer is not null && tempRead > 0)
{
// When transitioning out of recording mode, add to rolling buffer
// so that future rewinds will work
var tempBuffer = buffer.Slice(read, tempRead).ToArray();
AddToRollingBuffer(tempBuffer, 0, tempRead);
}
streamPosition += tempRead;
2026-02-04 08:20:05 +00:00
_logicalPosition = streamPosition;
2026-01-28 16:50:35 +00:00
read += tempRead;
}
2026-01-29 15:23:53 +00:00
if (bufferStream.Position == bufferStream.Length)
2026-01-28 16:50:35 +00:00
{
2026-01-29 15:23:53 +00:00
isRewound = false;
2026-01-28 16:50:35 +00:00
}
return read;
}
read = await stream.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
if (IsRecording)
{
2026-01-29 15:23:53 +00:00
await bufferStream
.WriteAsync(buffer.Slice(0, read), cancellationToken)
.ConfigureAwait(false);
2026-01-28 16:50:35 +00:00
}
streamPosition += read;
2026-02-04 08:20:05 +00:00
_logicalPosition = streamPosition;
2026-01-28 16:50:35 +00:00
return read;
}
2026-02-04 08:20:05 +00:00
/// <summary>
/// Async version of ReadWithRollingBuffer for Memory&lt;byte&gt;.
/// </summary>
private async ValueTask<int> ReadWithRollingBufferAsync(
Memory<byte> buffer,
CancellationToken cancellationToken
)
{
int totalRead = 0;
int count = buffer.Length;
int offset = 0;
// If logical position is behind stream position, read from rolling buffer first
while (count > 0 && _logicalPosition < streamPosition)
{
long bytesFromEnd = streamPosition - _logicalPosition;
if (bytesFromEnd > _rollingBufferLength)
{
throw new InvalidOperationException(
"Logical position is outside rolling buffer range."
);
}
int bufferIndex = (int)(
(_rollingBufferWritePos - bytesFromEnd + _rollingBufferSize) % _rollingBufferSize
);
int availableFromBuffer = (int)Math.Min(bytesFromEnd, count);
int firstPart = Math.Min(availableFromBuffer, _rollingBufferSize - bufferIndex);
_rollingBuffer.AsSpan(bufferIndex, firstPart).CopyTo(buffer.Span.Slice(offset));
if (firstPart < availableFromBuffer)
{
_rollingBuffer
.AsSpan(0, availableFromBuffer - firstPart)
.CopyTo(buffer.Span.Slice(offset + firstPart));
}
totalRead += availableFromBuffer;
offset += availableFromBuffer;
count -= availableFromBuffer;
_logicalPosition += availableFromBuffer;
}
// If more data needed, read from underlying stream
if (count > 0)
{
int read = await stream
.ReadAsync(buffer.Slice(offset, count), cancellationToken)
.ConfigureAwait(false);
if (read > 0)
{
// AddToRollingBuffer expects byte[], so we need to copy
var tempBuffer = buffer.Slice(offset, read).ToArray();
AddToRollingBuffer(tempBuffer, 0, read);
streamPosition += read;
_logicalPosition += read;
totalRead += read;
}
}
return totalRead;
}
2026-01-28 16:50:35 +00:00
#endif
2026-01-29 15:23:53 +00:00
public override Task WriteAsync(
byte[] buffer,
int offset,
int count,
CancellationToken cancellationToken
) => throw new NotSupportedException();
#if !LEGACY_DOTNET
public override ValueTask WriteAsync(
ReadOnlyMemory<byte> buffer,
CancellationToken cancellationToken = default
) => throw new NotSupportedException();
#endif
public override Task FlushAsync(CancellationToken cancellationToken) =>
throw new NotSupportedException();
public override async Task CopyToAsync(
Stream destination,
int bufferSize,
CancellationToken cancellationToken
)
{
byte[] buffer = new byte[bufferSize];
int bytesRead;
while ((bytesRead = await ReadAsync(buffer, 0, buffer.Length, cancellationToken)) != 0)
{
await destination.WriteAsync(buffer, 0, bytesRead, cancellationToken);
}
}
#if !LEGACY_DOTNET
public override async ValueTask DisposeAsync()
{
if (!isDisposed)
{
isDisposed = true;
await stream.DisposeAsync();
2026-02-04 08:20:05 +00:00
if (_rollingBuffer is not null)
{
System.Buffers.ArrayPool<byte>.Shared.Return(_rollingBuffer);
_rollingBuffer = null;
}
2026-01-29 15:23:53 +00:00
}
}
#endif
2026-01-28 16:50:35 +00:00
}