Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<Task>)ReceiveData, "TcpNetConReceiveData");
}
Expand Down Expand Up @@ -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()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
+ }
Expand Down
17 changes: 13 additions & 4 deletions sources/VintagestoryLib/Vintagestory.Server/StratumConfig.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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);
}
}

Expand Down
198 changes: 182 additions & 16 deletions sources/VintagestoryLib/Vintagestory.Server/StratumSendQueue.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand All @@ -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<byte[]>(new UnboundedChannelOptions
{
SingleReader = true,
Expand All @@ -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<Task>)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<byte[]> 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)
{
Expand All @@ -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;
}
Expand All @@ -128,17 +272,39 @@ 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);
}

private async Task<bool> SendAsync(byte[] buffer, int length)
{
try
{
await socket.SendAsync(new ReadOnlyMemory<byte>(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<byte>(buffer, offset, length - offset), SocketFlags.None, cancellationToken).ConfigureAwait(false);
if (sent <= 0)
{
connection.InvokeDisconnected();
return false;
}

offset += sent;
}

return true;
}
catch (OperationCanceledException)
Expand Down
Loading