Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Directory.Build.props
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
<Project>
<PropertyGroup>
<TargetFramework>net9.0</TargetFramework>
<VersionPrefix>0.0.4</VersionPrefix>
<VersionPrefix>0.0.5</VersionPrefix>
<VersionSuffix>rc1</VersionSuffix>
<Authors>Paul Bleess</Authors>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
Expand Down
13 changes: 13 additions & 0 deletions src/PipeExtensions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
using System.IO.Pipelines;

namespace Yamux
{
public static class PipeExtensions
{
/// <summary>
/// Creates a yamux session from the duplex pipe
/// </summary>
public static Session AsYamuxSession(this IDuplexPipe pipe, bool isClient, bool keepOpen = false, SessionOptions? options = null)
=> new Session(new PipePeer(pipe), isClient, keepOpen, options);
}
}
66 changes: 66 additions & 0 deletions src/PipePeer.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
using System.IO.Pipelines;

namespace Yamux
{
public class PipePeer : ITransport
{
private readonly PipeReader _reader;
private readonly PipeWriter _writer;

public PipePeer(IDuplexPipe pipe)
{
ArgumentNullException.ThrowIfNull(pipe);
_reader = pipe.Input;
_writer = pipe.Output;
}

public async ValueTask<int> ReadAsync(Memory<byte> data, CancellationToken cancel)
{
if (data.IsEmpty)
return 0;

var result = await _reader.ReadAsync(cancel);
var buffer = result.Buffer;

if (buffer.IsEmpty && result.IsCompleted)
{
_reader.AdvanceTo(buffer.End);
return 0;
}

var len = (int)Math.Min(buffer.Length, data.Length);
var remaining = len;
foreach (var segment in buffer)
{
if (remaining <= 0)
break;
var toCopy = Math.Min(segment.Length, remaining);
segment.Span[..toCopy].CopyTo(data.Span[(len - remaining)..]);
remaining -= toCopy;
}
_reader.AdvanceTo(buffer.GetPosition(len));

return len;
}

public async ValueTask WriteAsync(ReadOnlyMemory<byte> data, CancellationToken cancel)
{
if (data.IsEmpty)
return;

await _writer.WriteAsync(data, cancel);
}

public void Close()
{
_reader.Complete();
_writer.Complete();
}

public void Dispose()
{
_reader.Complete();
_writer.Complete();
}
}
}
98 changes: 98 additions & 0 deletions test/Yamux.Tests/PipeTransportTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
using AwesomeAssertions;
using Bogus;
using System.Buffers;
using System.IO.Pipelines;
using System.Text;

namespace Yamux.Tests;

public class PipeTransportTests
{
private static (IDuplexPipe Client, IDuplexPipe Server) CreatePipePair()
{
var clientPipe = new Pipe();
var serverPipe = new Pipe();

var client = new DuplexPipe(serverPipe.Reader, clientPipe.Writer);
var server = new DuplexPipe(clientPipe.Reader, serverPipe.Writer);

return (client, server);
}

[Fact]
public async Task SingleOneWayTest()
{
var faker = new Faker();
var data = faker.Random.Chars(count: 1024 * 750);
var buffer = Encoding.UTF8.GetBytes(data).AsMemory();
var result = new byte[buffer.Length].AsMemory();

(var client, var server) = CreatePipePair();

var serverTask = Task.Run(async () =>
{
await using var serverSession = server.AsYamuxSession(false);
serverSession.Start();

using var channel = await serverSession.AcceptAsync();

long index = 0;
ReadResult res;
do
{
res = await channel.Input.ReadAsync();
if (res.Buffer.Length > 0)
{
res.Buffer.CopyTo(result.Slice((int)index, (int)res.Buffer.Length).Span);
index += res.Buffer.Length;
channel.Input.AdvanceTo(res.Buffer.End, res.Buffer.End);
}
} while (!res.IsCanceled && !res.IsCompleted);

channel.Close();
await channel.WhenRemoteCloseAsync(TimeSpan.FromSeconds(1));
channel.Dispose();
});

var clientTask = Task.Run(async () =>
{
await using var clientSession = client.AsYamuxSession(true);
clientSession.Start();

using var channel = await clientSession.OpenChannelAsync();

var size = buffer.Length;
int current = 0;
int chunkSize = 1024 * 4;

while (current < size)
{
if (buffer.Length > chunkSize)
{
var end = current + chunkSize;
var slice = buffer.Slice(current, end >= buffer.Length ? buffer.Length - current : chunkSize);
await channel.WriteAsync(slice, CancellationToken.None);
}
current += chunkSize;
}

channel.Close();
await channel.WhenRemoteCloseAsync(TimeSpan.FromSeconds(1));
});

await Task.WhenAll(serverTask, clientTask);
result.ToArray().Should().BeEquivalentTo(buffer.ToArray());
}

private sealed class DuplexPipe : IDuplexPipe
{
public DuplexPipe(PipeReader input, PipeWriter output)
{
Input = input;
Output = output;
}

public PipeReader Input { get; }
public PipeWriter Output { get; }
}
}
Loading