From 990e6fe188059765ba7bfb98b46b57d6f59d5b85 Mon Sep 17 00:00:00 2001 From: Kamron Batman <3953314+kamronbatman@users.noreply.github.com> Date: Sun, 25 Oct 2020 17:40:08 -0700 Subject: [PATCH] Updates Pipe & Makes NetState more testable (#288) Updates pipe eliminate result and segments from being allocated. Bumps release version --- Projects/Server.Tests/Network/PipeTests.cs | 90 +++++++----- Projects/Server/Buffers/CircularBuffer.cs | 12 +- .../Server/Buffers/CircularBufferReader.cs | 3 +- .../Server/Buffers/CircularBufferWriter.cs | 4 - Projects/Server/Items/Item.cs | 2 - Projects/Server/Main.cs | 5 - Projects/Server/Mobiles/Mobile.cs | 2 - Projects/Server/Network/NetState/NetState.cs | 104 +++++++------ Projects/Server/Network/NetworkCompression.cs | 6 +- Projects/Server/Network/Packet.cs | 2 - Projects/Server/Network/PacketHandlers.cs | 10 +- Projects/Server/Network/PacketUtilities.cs | 31 ++++ Projects/Server/Network/Pipe.cs | 137 ++++++++---------- Projects/Server/Timer/Timer.cs | 9 -- 14 files changed, 219 insertions(+), 198 deletions(-) create mode 100644 Projects/Server/Network/PacketUtilities.cs diff --git a/Projects/Server.Tests/Network/PipeTests.cs b/Projects/Server.Tests/Network/PipeTests.cs index 11b870362..fc2b31638 100644 --- a/Projects/Server.Tests/Network/PipeTests.cs +++ b/Projects/Server.Tests/Network/PipeTests.cs @@ -26,7 +26,7 @@ namespace Server.Tests.Network DelayedExecute(() => { // Write some data into the pipe - var buffer = writer.GetAvailable(); + writer.GetAvailable(out var buffer); Assert.True(buffer.Length == 99); buffer.CopyFrom(new byte[] { 1 }); @@ -37,9 +37,10 @@ namespace Server.Tests.Network writer.Flush(); }); - var result = await reader; + var segments = new ArraySegment[2]; + (await reader).TryRead(segments); - Assert.True(result.Buffer[0].Count == 3); + Assert.True(segments[0].Count == 3); } private bool _signal; @@ -53,11 +54,13 @@ namespace Server.Tests.Network while (count < 0x8000000) { - var result = reader.TryRead(); + reader.TryRead(out var buffer); - for (int i = 0; i < result.Buffer[0].Count; i++) + var first = buffer.GetSpan(0); + + for (int i = 0; i < first.Length; i++) { - Assert.True(result.Buffer[0][i] == expected_value); + Assert.True(first[i] == expected_value); count++; if (count == 0x1000) @@ -66,9 +69,10 @@ namespace Server.Tests.Network } } - for (int i = 0; i < result.Buffer[1].Count; i++) + var second = buffer.GetSpan(1); + for (int i = 0; i < second.Length; i++) { - Assert.True(result.Buffer[1][i] == expected_value); + Assert.True(second[i] == expected_value); count++; if (count == 0x1000) @@ -77,14 +81,14 @@ namespace Server.Tests.Network } } - reader.Advance((uint)result.Length); + reader.Advance((uint)buffer.Length); } _signal = true; } [Fact] - public async void Threading() + public void Threading() { var pipe = new Pipe(new byte[0x1001]); @@ -97,15 +101,19 @@ namespace Server.Tests.Network while (count < 0x8000000) { - var result = writer.GetAvailable(); + writer.GetAvailable(out var buffer); - if (result.Length < 16) + if (buffer.Length < 16) { continue; } - result.CopyFrom(new[] { expected_value, expected_value, expected_value, expected_value, expected_value, expected_value, expected_value, expected_value, - expected_value, expected_value, expected_value, expected_value, expected_value, expected_value, expected_value, expected_value }); + buffer.CopyFrom(new[] { + expected_value, expected_value, expected_value, expected_value, + expected_value, expected_value, expected_value, expected_value, + expected_value, expected_value, expected_value, expected_value, + expected_value, expected_value, expected_value, expected_value + }); writer.Advance(16); count += 16; @@ -129,32 +137,33 @@ namespace Server.Tests.Network var reader = pipe.Reader; var writer = pipe.Writer; - var result = writer.GetAvailable(); - Assert.True(result.Length == 9); + writer.GetAvailable(out var buffer); + + Assert.True(buffer.Length == 9); Assert.True(reader.GetAvailable() == 0); - result = reader.TryRead(); - Assert.True(result.Length == 0); + reader.TryRead(out buffer); + Assert.True(buffer.Length == 0); writer.Advance(7); - result = writer.GetAvailable(); - Assert.True(result.Length == 2); + writer.GetAvailable(out buffer); + Assert.True(buffer.Length == 2); Assert.True(reader.GetAvailable() == 7); - result = reader.TryRead(); - Assert.True(result.Length == 7); + reader.TryRead(out buffer); + Assert.True(buffer.Length == 7); reader.Advance(4); - result = writer.GetAvailable(); - Assert.True(result.Length == 6); + writer.GetAvailable(out buffer); + Assert.True(buffer.Length == 6); Assert.True(reader.GetAvailable() == 3); - result = reader.TryRead(); - Assert.True(result.Length == 3); + reader.TryRead(out buffer); + Assert.True(buffer.Length == 3); writer.Advance(3); - result = writer.GetAvailable(); - Assert.True(result.Length == 3); + writer.GetAvailable(out buffer); + Assert.True(buffer.Length == 3); Assert.True(reader.GetAvailable() == 6); - result = reader.TryRead(); - Assert.True(result.Length == 6); + reader.TryRead(out buffer); + Assert.True(buffer.Length == 6); } @@ -184,7 +193,7 @@ namespace Server.Tests.Network var reader = pipe.Reader; var writer = pipe.Writer; - var buffer = writer.GetAvailable(); + writer.GetAvailable(out var buffer); Assert.True(buffer.Length == 9); buffer.CopyFrom(new byte[] { 0, 1, 2, 3, 4, 5, 6, 7, 8 }); @@ -194,21 +203,24 @@ namespace Server.Tests.Network Assert.True(reader.GetAvailable() == 9); - buffer = reader.TryRead(); + reader.TryRead(out buffer); + + var first = buffer.GetSpan(0); for (int i = 0; i < 9; i++) { - Assert.True(buffer.Buffer[0][i] == i); + Assert.True(first[i] == i); } reader.Advance(4); - buffer = reader.TryRead(); + reader.TryRead(out buffer); + first = buffer.GetSpan(0); Assert.True(buffer.Length == 5); - Assert.True(buffer.Buffer[0][0] == 4); - Assert.True(buffer.Buffer[0][1] == 5); - Assert.True(buffer.Buffer[0][2] == 6); - Assert.True(buffer.Buffer[0][3] == 7); - Assert.True(buffer.Buffer[0][4] == 8); + Assert.True(first[0] == 4); + Assert.True(first[1] == 5); + Assert.True(first[2] == 6); + Assert.True(first[3] == 7); + Assert.True(first[4] == 8); } } } diff --git a/Projects/Server/Buffers/CircularBuffer.cs b/Projects/Server/Buffers/CircularBuffer.cs index b34240490..a5e25059a 100644 --- a/Projects/Server/Buffers/CircularBuffer.cs +++ b/Projects/Server/Buffers/CircularBuffer.cs @@ -15,7 +15,7 @@ namespace System.Buffers { - public readonly ref struct CircularBuffer where T : struct + public readonly ref struct CircularBuffer { private readonly Span _first; private readonly Span _second; @@ -126,5 +126,15 @@ namespace System.Buffers return new CircularBuffer(first, second); } + + public Span GetSpan(int index) + { + if (index < 0 || index > 1) + { + throw new ArgumentOutOfRangeException(nameof(index)); + } + + return index == 0 ? _first : _second; + } } } diff --git a/Projects/Server/Buffers/CircularBufferReader.cs b/Projects/Server/Buffers/CircularBufferReader.cs index 80d2e3c25..711cbed17 100644 --- a/Projects/Server/Buffers/CircularBufferReader.cs +++ b/Projects/Server/Buffers/CircularBufferReader.cs @@ -14,6 +14,7 @@ *************************************************************************/ using System; +using System.Buffers; using System.Buffers.Binary; using System.IO; using System.Runtime.CompilerServices; @@ -30,7 +31,7 @@ namespace Server.Network public int Position { get; private set; } public int Remaining => Length - Position; - public CircularBufferReader(ArraySegment[] buffers) : this(buffers[0], buffers[1]) + public CircularBufferReader(ref CircularBuffer buffer) : this(buffer.GetSpan(0), buffer.GetSpan(1)) { } diff --git a/Projects/Server/Buffers/CircularBufferWriter.cs b/Projects/Server/Buffers/CircularBufferWriter.cs index 542c72b8a..5ec3742ee 100644 --- a/Projects/Server/Buffers/CircularBufferWriter.cs +++ b/Projects/Server/Buffers/CircularBufferWriter.cs @@ -29,10 +29,6 @@ namespace System.Buffers public int Length { get; } public int Position { get; private set; } - public CircularBufferWriter(ArraySegment[] buffers) : this(buffers[0], buffers[1]) - { - } - public CircularBufferWriter(Span first, Span second) { _first = first; diff --git a/Projects/Server/Items/Item.cs b/Projects/Server/Items/Item.cs index e9d2bf47f..5a4b10f8d 100644 --- a/Projects/Server/Items/Item.cs +++ b/Projects/Server/Items/Item.cs @@ -3361,8 +3361,6 @@ namespace Server m_DeltaQueue.Add(this); } } - - Core.Set(); } public void RemDelta(ItemDelta flags) diff --git a/Projects/Server/Main.cs b/Projects/Server/Main.cs index 24694cf2c..fa4c8833c 100644 --- a/Projects/Server/Main.cs +++ b/Projects/Server/Main.cs @@ -325,11 +325,6 @@ namespace Server Console.WriteLine("done"); } - public static void Set() - { - // m_Signal.Set(); - } - public static void Main(string[] args) { AppDomain.CurrentDomain.UnhandledException += CurrentDomain_UnhandledException; diff --git a/Projects/Server/Mobiles/Mobile.cs b/Projects/Server/Mobiles/Mobile.cs index 40afe92c1..2e56a79b7 100644 --- a/Projects/Server/Mobiles/Mobile.cs +++ b/Projects/Server/Mobiles/Mobile.cs @@ -8524,8 +8524,6 @@ namespace Server m_DeltaQueue.Enqueue(this); } } - - Core.Set(); } public static void ProcessDeltaQueue() diff --git a/Projects/Server/Network/NetState/NetState.cs b/Projects/Server/Network/NetState/NetState.cs index da985d191..2aeaab148 100644 --- a/Projects/Server/Network/NetState/NetState.cs +++ b/Projects/Server/Network/NetState/NetState.cs @@ -33,9 +33,9 @@ namespace Server.Network { public delegate void NetStateCreatedCallback(NetState ns); - public delegate void EncodePacket(CircularBuffer buffer, ref int length); + public delegate void EncodePacket(ref CircularBuffer buffer, ref int length); - public partial class NetState : IComparable + public partial class NetState : IComparable, 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 _recvPipe; private byte[] _sendBuffer; - private Pipe _sendPipe; private long m_NextCheckActivity; private volatile bool m_Running; private readonly Thread _sendThread; @@ -89,9 +87,9 @@ namespace Server.Network Menus = new List(); Trades = new List(); _recvBuffer = new byte[RecvPipeSize]; - _recvPipe = new Pipe(_recvBuffer); + RecvPipe = new Pipe(_recvBuffer); _sendBuffer = new byte[SendPipeSize]; - _sendPipe = new Pipe(_sendBuffer); + SendPipe = new Pipe(_sendBuffer); m_NextCheckActivity = Core.TickCount + 30000; _sendThread = sendThread ?? Core.Thread; @@ -142,6 +140,10 @@ namespace Server.Network public bool Seeded { get; set; } + public Pipe RecvPipe { get; private set; } + + public Pipe 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.Result GetAvailableSendPipe() => _recvPipe.Writer.GetAvailable(); + public bool GetAvailableSendPipe(out CircularBuffer buffer) => SendPipe.Writer.GetAvailable(out buffer); - public virtual void Send(CircularBuffer buffer, int length) + public virtual void Send(ref CircularBuffer 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[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[] buffer, ref int length) { CircularBuffer cBuffer = new CircularBuffer(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[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(); diff --git a/Projects/Server/Network/NetworkCompression.cs b/Projects/Server/Network/NetworkCompression.cs index 7c78d0d2c..49cfe1646 100644 --- a/Projects/Server/Network/NetworkCompression.cs +++ b/Projects/Server/Network/NetworkCompression.cs @@ -61,12 +61,12 @@ namespace Server.Network 0x4, 0x00D }; - public static void Compress(CircularBuffer buffer, ref int length) + public static void Compress(ref CircularBuffer buffer, ref int length) { - length = Compress(buffer, length, buffer); + length = Compress(ref buffer, length, ref buffer); } - public static int Compress(CircularBuffer input, int inputLength, CircularBuffer output) + public static int Compress(ref CircularBuffer input, int inputLength, ref CircularBuffer output) { if (inputLength > DefiniteOverflow) { diff --git a/Projects/Server/Network/Packet.cs b/Projects/Server/Network/Packet.cs index 9676e6cc8..0eaf16a34 100644 --- a/Projects/Server/Network/Packet.cs +++ b/Projects/Server/Network/Packet.cs @@ -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(); diff --git a/Projects/Server/Network/PacketHandlers.cs b/Projects/Server/Network/PacketHandlers.cs index e1386bdde..350525664 100644 --- a/Projects/Server/Network/PacketHandlers.cs +++ b/Projects/Server/Network/PacketHandlers.cs @@ -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[] segments) + public static int ProcessPacket(NetState ns, ref CircularBuffer buffer) { - var reader = new CircularBufferReader(segments); + var reader = new CircularBufferReader(ref buffer); var packetId = reader.ReadByte(); diff --git a/Projects/Server/Network/PacketUtilities.cs b/Projects/Server/Network/PacketUtilities.cs new file mode 100644 index 000000000..0fda88dce --- /dev/null +++ b/Projects/Server/Network/PacketUtilities.cs @@ -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 . * + *************************************************************************/ + +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); + } + } +} diff --git a/Projects/Server/Network/Pipe.cs b/Projects/Server/Network/Pipe.cs index cf11309ac..4ab0db2f8 100644 --- a/Projects/Server/Network/Pipe.cs +++ b/Projects/Server/Network/Pipe.cs @@ -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 { - public struct Result - { - public ArraySegment[] 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 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[segments]; - } - } - public class PipeWriter { private readonly Pipe _pipe; public PipeWriter(Pipe pipe) => _pipe = pipe; - public Result GetAvailable() + public bool GetAvailable(ArraySegment[] segments) { var read = _pipe._readIdx; var write = _pipe._writeIdx; - var result = new Result(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.Empty : new ArraySegment(_pipe._buffer, (int)write, (int)sz); - result.Buffer[1] = readZero ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, 0, (int)read - 1); + segments[0] = sz == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, (int)write, (int)sz); + segments[1] = readZero ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, 0, (int)read - 1); } else { var sz = read - write - 1; - result.Buffer[0] = sz == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, (int)write, (int)sz); - result.Buffer[1] = ArraySegment.Empty; + segments[0] = sz == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, (int)write, (int)sz); + segments[1] = ArraySegment.Empty; } - return result; + return !_pipe._closed; + } + + public bool GetAvailable(out CircularBuffer buffer) + { + var read = _pipe._readIdx; + var write = _pipe._writeIdx; + + Span first; + Span second; + + if (read <= write) + { + var readZero = read == 0; + var sz = _pipe.Size - write - (readZero ? 1 : 0); + + first = sz == 0 ? Span.Empty : _pipe._buffer.AsSpan((int)write, (int)sz); + second = readZero ? Span.Empty : _pipe._buffer.AsSpan(0, (int)read - 1); + } + else + { + var sz = read - write - 1; + + first = sz == 0 ? Span.Empty : _pipe._buffer.AsSpan((int)write, (int)sz); + second = Span.Empty; + } + + buffer = new CircularBuffer(first, second); + + return !_pipe._closed; } public void Advance(uint count) @@ -209,7 +183,7 @@ namespace Server.Network } } - public class PipeReader : IPipeTask> + public class PipeReader : IPipeTask> { private readonly Pipe _pipe; @@ -229,35 +203,46 @@ namespace Server.Network return write + _pipe.Size - read; } - public Result TryRead() + public bool TryRead(out CircularBuffer buffer) { var read = _pipe._readIdx; var write = _pipe._writeIdx; - var result = new Result(2) { IsClosed = _pipe._closed }; + Span first; + Span second; if (read <= write) { - result.Buffer[0] = write - read == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, (int)read, (int)(write - read)); - result.Buffer[1] = ArraySegment.Empty; + first = write - read == 0 ? Span.Empty : _pipe._buffer.AsSpan((int)read, (int)(write - read)); + second = Span.Empty; } else { - result.Buffer[0] = _pipe.Size - read == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, (int)read, (int)(_pipe.Size - read)); - result.Buffer[1] = write == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, 0, (int)write); + first = _pipe.Size - read == 0 ? Span.Empty : _pipe._buffer.AsSpan((int)read, (int)(_pipe.Size - read)); + second = write == 0 ? Span.Empty : _pipe._buffer.AsSpan(0, (int)write); } - return result; + buffer = new CircularBuffer(first, second); + return !_pipe._closed; } - public IPipeTask> Read() + public bool TryRead(ArraySegment[] 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.Empty : new ArraySegment(_pipe._buffer, (int)read, (int)(write - read)); + segments[1] = ArraySegment.Empty; + } + else + { + segments[0] = _pipe.Size - read == 0 ? ArraySegment.Empty : new ArraySegment(_pipe._buffer, (int)read, (int)(_pipe.Size - read)); + segments[1] = write == 0 ? ArraySegment.Empty : new ArraySegment(_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> GetAwaiter() => this; + public IPipeTask> GetAwaiter() => this; public bool IsCompleted { @@ -319,7 +304,7 @@ namespace Server.Network } } - public Result GetResult() => TryRead(); + public PipeReader GetResult() => this; public void OnCompleted(Action continuation) => _pipe._readerContinuation = continuation; diff --git a/Projects/Server/Timer/Timer.cs b/Projects/Server/Timer/Timer.cs index f0c5e49c1..eae982389 100644 --- a/Projects/Server/Timer/Timer.cs +++ b/Projects/Server/Timer/Timer.cs @@ -402,8 +402,6 @@ namespace Server ProcessChanged(); - var loaded = false; - for (var i = 0; i < m_Timers.Length; i++) { var now = Core.TickCount; @@ -427,8 +425,6 @@ namespace Server m_Queue.Enqueue(t); } - loaded = true; - if (t.m_Count != 0 && ++t.m_Index >= t.m_Count) { t.Stop(); @@ -441,11 +437,6 @@ namespace Server } } - if (loaded) - { - Core.Set(); - } - m_Signal.WaitOne(1, false); } }