From f7ab35529de74fd9b5cd18843ef5331b31790e6c Mon Sep 17 00:00:00 2001 From: leroysquad <263533951+leroysquad@users.noreply.github.com> Date: Wed, 30 Sep 2026 17:52:22 -0700 Subject: [PATCH 1/4] Send packets through the per-connection queue by default. Direct Socket.SendAsync on the gameplay thread was half of main-thread samples at 1000 players on a 130k-chunk world. The existing drain already runs on the thread pool. Co-authored-by: Cursor --- .../Vintagestory.Server/StratumConfig.cs | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs index e9d36c7..8062f6d 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs @@ -632,11 +632,12 @@ public void EnsurePopulated() internal class StratumNetworkConfig { - // Off by default. The queue removes the packet reordering bug from the old - // StratumNetworkFlush (disabled after client crashes, see PR #138), but it is new code - // on the hottest path in the server. Soak on a community server before flipping the - // shipped default. - public bool SendQueueEnabled { get; set; } = false; + // On by default. With the queue off, TcpNetConnection.Send calls Socket.SendAsync on + // the gameplay thread. A 1000-player capture on a 130k-chunk world spent about half of + // ServerMain.Process samples in SocketAsyncEventArgs.DoOperationSendSingleBuffer. + // The queue is the only sender for the connection and runs on the thread pool. Set + // this false to restore the direct send path. See #325. + public bool SendQueueEnabled { get; set; } = true; // Packets at or above this size skip coalescing and go out alone. Matches the TCP MTU // assumption the old flush buffer used. From 7b9b567388c0c0eac9e429392e888e370f10e31c Mon Sep 17 00:00:00 2001 From: leroysquad <263533951+leroysquad@users.noreply.github.com> Date: Wed, 30 Sep 2026 22:26:05 -0700 Subject: [PATCH 2/4] fix: bound the send queue and finish short socket writes A slow client could retain unbounded packet arrays, and a short Socket.SendAsync write could split a length-prefixed frame. Co-authored-by: Cursor --- .../Vintagestory.Server/StratumConfig.cs | 9 ++++ .../Vintagestory.Server/StratumSendQueue.cs | 54 +++++++++++++++++-- 2 files changed, 60 insertions(+), 3 deletions(-) diff --git a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs index 8062f6d..cfeaf21 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs @@ -647,10 +647,19 @@ internal class StratumNetworkConfig // a single SendAsync call carries. public int CoalesceLimitBytes { get; set; } = 65536; + // Bytes StratumSendQueue may retain for one connection. Chunk scheduling already + // slows that client at OutboundPressurePendingBytesHardLimit (1 MiB) and still + // sends a minimum budget there, so this cap sits above that line. An enqueue that + // would grow a non-empty queue past the cap is refused and the connection is + // disconnected. One packet is still accepted when the queue is empty, so a single + // large send is not dropped on the floor. See #345. + public int MaxPendingBytes { get; set; } = 8 * 1024 * 1024; + public void EnsureSane() { LargeThresholdBytes = Math.Max(64, LargeThresholdBytes); CoalesceLimitBytes = Math.Max(LargeThresholdBytes, CoalesceLimitBytes); + MaxPendingBytes = Math.Max(CoalesceLimitBytes, MaxPendingBytes); } } diff --git a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs index 0be92a7..e2da223 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs @@ -30,9 +30,11 @@ internal sealed class StratumSendQueue private readonly CancellationToken cancellationToken; private readonly int largeThreshold; private readonly int coalesceLimit; + private readonly long maxPendingBytes; private int pendingCount; private long pendingBytes; + private int overflowDisconnectStarted; public int PendingCount => Volatile.Read(ref pendingCount); @@ -46,6 +48,7 @@ public StratumSendQueue(TcpNetConnection connection, Socket socket, Cancellation StratumNetworkConfig config = StratumRuntime.Config.Performance.Network; largeThreshold = config.LargeThresholdBytes; coalesceLimit = config.CoalesceLimitBytes; + maxPendingBytes = config.MaxPendingBytes; channel = Channel.CreateUnbounded(new UnboundedChannelOptions { SingleReader = true, @@ -63,15 +66,44 @@ public void Start() // they will not mutate again (a fresh array per send), so no copy is needed here. public void Enqueue(byte[] dataWithLength) { + if (Volatile.Read(ref overflowDisconnectStarted) != 0) + { + return; + } + + int length = dataWithLength.Length; + long queued = Interlocked.Add(ref pendingBytes, length); + // Queue was already holding data and this packet would pass the cap. Roll it + // back and disconnect. Dropping the packet without closing would split the + // length-prefixed stream. Blocking here would put the gameplay thread back + // on the slow client. An empty queue still accepts one packet so a single + // large send is not refused. + if (queued > maxPendingBytes && queued - length > 0) + { + Interlocked.Add(ref pendingBytes, -length); + DisconnectOverflow(); + return; + } + Interlocked.Increment(ref pendingCount); - Interlocked.Add(ref pendingBytes, dataWithLength.Length); if (!channel.Writer.TryWrite(dataWithLength)) { // Writer already completed (connection closing). Drop, matches vanilla behavior // of a send attempted after Close()/Dispose(). Interlocked.Decrement(ref pendingCount); - Interlocked.Add(ref pendingBytes, -dataWithLength.Length); + Interlocked.Add(ref pendingBytes, -length); + } + } + + private void DisconnectOverflow() + { + if (Interlocked.Exchange(ref overflowDisconnectStarted, 1) != 0) + { + return; } + + StratumRuntime.LogWarning("StratumSendQueue disconnected a connection: pending bytes exceeded MaxPendingBytes (" + maxPendingBytes + ")."); + connection.InvokeDisconnected(); } public void Complete() @@ -138,7 +170,23 @@ private async Task SendAsync(byte[] buffer, int length) { try { - await socket.SendAsync(new ReadOnlyMemory(buffer, 0, length), SocketFlags.None, cancellationToken).ConfigureAwait(false); + // This overload returns the number of bytes accepted. A short write leaves a + // truncated length-prefixed frame, and the next packet would be parsed as the + // rest of it. Retry the tail. A non-positive count means the socket took + // nothing; disconnect instead of spinning. + int offset = 0; + while (offset < length) + { + int sent = await socket.SendAsync(new ReadOnlyMemory(buffer, offset, length - offset), SocketFlags.None, cancellationToken).ConfigureAwait(false); + if (sent <= 0) + { + connection.InvokeDisconnected(); + return false; + } + + offset += sent; + } + return true; } catch (OperationCanceledException) From 56f7e92603a6426808e917b9873d0fb4c73588ea Mon Sep 17 00:00:00 2001 From: leroysquad <263533951+leroysquad@users.noreply.github.com> Date: Thu, 1 Oct 2026 22:50:44 -0700 Subject: [PATCH 3/4] fix: flush the send queue before socket shutdown A disconnect reason was only enqueued, and Shutdown closed the socket before the drain could write it. Co-authored-by: Cursor --- .../TcpNetConnection.cs.patch | 17 +++- .../Vintagestory.Server/StratumSendQueue.cs | 93 ++++++++++++++++--- .../SendQueueDisconnectScenarios.cs | 49 ++++++++++ 3 files changed, 143 insertions(+), 16 deletions(-) create mode 100644 tests/StratumScenarios/SendQueueDisconnectScenarios.cs diff --git a/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch b/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch index 6f09bce..bd5b962 100644 --- a/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch +++ b/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch @@ -35,8 +35,8 @@ index ae86835..34c0f3d 100644 cts = new CancellationTokenSource(); + if (StratumRuntime.Config.Performance.Network.SendQueueEnabled) + { ++ // The drain task starts on the first enqueued packet, not here. + stratumSendQueue = new StratumSendQueue(this, TcpSocket, cts.Token); -+ stratumSendQueue.Start(); + } TyronThreadPool.QueueTask((Func)ReceiveData, "TcpNetConReceiveData"); } @@ -88,19 +88,26 @@ index ae86835..34c0f3d 100644 catch { InvokeDisconnected(); -@@ -281,10 +313,11 @@ public class TcpNetConnection : NetConnection +@@ -269,4 +301,6 @@ public class TcpNetConnection : NetConnection + public override void Shutdown() + { ++ // Stratum: put queued bytes, including a disconnect reason, on the socket before FIN. ++ stratumSendQueue?.FlushPending(250); + if (TcpSocket == null) + { +@@ -281,10 +315,11 @@ public class TcpNetConnection : NetConnection } } - + public override void Close() { -+ stratumSendQueue?.Complete(); ++ stratumSendQueue?.FlushPending(250); try { cts?.Cancel(); } catch -@@ -302,10 +335,11 @@ public class TcpNetConnection : NetConnection +@@ -302,10 +337,11 @@ public class TcpNetConnection : NetConnection internal void InvokeDisconnected() { diff --git a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs index e2da223..67929f8 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs @@ -32,6 +32,10 @@ internal sealed class StratumSendQueue private readonly int coalesceLimit; private readonly long maxPendingBytes; + private readonly object startGate = new object(); + private Task drainTask; + private int closed; + private int pendingCount; private long pendingBytes; private int overflowDisconnectStarted; @@ -57,9 +61,57 @@ public StratumSendQueue(TcpNetConnection connection, Socket socket, Cancellation }); } - public void Start() + // The drain task and its coalesce buffer stay unallocated until the first packet. + // A connection that never sends (a scanner holding the socket open) then costs the + // socket alone, not a 64 KiB buffer and a thread-pool task. + private void EnsureDrainStarted() { - TyronThreadPool.QueueTask((Func)DrainAsync, "StratumSendQueueDrain"); + if (drainTask != null) + { + return; + } + + var done = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + drainTask = done.Task; + TyronThreadPool.QueueTask(async () => + { + try + { + await DrainAsync().ConfigureAwait(false); + done.TrySetResult(); + } + catch (Exception ex) + { + done.TrySetException(ex); + } + }, "StratumSendQueueDrain"); + } + + // Finish bytes already queued, then return. Close and Shutdown call this before they + // cancel the socket, so a disconnect reason enqueued just above them still reaches + // the client. The wait is bounded: a wedged socket does not hold the caller forever. + public void FlushPending(int timeoutMs) + { + Task task; + lock (startGate) + { + closed = 1; + channel.Writer.TryComplete(); + task = drainTask; + } + + if (task == null) + { + return; + } + + try + { + task.Wait(Math.Max(0, timeoutMs)); + } + catch (AggregateException) + { + } } // dataWithLength must already carry the 4-byte length prefix. Callers hand off a buffer @@ -86,12 +138,18 @@ public void Enqueue(byte[] dataWithLength) } Interlocked.Increment(ref pendingCount); - if (!channel.Writer.TryWrite(dataWithLength)) + lock (startGate) { - // Writer already completed (connection closing). Drop, matches vanilla behavior - // of a send attempted after Close()/Dispose(). - Interlocked.Decrement(ref pendingCount); - Interlocked.Add(ref pendingBytes, -length); + if (closed != 0 || !channel.Writer.TryWrite(dataWithLength)) + { + // Writer already completed (connection closing). Drop, matches vanilla behavior + // of a send attempted after Close()/Dispose(). + Interlocked.Decrement(ref pendingCount); + Interlocked.Add(ref pendingBytes, -length); + return; + } + + EnsureDrainStarted(); } } @@ -102,7 +160,19 @@ private void DisconnectOverflow() return; } - StratumRuntime.LogWarning("StratumSendQueue disconnected a connection: pending bytes exceeded MaxPendingBytes (" + maxPendingBytes + ")."); + string player = connection.client?.PlayerName; + if (string.IsNullOrEmpty(player)) + { + player = "unidentified"; + } + + string address = connection.Address; + if (string.IsNullOrEmpty(address)) + { + address = connection.TcpSocket?.RemoteEndPoint?.ToString() ?? "unknown"; + } + + StratumRuntime.LogWarning("StratumSendQueue disconnected " + player + " at " + address + ": pending bytes exceeded Performance.Network.MaxPendingBytes (" + maxPendingBytes + ")."); connection.InvokeDisconnected(); } @@ -114,13 +184,13 @@ public void Complete() private async Task DrainAsync() { ChannelReader reader = channel.Reader; - byte[] coalesceBuffer = new byte[coalesceLimit]; + byte[] coalesceBuffer = null; try { while (await reader.WaitToReadAsync(cancellationToken).ConfigureAwait(false)) { int coalescedLength = 0; - while (coalescedLength < coalesceBuffer.Length && reader.TryRead(out byte[] packet)) + while (coalescedLength < coalesceLimit && reader.TryRead(out byte[] packet)) { if (packet.Length >= largeThreshold) { @@ -137,7 +207,8 @@ private async Task DrainAsync() continue; } - if (coalescedLength + packet.Length > coalesceBuffer.Length) + coalesceBuffer ??= new byte[coalesceLimit]; + if (coalescedLength + packet.Length > coalesceLimit) { if (!await SendAsync(coalesceBuffer, coalescedLength).ConfigureAwait(false)) return; coalescedLength = 0; diff --git a/tests/StratumScenarios/SendQueueDisconnectScenarios.cs b/tests/StratumScenarios/SendQueueDisconnectScenarios.cs new file mode 100644 index 0000000..2c15d52 --- /dev/null +++ b/tests/StratumScenarios/SendQueueDisconnectScenarios.cs @@ -0,0 +1,49 @@ +using System.Net; +using System.Net.Sockets; +using System.Reflection; +using System.Text; +using Atlas.XUnit; +using Xunit; + +namespace StratumScenarios; + +/// +/// With the send queue on, DisconnectPlayer only enqueues the reason and CloseConnection +/// shuts the socket down immediately. The reason has to be on the wire before that FIN. +/// +public class SendQueueDisconnectScenarios : AtlasScenarioBase +{ + [AtlasScenario(TimeoutMs = 60_000)] + public void Shutdown_Should_DeliverQueuedBytes_When_SendQueueEnabled() + { + Type connectionType = Type.GetType("Vintagestory.Server.Network.TcpNetConnection, VintagestoryLib") + ?? throw new InvalidOperationException("TcpNetConnection was not in the loaded lib"); + + using var listener = new TcpListener(IPAddress.Loopback, 0); + listener.Start(); + int port = ((IPEndPoint)listener.LocalEndpoint).Port; + + using var client = new TcpClient(); + client.Connect(IPAddress.Loopback, port); + using Socket accepted = listener.AcceptSocket(); + client.Client.ReceiveTimeout = 2000; + + object connection = Activator.CreateInstance(connectionType, accepted)!; + connectionType.GetMethod("StartReceiving")!.Invoke(connection, null); + object queue = connectionType.GetField("stratumSendQueue", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(connection) + ?? throw new InvalidOperationException("SendQueueEnabled did not attach a queue"); + + byte[] payload = Encoding.ASCII.GetBytes("kick-reason"); + MethodInfo send = connectionType.GetMethod("Send", new[] { typeof(byte[]), typeof(bool) })!; + send.Invoke(connection, new object[] { payload, false }); + connectionType.GetMethod("Shutdown")!.Invoke(connection, null); + + byte[] buffer = new byte[64]; + int read = client.Client.Receive(buffer); + string text = Encoding.ASCII.GetString(buffer, 0, read); + Assert.Contains("kick-reason", text, StringComparison.Ordinal); + + connectionType.GetMethod("Close")!.Invoke(connection, null); + GC.KeepAlive(queue); + } +} From 3acd1f96f62165acc8b837c88230b0e377419835 Mon Sep 17 00:00:00 2001 From: leroysquad <263533951+leroysquad@users.noreply.github.com> Date: Sat, 3 Oct 2026 21:18:20 -0700 Subject: [PATCH 4/4] fix: stop send-queue shutdown from stalling the caller The queue stays off until it has a load run on this code. Shutdown returns immediately and closes the socket after the drain, query sockets always close, and one oversized packet is not counted against the byte cap. --- .../TcpNetConnection.cs.patch | 28 +++- .../Vintagestory.Server/ServerMain.cs.patch | 9 +- .../Vintagestory.Server/StratumConfig.cs | 23 ++- .../Vintagestory.Server/StratumSendQueue.cs | 91 +++++++++--- .../SendQueueDisconnectScenarios.cs | 132 ++++++++++++++---- .../stratum-performance.json | 5 + .../stratum-sendqueue-on/stratum.json | 3 + 7 files changed, 222 insertions(+), 69 deletions(-) create mode 100644 tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum-performance.json create mode 100644 tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum.json diff --git a/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch b/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch index bd5b962..6ff211d 100644 --- a/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch +++ b/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch @@ -88,26 +88,42 @@ index ae86835..34c0f3d 100644 catch { InvokeDisconnected(); -@@ -269,4 +301,6 @@ public class TcpNetConnection : NetConnection +@@ -269,4 +301,21 @@ public class TcpNetConnection : NetConnection public override void Shutdown() { -+ // Stratum: put queued bytes, including a disconnect reason, on the socket before FIN. -+ stratumSendQueue?.FlushPending(250); ++ // Stratum: complete the writer and return. The FIN runs after the drain ++ // finishes, or after 250 ms, whichever is first. Waiting here stalled the ++ // main thread for the full timeout when the peer never read. ++ if (stratumSendQueue != null) ++ { ++ stratumSendQueue.ScheduleShutdown(250, delegate ++ { ++ try ++ { ++ TcpSocket?.Shutdown(SocketShutdown.Both); ++ } ++ catch ++ { ++ } ++ }); ++ return; ++ } if (TcpSocket == null) { -@@ -281,10 +315,11 @@ public class TcpNetConnection : NetConnection +@@ -281,10 +329,12 @@ public class TcpNetConnection : NetConnection } } public override void Close() { -+ stratumSendQueue?.FlushPending(250); ++ // Stratum: do not wait. Cancel is the hard bound. Shutdown already scheduled the FIN. ++ stratumSendQueue?.Complete(); try { cts?.Cancel(); } catch -@@ -302,10 +337,11 @@ public class TcpNetConnection : NetConnection +@@ -302,10 +353,11 @@ public class TcpNetConnection : NetConnection internal void InvokeDisconnected() { diff --git a/patches/VintagestoryLib/Vintagestory.Server/ServerMain.cs.patch b/patches/VintagestoryLib/Vintagestory.Server/ServerMain.cs.patch index 170027d..628503f 100644 --- a/patches/VintagestoryLib/Vintagestory.Server/ServerMain.cs.patch +++ b/patches/VintagestoryLib/Vintagestory.Server/ServerMain.cs.patch @@ -1122,7 +1122,7 @@ index 6f69cd3..a654aca 100644 public IPlayer PlayerByUid(string playerUid) { if (playerUid == null) -@@ -3414,44 +3895,63 @@ public sealed class ServerMain : GameMain, IServerWorldAccessor, IWorldAccessor, +@@ -3414,44 +3895,62 @@ public sealed class ServerMain : GameMain, IServerWorldAccessor, IWorldAccessor, { return; } @@ -1160,10 +1160,9 @@ index 6f69cd3..a654aca 100644 + if (requestId == Packet_ClientIdEnum.ServerQuery) + { + connection.SendPreparedPacket(StratumServerListQuery.Instance.PreparedPacket, false, Logger); -+ if (!StratumRuntime.Config.Performance.Network.SendQueueEnabled) -+ { -+ connection.Shutdown(); -+ } ++ // Stratum: always close the query socket. Shutdown flushes the answer ++ // and returns without waiting, so the packet thread is not stalled. ++ connection.Shutdown(); + TotalReceivedBytes += msg.messageLength; + return; + } diff --git a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs index cfeaf21..61f9d3b 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs @@ -632,12 +632,11 @@ public void EnsurePopulated() internal class StratumNetworkConfig { - // On by default. With the queue off, TcpNetConnection.Send calls Socket.SendAsync on - // the gameplay thread. A 1000-player capture on a 130k-chunk world spent about half of - // ServerMain.Process samples in SocketAsyncEventArgs.DoOperationSendSingleBuffer. - // The queue is the only sender for the connection and runs on the thread pool. Set - // this false to restore the direct send path. See #325. - public bool SendQueueEnabled { get; set; } = true; + // Off by default. The queue is the only sender for the connection and runs on the + // thread pool. The short-write loop, MaxPendingBytes disconnect, and shutdown flush + // are in place. The shipped default stays off until a 1000-bot run, a slow real + // client, and a play session are recorded on this code. Set this true to enable it. + public bool SendQueueEnabled { get; set; } = false; // Packets at or above this size skip coalescing and go out alone. Matches the TCP MTU // assumption the old flush buffer used. @@ -647,12 +646,12 @@ internal class StratumNetworkConfig // a single SendAsync call carries. public int CoalesceLimitBytes { get; set; } = 65536; - // Bytes StratumSendQueue may retain for one connection. Chunk scheduling already - // slows that client at OutboundPressurePendingBytesHardLimit (1 MiB) and still - // sends a minimum budget there, so this cap sits above that line. An enqueue that - // would grow a non-empty queue past the cap is refused and the connection is - // disconnected. One packet is still accepted when the queue is empty, so a single - // large send is not dropped on the floor. See #345. + // Bytes StratumSendQueue may retain for one connection, not counting one accepted + // packet that is already at or above this cap. That oversized packet is not added + // to the pending total, so the packets behind it are judged on their own size. + // A second packet that large, or an enqueue that would grow an already non-empty + // queue past the cap, disconnects the connection. Chunk scheduling already slows + // that client at OutboundPressurePendingBytesHardLimit (1 MiB). See #345. public int MaxPendingBytes { get; set; } = 8 * 1024 * 1024; public void EnsureSane() diff --git a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs index 67929f8..3201685 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs @@ -38,6 +38,7 @@ internal sealed class StratumSendQueue private int pendingCount; private long pendingBytes; + private int exemptInFlight; private int overflowDisconnectStarted; public int PendingCount => Volatile.Read(ref pendingCount); @@ -87,10 +88,11 @@ private void EnsureDrainStarted() }, "StratumSendQueueDrain"); } - // Finish bytes already queued, then return. Close and Shutdown call this before they - // cancel the socket, so a disconnect reason enqueued just above them still reaches - // the client. The wait is bounded: a wedged socket does not hold the caller forever. - public void FlushPending(int timeoutMs) + // Stop accepting packets and let the drain finish what is already queued. Returns + // immediately. shutdownSocket runs after the drain completes, or after timeoutMs, + // whichever comes first. The caller must not wait: DisconnectPlayer runs on the + // main thread, and a peer that never reads would otherwise stall that thread. + public void ScheduleShutdown(int timeoutMs, Action shutdownSocket) { Task task; lock (startGate) @@ -102,14 +104,28 @@ public void FlushPending(int timeoutMs) if (task == null) { + shutdownSocket(); return; } + _ = ObserveShutdown(task, Math.Max(0, timeoutMs), shutdownSocket); + } + + private static async Task ObserveShutdown(Task drain, int timeoutMs, Action shutdownSocket) + { + try + { + await Task.WhenAny(drain, Task.Delay(timeoutMs)).ConfigureAwait(false); + } + catch (Exception) + { + } + try { - task.Wait(Math.Max(0, timeoutMs)); + shutdownSocket(); } - catch (AggregateException) + catch (Exception) { } } @@ -124,17 +140,39 @@ public void Enqueue(byte[] dataWithLength) } int length = dataWithLength.Length; - long queued = Interlocked.Add(ref pendingBytes, length); - // Queue was already holding data and this packet would pass the cap. Roll it - // back and disconnect. Dropping the packet without closing would split the - // length-prefixed stream. Blocking here would put the gameplay thread back - // on the slow client. An empty queue still accepts one packet so a single - // large send is not refused. - if (queued > maxPendingBytes && queued - length > 0) + bool exempt = length >= maxPendingBytes; + if (exempt) { - Interlocked.Add(ref pendingBytes, -length); - DisconnectOverflow(); - return; + // One packet at or above the cap is accepted when the queue holds nothing + // else, and it is not added to pendingBytes. The packets behind it are + // judged on their own size. A second packet that large disconnects. + if (Interlocked.Read(ref pendingBytes) > 0) + { + DisconnectOverflow(); + return; + } + if (Interlocked.Increment(ref exemptInFlight) != 1) + { + Interlocked.Decrement(ref exemptInFlight); + DisconnectOverflow(); + return; + } + if (Interlocked.Read(ref pendingBytes) > 0) + { + Interlocked.Decrement(ref exemptInFlight); + DisconnectOverflow(); + return; + } + } + else + { + long queued = Interlocked.Add(ref pendingBytes, length); + if (queued > maxPendingBytes) + { + Interlocked.Add(ref pendingBytes, -length); + DisconnectOverflow(); + return; + } } Interlocked.Increment(ref pendingCount); @@ -144,8 +182,7 @@ public void Enqueue(byte[] dataWithLength) { // Writer already completed (connection closing). Drop, matches vanilla behavior // of a send attempted after Close()/Dispose(). - Interlocked.Decrement(ref pendingCount); - Interlocked.Add(ref pendingBytes, -length); + Release(length); return; } @@ -178,7 +215,11 @@ private void DisconnectOverflow() public void Complete() { - channel.Writer.TryComplete(); + lock (startGate) + { + closed = 1; + channel.Writer.TryComplete(); + } } private async Task DrainAsync() @@ -202,7 +243,7 @@ private async Task DrainAsync() coalescedLength = 0; } - Decrement(packet.Length); + Release(packet.Length); if (!await SendAsync(packet, packet.Length).ConfigureAwait(false)) return; continue; } @@ -214,7 +255,7 @@ private async Task DrainAsync() coalescedLength = 0; } - Decrement(packet.Length); + Release(packet.Length); Buffer.BlockCopy(packet, 0, coalesceBuffer, coalescedLength, packet.Length); coalescedLength += packet.Length; } @@ -231,9 +272,15 @@ private async Task DrainAsync() } } - private void Decrement(int bytes) + private void Release(int bytes) { Interlocked.Decrement(ref pendingCount); + if (bytes >= maxPendingBytes) + { + Interlocked.Decrement(ref exemptInFlight); + return; + } + Interlocked.Add(ref pendingBytes, -bytes); } diff --git a/tests/StratumScenarios/SendQueueDisconnectScenarios.cs b/tests/StratumScenarios/SendQueueDisconnectScenarios.cs index 2c15d52..416cf90 100644 --- a/tests/StratumScenarios/SendQueueDisconnectScenarios.cs +++ b/tests/StratumScenarios/SendQueueDisconnectScenarios.cs @@ -1,3 +1,5 @@ +using System; +using System.Diagnostics; using System.Net; using System.Net.Sockets; using System.Reflection; @@ -8,42 +10,124 @@ namespace StratumScenarios; /// -/// With the send queue on, DisconnectPlayer only enqueues the reason and CloseConnection -/// shuts the socket down immediately. The reason has to be on the wire before that FIN. +/// With the send queue on, DisconnectPlayer only enqueues the reason. Shutdown must +/// put that reason on the wire before the FIN, and it must return without waiting +/// on a peer that never reads. +/// The fixture forces the queue on. The shipped default is off. /// +[AtlasDataFiles("fixtures/stratum-sendqueue-on", TargetPath = "")] public class SendQueueDisconnectScenarios : AtlasScenarioBase { [AtlasScenario(TimeoutMs = 60_000)] public void Shutdown_Should_DeliverQueuedBytes_When_SendQueueEnabled() { - Type connectionType = Type.GetType("Vintagestory.Server.Network.TcpNetConnection, VintagestoryLib") + using ConnectionHarness harness = ConnectionHarness.Open(); + byte[] payload = Encoding.ASCII.GetBytes("kick-reason"); + harness.Send(payload); + harness.Shutdown(); + + byte[] buffer = new byte[64]; + int read = harness.Client.Client.Receive(buffer); + string text = Encoding.ASCII.GetString(buffer, 0, read); + Assert.Contains("kick-reason", text, StringComparison.Ordinal); + harness.Close(); + } + + [AtlasScenario(TimeoutMs = 60_000)] + public void Shutdown_Should_ReturnImmediately_When_PeerNeverReads() + { + using ConnectionHarness harness = ConnectionHarness.Open(); + harness.Accepted.SendBufferSize = 1024; + byte[] backlog = new byte[256 * 1024]; + harness.Send(backlog); + + Stopwatch watch = Stopwatch.StartNew(); + harness.Shutdown(); + watch.Stop(); + Assert.True(watch.ElapsedMilliseconds < 50, "Shutdown blocked for " + watch.ElapsedMilliseconds + " ms"); + harness.Close(); + } + + private sealed class ConnectionHarness : IDisposable + { + private static readonly Type ConnectionType = Type.GetType("Vintagestory.Server.Network.TcpNetConnection, VintagestoryLib") ?? throw new InvalidOperationException("TcpNetConnection was not in the loaded lib"); - using var listener = new TcpListener(IPAddress.Loopback, 0); - listener.Start(); - int port = ((IPEndPoint)listener.LocalEndpoint).Port; + private readonly TcpListener listener; + private readonly object connection; + private bool sendQueueForced; - using var client = new TcpClient(); - client.Connect(IPAddress.Loopback, port); - using Socket accepted = listener.AcceptSocket(); - client.Client.ReceiveTimeout = 2000; + public TcpClient Client { get; } - object connection = Activator.CreateInstance(connectionType, accepted)!; - connectionType.GetMethod("StartReceiving")!.Invoke(connection, null); - object queue = connectionType.GetField("stratumSendQueue", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(connection) - ?? throw new InvalidOperationException("SendQueueEnabled did not attach a queue"); + public Socket Accepted { get; } - byte[] payload = Encoding.ASCII.GetBytes("kick-reason"); - MethodInfo send = connectionType.GetMethod("Send", new[] { typeof(byte[]), typeof(bool) })!; - send.Invoke(connection, new object[] { payload, false }); - connectionType.GetMethod("Shutdown")!.Invoke(connection, null); + private ConnectionHarness(TcpListener listener, TcpClient client, Socket accepted, object connection) + { + this.listener = listener; + Client = client; + Accepted = accepted; + this.connection = connection; + } - byte[] buffer = new byte[64]; - int read = client.Client.Receive(buffer); - string text = Encoding.ASCII.GetString(buffer, 0, read); - Assert.Contains("kick-reason", text, StringComparison.Ordinal); + public static ConnectionHarness Open() + { + SetSendQueueEnabled(true); + TcpListener listener = new TcpListener(IPAddress.Loopback, 0); + listener.Start(); + int port = ((IPEndPoint)listener.LocalEndpoint).Port; + TcpClient client = new TcpClient(); + client.Connect(IPAddress.Loopback, port); + Socket accepted = listener.AcceptSocket(); + client.Client.ReceiveTimeout = 2000; + object connection = Activator.CreateInstance(ConnectionType, accepted)!; + ConnectionType.GetMethod("StartReceiving")!.Invoke(connection, null); + object queue = ConnectionType.GetField("stratumSendQueue", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(connection) + ?? throw new InvalidOperationException("SendQueueEnabled did not attach a queue"); + GC.KeepAlive(queue); + return new ConnectionHarness(listener, client, accepted, connection) { sendQueueForced = true }; + } + + public void Send(byte[] payload) + { + MethodInfo send = ConnectionType.GetMethod("Send", new[] { typeof(byte[]), typeof(bool) })!; + send.Invoke(connection, new object[] { payload, false }); + } + + public void Shutdown() + { + ConnectionType.GetMethod("Shutdown")!.Invoke(connection, null); + } + + public void Close() + { + ConnectionType.GetMethod("Close")!.Invoke(connection, null); + } + + public void Dispose() + { + try + { + Close(); + } + catch (TargetInvocationException) + { + } + Client.Dispose(); + listener.Stop(); + if (sendQueueForced) + { + SetSendQueueEnabled(false); + } + } - connectionType.GetMethod("Close")!.Invoke(connection, null); - GC.KeepAlive(queue); + private static void SetSendQueueEnabled(bool enabled) + { + Type runtime = Type.GetType("Vintagestory.Server.StratumRuntime, VintagestoryLib") + ?? throw new InvalidOperationException("StratumRuntime was not in the loaded lib"); + object config = runtime.GetProperty("Config")!.GetValue(null)!; + object performance = config.GetType().GetProperty("Performance")!.GetValue(config)!; + object network = performance.GetType().GetProperty("Network")!.GetValue(performance)!; + network.GetType().GetProperty("SendQueueEnabled")!.SetValue(network, enabled); + } } } diff --git a/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum-performance.json b/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum-performance.json new file mode 100644 index 0000000..f825eea --- /dev/null +++ b/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum-performance.json @@ -0,0 +1,5 @@ +{ + "Network": { + "SendQueueEnabled": true + } +} diff --git a/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum.json b/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum.json new file mode 100644 index 0000000..dab506f --- /dev/null +++ b/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum.json @@ -0,0 +1,3 @@ +{ + "ConfigVersion": 4 +}