feat: add system.io.pipeline support to pipe io

This commit is contained in:
Louis Seubert 2026-08-25 22:15:40 +02:00
commit 165089987e
Signed by: louis9902
GPG key ID: 4B9DB28F826553BD
3 changed files with 87 additions and 0 deletions

View file

@ -2,6 +2,7 @@
// SPDX-License-Identifier: EUPL-1.2 // SPDX-License-Identifier: EUPL-1.2
using System.Text; using System.Text;
using System.IO.Pipelines;
using Geekeey.Process.Buffered; using Geekeey.Process.Buffered;
@ -64,6 +65,23 @@ internal sealed class PipingTests
await Assert.That(result.StandardOutput.Trim()).IsEqualTo("Hello World!"); await Assert.That(result.StandardOutput.Trim()).IsEqualTo("Hello World!");
} }
[Test]
public async Task I_can_execute_a_command_and_pipe_the_stdin_from_a_pipe_reader()
{
var pipe = new Pipe();
await pipe.Writer.WriteAsync("Hello World!"u8.ToArray());
await pipe.Writer.CompleteAsync();
var cmd = PipeSource.FromPipeReader(pipe.Reader) |
new Command(Testing.Fixture.Program.FilePath)
.WithArguments("echo-stdin");
var result = await cmd.ExecuteBufferedAsync();
await Assert.That(result.StandardOutput.Trim()).IsEqualTo("Hello World!");
await pipe.Reader.CompleteAsync();
}
[Test] [Test]
public async Task I_can_execute_a_command_and_pipe_the_stdin_from_a_stream_with_a_custom_buffer_size() public async Task I_can_execute_a_command_and_pipe_the_stdin_from_a_stream_with_a_custom_buffer_size()
{ {
@ -570,6 +588,37 @@ internal sealed class PipingTests
} }
} }
[Test]
public async Task I_can_execute_a_command_and_pipe_the_stdout_to_a_pipe_writer()
{
// Arrange
var pipe = new Pipe();
var cmd = new Command(Testing.Fixture.Program.FilePath)
.WithArguments(["generate", "blob", "--length", "100000"]) |
PipeTarget.ToPipeWriter(pipe.Writer);
// Act
using var output = new MemoryStream();
var executionTask = cmd.ExecuteAsync();
while (true)
{
var result = await pipe.Reader.ReadAsync();
foreach (var segment in result.Buffer)
{
await output.WriteAsync(segment);
}
pipe.Reader.AdvanceTo(result.Buffer.End);
if (result.IsCompleted)
{
break;
}
}
await executionTask;
await pipe.Reader.CompleteAsync();
await Assert.That(output.Length).IsEqualTo(100_000);
}
[Test] [Test]
public async Task I_can_execute_a_command_and_pipe_the_stdout_into_multiple_hierarchical_targets() public async Task I_can_execute_a_command_and_pipe_the_stdout_into_multiple_hierarchical_targets()
{ {

View file

@ -2,6 +2,7 @@
// SPDX-License-Identifier: EUPL-1.2 // SPDX-License-Identifier: EUPL-1.2
using System.Text; using System.Text;
using System.IO.Pipelines;
namespace Geekeey.Process; namespace Geekeey.Process;
@ -75,6 +76,14 @@ public abstract partial class PipeSource
return Create((target, token) => stream.CopyToAsync(target, bufferSize, token)); return Create((target, token) => stream.CopyToAsync(target, bufferSize, token));
} }
/// <summary>
/// Creates a pipe source that reads from the specified pipe reader.
/// </summary>
public static PipeSource FromPipeReader(PipeReader reader)
{
return Create(reader.CopyToAsync);
}
/// <summary> /// <summary>
/// Creates a pipe source that reads from the specified file. /// Creates a pipe source that reads from the specified file.
/// </summary> /// </summary>

View file

@ -2,6 +2,7 @@
// SPDX-License-Identifier: EUPL-1.2 // SPDX-License-Identifier: EUPL-1.2
using System.Buffers; using System.Buffers;
using System.IO.Pipelines;
using System.Text; using System.Text;
namespace Geekeey.Process; namespace Geekeey.Process;
@ -166,6 +167,34 @@ public partial class PipeTarget
await origin.CopyToAsync(stream, bufferSize, cancellationToken)); await origin.CopyToAsync(stream, bufferSize, cancellationToken));
} }
/// <summary>
/// Creates a pipe target that writes to the specified pipe writer.
/// </summary>
public static PipeTarget ToPipeWriter(PipeWriter writer, bool leaveOpen = false)
{
return Create(async (origin, cancellationToken) =>
{
var exception = default(Exception);
try
{
await origin.CopyToAsync(writer, cancellationToken);
if (!leaveOpen)
{
await writer.CompleteAsync();
}
}
catch (Exception caught)
{
exception = caught;
throw;
}
finally
{
await writer.CompleteAsync(exception);
}
});
}
/// <summary> /// <summary>
/// Creates a pipe target that writes to the specified file. /// Creates a pipe target that writes to the specified file.
/// </summary> /// </summary>