diff --git a/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch b/patches/VintagestoryLib/Vintagestory.Server.Network/TcpNetConnection.cs.patch index 6f09bcea..6ff211dd 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,42 @@ index ae86835..34c0f3d 100644 catch { InvokeDisconnected(); -@@ -281,10 +313,11 @@ public class TcpNetConnection : NetConnection +@@ -269,4 +301,21 @@ public class TcpNetConnection : NetConnection + public override void Shutdown() + { ++ // 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 +329,12 @@ public class TcpNetConnection : NetConnection } } - + public override void Close() { ++ // Stratum: do not wait. Cancel is the hard bound. Shutdown already scheduled the FIN. + stratumSendQueue?.Complete(); try { cts?.Cancel(); } catch -@@ -302,10 +335,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 170027db..628503ff 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 e9d36c70..61f9d3bd 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs @@ -632,10 +632,10 @@ 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. + // 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 @@ -646,10 +646,19 @@ internal class StratumNetworkConfig // a single SendAsync call carries. public int CoalesceLimitBytes { get; set; } = 65536; + // 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() { 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 0be92a76..32016859 100644 --- a/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs +++ b/sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs @@ -30,9 +30,16 @@ internal sealed class StratumSendQueue private readonly CancellationToken cancellationToken; private readonly int largeThreshold; 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 exemptInFlight; + private int overflowDisconnectStarted; public int PendingCount => Volatile.Read(ref pendingCount); @@ -46,6 +53,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, @@ -54,41 +62,176 @@ 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"); + } + + // 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) + { + closed = 1; + channel.Writer.TryComplete(); + task = drainTask; + } + + 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 + { + shutdownSocket(); + } + catch (Exception) + { + } } // dataWithLength must already carry the 4-byte length prefix. Callers hand off a buffer // 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; + bool exempt = length >= maxPendingBytes; + if (exempt) + { + // 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); - Interlocked.Add(ref pendingBytes, dataWithLength.Length); - 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, -dataWithLength.Length); + if (closed != 0 || !channel.Writer.TryWrite(dataWithLength)) + { + // Writer already completed (connection closing). Drop, matches vanilla behavior + // of a send attempted after Close()/Dispose(). + Release(length); + return; + } + + EnsureDrainStarted(); } } + private void DisconnectOverflow() + { + if (Interlocked.Exchange(ref overflowDisconnectStarted, 1) != 0) + { + return; + } + + 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(); + } + public void Complete() { - channel.Writer.TryComplete(); + lock (startGate) + { + closed = 1; + channel.Writer.TryComplete(); + } } 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) { @@ -100,18 +243,19 @@ private async Task DrainAsync() coalescedLength = 0; } - Decrement(packet.Length); + Release(packet.Length); if (!await SendAsync(packet, packet.Length).ConfigureAwait(false)) return; 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; } - Decrement(packet.Length); + Release(packet.Length); Buffer.BlockCopy(packet, 0, coalesceBuffer, coalescedLength, packet.Length); coalescedLength += packet.Length; } @@ -128,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); } @@ -138,7 +288,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) diff --git a/tests/StratumScenarios/SendQueueDisconnectScenarios.cs b/tests/StratumScenarios/SendQueueDisconnectScenarios.cs new file mode 100644 index 00000000..416cf90d --- /dev/null +++ b/tests/StratumScenarios/SendQueueDisconnectScenarios.cs @@ -0,0 +1,133 @@ +using System; +using System.Diagnostics; +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. 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() + { + 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"); + + private readonly TcpListener listener; + private readonly object connection; + private bool sendQueueForced; + + public TcpClient Client { get; } + + public Socket Accepted { get; } + + private ConnectionHarness(TcpListener listener, TcpClient client, Socket accepted, object connection) + { + this.listener = listener; + Client = client; + Accepted = accepted; + this.connection = connection; + } + + 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); + } + } + + 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 00000000..f825eeaf --- /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 00000000..dab506ff --- /dev/null +++ b/tests/StratumScenarios/fixtures/stratum-sendqueue-on/stratum.json @@ -0,0 +1,3 @@ +{ + "ConfigVersion": 4 +}