Updates Pipe & Makes NetState more testable (#288)

Updates pipe eliminate result and segments from being allocated.

Bumps release version
This commit is contained in:
Kamron Batman 2020-10-25 17:40:08 -07:00 • committed by GitHub
parent 293e691539
commit 990e6fe188
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
14 changed files with 219 additions and 198 deletions

View file

@ -33,9 +33,9 @@ namespace Server.Network
{
public delegate void NetStateCreatedCallback(NetState ns);
public delegate void EncodePacket(CircularBuffer<byte> buffer, ref int length);
public delegate void EncodePacket(ref CircularBuffer<byte> buffer, ref int length);
public partial class NetState : IComparable<NetState>
public partial class NetState : IComparable<NetState>, IDisposable
{
private static int RecvPipeSize = 1024 * 64;
private static int SendPipeSize = 1024 * 256;
@ -52,9 +52,7 @@ namespace Server.Network
private int m_Disposing;
private ClientVersion m_Version;
private byte[] _recvBuffer;
private Pipe<byte> _recvPipe;
private byte[] _sendBuffer;
private Pipe<byte> _sendPipe;
private long m_NextCheckActivity;
private volatile bool m_Running;
private readonly Thread _sendThread;
@ -89,9 +87,9 @@ namespace Server.Network
Menus = new List<IMenu>();
Trades = new List<SecureTrade>();
_recvBuffer = new byte[RecvPipeSize];
_recvPipe = new Pipe<byte>(_recvBuffer);
RecvPipe = new Pipe<byte>(_recvBuffer);
_sendBuffer = new byte[SendPipeSize];
_sendPipe = new Pipe<byte>(_sendBuffer);
SendPipe = new Pipe<byte>(_sendBuffer);
m_NextCheckActivity = Core.TickCount + 30000;
_sendThread = sendThread ?? Core.Thread;
@ -142,6 +140,10 @@ namespace Server.Network
public bool Seeded { get; set; }
public Pipe<byte> RecvPipe { get; private set; }
public Pipe<byte> SendPipe { get; private set; }
public Socket Connection { get; private set; }
public bool CompressionEnabled { get; set; }
@ -370,29 +372,27 @@ namespace Server.Network
NetworkState.Resume(ref m_NetworkState);
}
public Pipe<byte>.Result<byte> GetAvailableSendPipe() => _recvPipe.Writer.GetAvailable();
public bool GetAvailableSendPipe(out CircularBuffer<byte> buffer) => SendPipe.Writer.GetAvailable(out buffer);
public virtual void Send(CircularBuffer<byte> buffer, int length)
public virtual void Send(ref CircularBuffer<byte> buffer, int length)
{
if (Connection == null || BlockAllPackets || buffer.Length == 0)
{
return;
}
#if DEBUG
var currentThread = Thread.CurrentThread;
if (currentThread != _sendThread)
{
Console.Error.WriteLine("Core: Attempted to send packet outside core thread! [{0}]", currentThread.ManagedThreadId);
#if DEBUG
throw new InvalidThreadException(nameof(Send));
#endif
throw new InvalidThreadException("Attempted to send packet outside send thread!");
}
#endif
try
{
_packetEncoder?.Invoke(buffer, ref length);
_sendPipe.Writer.Advance((uint)length);
_packetEncoder?.Invoke(ref buffer, ref length);
SendPipe.Writer.Advance((uint)length);
}
catch (Exception ex)
{
@ -412,17 +412,16 @@ namespace Server.Network
return;
}
#if DEBUG
var currentThread = Thread.CurrentThread;
if (currentThread != _sendThread)
{
Console.Error.WriteLine("Core: Attempted to send packet outside send thread! [{0}]", currentThread.ManagedThreadId);
#if DEBUG
throw new InvalidThreadException(nameof(Send));
#endif
throw new InvalidThreadException("Attempted to send packet outside send thread!");
}
#endif
var writer = _sendPipe.Writer;
var writer = SendPipe.Writer;
try
{
@ -430,11 +429,15 @@ namespace Server.Network
if (buffer.Length > 0 && length > 0)
{
var result = writer.GetAvailable();
if (result.Length >= length)
if (!GetAvailableSendPipe(out var pipeBuffer))
{
result.CopyFrom(buffer.AsSpan(0, length));
p.OnSend();
return;
}
if (pipeBuffer.Length >= length)
{
pipeBuffer.CopyFrom(buffer.AsSpan(0, length));
writer.Advance((uint)length);
// Flush at the end of the game loop
@ -449,8 +452,6 @@ namespace Server.Network
{
WriteConsole("Didn't write anything!");
}
p.OnSend();
}
catch (Exception ex)
{
@ -460,6 +461,10 @@ namespace Server.Network
#endif
Dispose();
}
finally
{
p.OnSend();
}
}
internal void Start()
@ -477,22 +482,24 @@ namespace Server.Network
private async void SendTask(object state)
{
var reader = _sendPipe.Reader;
var reader = SendPipe.Reader;
var segments = new ArraySegment<byte>[2];
try
{
while (m_Running)
{
var result = await reader.Read();
if (!(await reader).TryRead(segments))
{
break;
}
if (result.Length <= 0)
if (segments[0].Count + segments[1].Count <= 0)
{
continue;
}
var buffer = result.Buffer;
var bytesWritten = await Connection.SendAsync(buffer, SocketFlags.None);
var bytesWritten = await Connection.SendAsync(segments, SocketFlags.None);
if (bytesWritten > 0)
{
@ -517,13 +524,14 @@ namespace Server.Network
private void DecodePacket(ArraySegment<byte>[] buffer, ref int length)
{
CircularBuffer<byte> cBuffer = new CircularBuffer<byte>(buffer);
_packetDecoder?.Invoke(cBuffer, ref length);
_packetDecoder?.Invoke(ref cBuffer, ref length);
}
private async void RecvTask(object state)
{
var socket = Connection;
var writer = _recvPipe.Writer;
var writer = RecvPipe.Writer;
var segments = new ArraySegment<byte>[2];
try
{
@ -534,20 +542,24 @@ namespace Server.Network
continue;
}
var result = writer.GetAvailable();
// TODO: Make awaitable
if (!writer.GetAvailable(segments))
{
break;
}
if (result.Length <= 0)
if (segments[0].Count + segments[1].Count <= 0)
{
continue;
}
var bytesWritten = await socket.ReceiveAsync(result.Buffer, SocketFlags.None);
var bytesWritten = await socket.ReceiveAsync(segments, SocketFlags.None);
if (bytesWritten <= 0)
{
break;
}
DecodePacket(result.Buffer, ref bytesWritten);
DecodePacket(segments, ref bytesWritten);
writer.Advance((uint)bytesWritten);
m_NextCheckActivity = Core.TickCount + 90000;
@ -587,19 +599,17 @@ namespace Server.Network
try
{
var reader = _recvPipe.Reader;
var reader = RecvPipe.Reader;
// Process as many packets as we can synchronously
while (true)
{
var result = reader.TryRead();
if (result.Length <= 0)
if (!reader.TryRead(out var buffer) || buffer.Length <= 0)
{
return;
}
var bytesProcessed = PacketHandlers.ProcessPacket(this, result.Buffer);
var bytesProcessed = PacketHandlers.ProcessPacket(this, ref buffer);
if (bytesProcessed <= 0)
{
@ -630,7 +640,7 @@ namespace Server.Network
{
if (Connection != null)
{
_sendPipe.Writer.Flush();
SendPipe.Writer.Flush();
}
}
@ -715,7 +725,7 @@ namespace Server.Network
return;
}
_sendPipe.Writer.Close();
SendPipe.Writer.Close();
try
{
@ -754,9 +764,9 @@ namespace Server.Network
ns.m_Running = false;
ns.Connection = null;
ns._recvBuffer = null;
ns._recvPipe = null;
ns.RecvPipe = null;
ns._sendBuffer = null;
ns._sendPipe = null;
ns.SendPipe = null;
ns.Gumps.Clear();
ns.Menus.Clear();
ns.HuePickers.Clear();

View file

@ -61,12 +61,12 @@ namespace Server.Network
0x4, 0x00D
};
public static void Compress(CircularBuffer<byte> buffer, ref int length)
public static void Compress(ref CircularBuffer<byte> buffer, ref int length)
{
length = Compress(buffer, length, buffer);
length = Compress(ref buffer, length, ref buffer);
}
public static int Compress(CircularBuffer<byte> input, int inputLength, CircularBuffer<byte> output)
public static int Compress(ref CircularBuffer<byte> input, int inputLength, ref CircularBuffer<byte> output)
{
if (inputLength > DefiniteOverflow)
{

View file

@ -89,8 +89,6 @@ namespace Server.Network
public void OnSend()
{
Core.Set(); // Is this still needed if this is done async?
if ((m_State & (State.Acquired | State.Static)) == 0)
{
Free();

View file

@ -14,6 +14,7 @@
*************************************************************************/
using System;
using System.Buffers;
using System.Collections.Generic;
using System.IO;
using Server.ContextMenus;
@ -46,8 +47,6 @@ namespace Server.Network
public static class PacketHandlers
{
public delegate void PlayCharCallback(NetState state, bool val);
private const int m_AuthIDWindowSize = 128;
private static readonly PacketHandler[] m_6017Handlers = new PacketHandler[0x100];
@ -154,9 +153,6 @@ namespace Server.Network
RegisterEncoded(0x32, true, QuestGumpRequest);
}
public static PlayCharCallback ThirdPartyAuthCallback { get; set; }
public static PlayCharCallback ThirdPartyHackedCallback { get; set; }
public static PacketHandler[] Handlers { get; } = new PacketHandler[0x100];
public static bool SingleClickProps { get; set; }
@ -281,9 +277,9 @@ namespace Server.Network
}
}
public static int ProcessPacket(NetState ns, ArraySegment<byte>[] segments)
public static int ProcessPacket(NetState ns, ref CircularBuffer<byte> buffer)
{
var reader = new CircularBufferReader(segments);
var reader = new CircularBufferReader(ref buffer);
var packetId = reader.ReadByte();

View file

@ -0,0 +1,31 @@
/*************************************************************************
* ModernUO *
* Copyright 2019-2020 - ModernUO Development Team *
* Email: hi@modernuo.com *
* File: PacketUtilities.cs *
* *
* This program is free software: you can redistribute it and/or modify *
* it under the terms of the GNU General Public License as published by *
* the Free Software Foundation, either version 3 of the License, or *
* (at your option) any later version. *
* *
* You should have received a copy of the GNU General Public License *
* along with this program. If not, see <http://www.gnu.org/licenses/>. *
*************************************************************************/
using System.Buffers;
using System.IO;
namespace Server.Network
{
public static class PacketUtilities
{
public static void WritePacketLength(this CircularBufferWriter writer)
{
var length = writer.Position;
writer.Seek(1, SeekOrigin.Begin);
writer.Write((ushort)length);
writer.Seek(length, SeekOrigin.Begin);
}
}
}

View file

@ -14,6 +14,7 @@
*************************************************************************/
using System;
using System.Buffers;
using System.Runtime.CompilerServices;
using System.Threading;
@ -32,90 +33,63 @@ namespace Server.Network
public class Pipe<T>
{
public struct Result<T>
{
public ArraySegment<T>[] Buffer { get; }
public bool IsClosed { get; set; }
public int Length
{
get
{
var length = 0;
for (int i = 0; i < Buffer.Length; i++)
{
length += Buffer[i].Count;
}
return length;
}
}
public void CopyFrom(ReadOnlySpan<T> bytes)
{
var remaining = bytes.Length;
var offset = 0;
if (remaining == 0)
{
return;
}
for (int i = 0; i < Buffer.Length; i++)
{
var buffer = Buffer[i];
var sz = Math.Min(remaining, buffer.Count);
bytes.Slice(offset, sz).CopyTo(buffer);
remaining -= sz;
offset += sz;
if (remaining == 0)
{
return;
}
}
throw new OutOfMemoryException();
}
public Result(int segments)
{
IsClosed = false;
Buffer = new ArraySegment<T>[segments];
}
}
public class PipeWriter<T>
{
private readonly Pipe<T> _pipe;
public PipeWriter(Pipe<T> pipe) => _pipe = pipe;
public Result<T> GetAvailable()
public bool GetAvailable(ArraySegment<T>[] segments)
{
var read = _pipe._readIdx;
var write = _pipe._writeIdx;
var result = new Result<T>(2) { IsClosed = _pipe._closed };
if (read <= write)
{
var readZero = read == 0;
var sz = _pipe.Size - write - (readZero ? 1 : 0);
result.Buffer[0] = sz == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)write, (int)sz);
result.Buffer[1] = readZero ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, 0, (int)read - 1);
segments[0] = sz == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)write, (int)sz);
segments[1] = readZero ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, 0, (int)read - 1);
}
else
{
var sz = read - write - 1;
result.Buffer[0] = sz == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)write, (int)sz);
result.Buffer[1] = ArraySegment<T>.Empty;
segments[0] = sz == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)write, (int)sz);
segments[1] = ArraySegment<T>.Empty;
}
return result;
return !_pipe._closed;
}
public bool GetAvailable(out CircularBuffer<T> buffer)
{
var read = _pipe._readIdx;
var write = _pipe._writeIdx;
Span<T> first;
Span<T> second;
if (read <= write)
{
var readZero = read == 0;
var sz = _pipe.Size - write - (readZero ? 1 : 0);
first = sz == 0 ? Span<T>.Empty : _pipe._buffer.AsSpan((int)write, (int)sz);
second = readZero ? Span<T>.Empty : _pipe._buffer.AsSpan(0, (int)read - 1);
}
else
{
var sz = read - write - 1;
first = sz == 0 ? Span<T>.Empty : _pipe._buffer.AsSpan((int)write, (int)sz);
second = Span<T>.Empty;
}
buffer = new CircularBuffer<T>(first, second);
return !_pipe._closed;
}
public void Advance(uint count)
@ -209,7 +183,7 @@ namespace Server.Network
}
}
public class PipeReader<T> : IPipeTask<Result<T>>
public class PipeReader<T> : IPipeTask<PipeReader<T>>
{
private readonly Pipe<T> _pipe;
@ -229,35 +203,46 @@ namespace Server.Network
return write + _pipe.Size - read;
}
public Result<T> TryRead()
public bool TryRead(out CircularBuffer<T> buffer)
{
var read = _pipe._readIdx;
var write = _pipe._writeIdx;
var result = new Result<T>(2) { IsClosed = _pipe._closed };
Span<T> first;
Span<T> second;
if (read <= write)
{
result.Buffer[0] = write - read == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)read, (int)(write - read));
result.Buffer[1] = ArraySegment<T>.Empty;
first = write - read == 0 ? Span<T>.Empty : _pipe._buffer.AsSpan((int)read, (int)(write - read));
second = Span<T>.Empty;
}
else
{
result.Buffer[0] = _pipe.Size - read == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)read, (int)(_pipe.Size - read));
result.Buffer[1] = write == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, 0, (int)write);
first = _pipe.Size - read == 0 ? Span<T>.Empty : _pipe._buffer.AsSpan((int)read, (int)(_pipe.Size - read));
second = write == 0 ? Span<T>.Empty : _pipe._buffer.AsSpan(0, (int)write);
}
return result;
buffer = new CircularBuffer<T>(first, second);
return !_pipe._closed;
}
public IPipeTask<Result<T>> Read()
public bool TryRead(ArraySegment<T>[] segments)
{
if (_pipe._awaitBeginning)
var read = _pipe._readIdx;
var write = _pipe._writeIdx;
if (read <= write)
{
throw new Exception("Double await on reader");
segments[0] = write - read == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)read, (int)(write - read));
segments[1] = ArraySegment<T>.Empty;
}
else
{
segments[0] = _pipe.Size - read == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, (int)read, (int)(_pipe.Size - read));
segments[1] = write == 0 ? ArraySegment<T>.Empty : new ArraySegment<T>(_pipe._buffer, 0, (int)write);
}
return this;
return !_pipe._closed;
}
public void Advance(uint count)
@ -303,7 +288,7 @@ namespace Server.Network
// The following makes it possible to await the reader. Do not use any of this directly.
public IPipeTask<Result<T>> GetAwaiter() => this;
public IPipeTask<PipeReader<T>> GetAwaiter() => this;
public bool IsCompleted
{
@ -319,7 +304,7 @@ namespace Server.Network
}
}
public Result<T> GetResult() => TryRead();
public PipeReader<T> GetResult() => this;
public void OnCompleted(Action continuation) => _pipe._readerContinuation = continuation;