Editor/Core/RelayPipe.cs
using System;
using System.IO;
using System.Threading;
using System.Threading.Channels;
using System.Threading.Tasks;

namespace TeamCreate;

/// <summary>
/// A byte stream that rides on discrete messages: what arrives is fed in with <see cref="Feed"/>, what
/// is written leaves through a send delegate in chunks no larger than a relay message. The session's
/// framing, encryption and sequencing sit above a <see cref="Stream"/> and never learn that the bytes
/// travel through a relay instead of straight over TCP.
/// </summary>
public sealed class RelayPipe : Stream
{
    // A reader that stops draining must not let the buffer grow without bound: past this the pipe closes,
    // which the session sees as a lost connection and recovers from, instead of the editor running out of memory.
    public const int MaxBufferedBytes = 32 * 1024 * 1024;

    private readonly Channel<byte[]> incoming = Channel.CreateUnbounded<byte[]>(new UnboundedChannelOptions { SingleReader = true, SingleWriter = false });

    private readonly Func<ReadOnlyMemory<byte>, CancellationToken, ValueTask> send;

    private readonly Action closed;

    private byte[] current = Array.Empty<byte>();

    private int offset;

    private long buffered;

    private int disposed;

    public RelayPipe(Func<ReadOnlyMemory<byte>, CancellationToken, ValueTask> send, Action closed = null)
    {
        this.send = send ?? throw new ArgumentNullException(nameof(send));
        this.closed = closed;
    }

    public bool IsClosed => Volatile.Read(ref disposed) != 0;

    public long BufferedBytes => Interlocked.Read(ref buffered);

    /// <summary>Queues received bytes. Returns false, and closes the pipe, when it is closed or over its buffer budget.</summary>
    public bool Feed(byte[] data)
    {
        if (data == null || data.Length == 0)
        {
            return !IsClosed;
        }
        if (IsClosed)
        {
            return false;
        }
        if (Interlocked.Add(ref buffered, data.Length) > MaxBufferedBytes)
        {
            Dispose();
            return false;
        }
        return incoming.Writer.TryWrite(data);
    }

    /// <summary>Marks the far end as gone without calling the close callback: the peer already knows.</summary>
    public void CompleteFromPeer()
    {
        if (Interlocked.Exchange(ref disposed, 1) == 0)
        {
            incoming.Writer.TryComplete();
        }
    }

    public override async ValueTask<int> ReadAsync(Memory<byte> destination, CancellationToken cancellationToken = default)
    {
        if (destination.Length == 0)
        {
            return 0;
        }
        while (true)
        {
            if (offset < current.Length)
            {
                int count = Math.Min(destination.Length, current.Length - offset);
                current.AsMemory(offset, count).CopyTo(destination);
                offset += count;
                Interlocked.Add(ref buffered, -count);
                return count;
            }
            try
            {
                current = await incoming.Reader.ReadAsync(cancellationToken).ConfigureAwait(false);
                offset = 0;
            }
            catch (ChannelClosedException)
            {
                return 0;
            }
        }
    }

    public override async ValueTask WriteAsync(ReadOnlyMemory<byte> source, CancellationToken cancellationToken = default)
    {
        while (source.Length > 0)
        {
            if (IsClosed)
            {
                throw new IOException("The relay connection is closed.");
            }
            int count = Math.Min(source.Length, RelayProtocol.MaxChunkBytes);
            await send(source.Slice(0, count), cancellationToken).ConfigureAwait(false);
            source = source.Slice(count);
        }
    }

    protected override void Dispose(bool disposing)
    {
        if (disposing && Interlocked.Exchange(ref disposed, 1) == 0)
        {
            incoming.Writer.TryComplete();
            try
            {
                closed?.Invoke();
            }
            catch
            {
                // A close notification that cannot be delivered changes nothing: the pipe is already closed.
            }
        }
        base.Dispose(disposing);
    }

    public override bool CanRead => true;

    public override bool CanWrite => true;

    public override bool CanSeek => false;

    public override long Length => throw new NotSupportedException();

    public override long Position
    {
        get => throw new NotSupportedException();
        set => throw new NotSupportedException();
    }

    public override void Flush()
    {
    }

    public override Task FlushAsync(CancellationToken cancellationToken) => Task.CompletedTask;

    public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException("The relay pipe is asynchronous only; a synchronous read would block the editor's frame.");

    public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException("The relay pipe is asynchronous only.");

    public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException();

    public override void SetLength(long value) => throw new NotSupportedException();
}