From 165089987e6630e6ad31ab192442e027c460689c Mon Sep 17 00:00:00 2001 From: Louis Seubert Date: Tue, 25 Aug 2026 22:15:40 +0200 Subject: [PATCH] feat: add system.io.pipeline support to pipe io --- src/process.tests/PipingTests.cs | 49 ++++++++++++++++++++++++++++++++ src/process/PipeSource.cs | 9 ++++++ src/process/PipeTarget.cs | 29 +++++++++++++++++++ 3 files changed, 87 insertions(+) diff --git a/src/process.tests/PipingTests.cs b/src/process.tests/PipingTests.cs index 82f23b3..1839cc1 100644 --- a/src/process.tests/PipingTests.cs +++ b/src/process.tests/PipingTests.cs @@ -2,6 +2,7 @@ // SPDX-License-Identifier: EUPL-1.2 using System.Text; +using System.IO.Pipelines; using Geekeey.Process.Buffered; @@ -64,6 +65,23 @@ internal sealed class PipingTests 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] 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] public async Task I_can_execute_a_command_and_pipe_the_stdout_into_multiple_hierarchical_targets() { diff --git a/src/process/PipeSource.cs b/src/process/PipeSource.cs index 4f00157..096e629 100644 --- a/src/process/PipeSource.cs +++ b/src/process/PipeSource.cs @@ -2,6 +2,7 @@ // SPDX-License-Identifier: EUPL-1.2 using System.Text; +using System.IO.Pipelines; namespace Geekeey.Process; @@ -75,6 +76,14 @@ public abstract partial class PipeSource return Create((target, token) => stream.CopyToAsync(target, bufferSize, token)); } + /// + /// Creates a pipe source that reads from the specified pipe reader. + /// + public static PipeSource FromPipeReader(PipeReader reader) + { + return Create(reader.CopyToAsync); + } + /// /// Creates a pipe source that reads from the specified file. /// diff --git a/src/process/PipeTarget.cs b/src/process/PipeTarget.cs index fbc1c97..ff82599 100644 --- a/src/process/PipeTarget.cs +++ b/src/process/PipeTarget.cs @@ -2,6 +2,7 @@ // SPDX-License-Identifier: EUPL-1.2 using System.Buffers; +using System.IO.Pipelines; using System.Text; namespace Geekeey.Process; @@ -166,6 +167,34 @@ public partial class PipeTarget await origin.CopyToAsync(stream, bufferSize, cancellationToken)); } + /// + /// Creates a pipe target that writes to the specified pipe writer. + /// + 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); + } + }); + } + /// /// Creates a pipe target that writes to the specified file. ///