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.
///