diff --git a/GenOnlineService/Constants.cs b/GenOnlineService/Constants.cs index a860c70..1dee6bf 100644 --- a/GenOnlineService/Constants.cs +++ b/GenOnlineService/Constants.cs @@ -30,6 +30,7 @@ using System.Text; using System.Text.Json; using System.Text.Json.Serialization; +using System.Threading.Channels; using System.Threading.Tasks; namespace GenOnlineService @@ -479,16 +480,14 @@ public static async Task CreateSession(AppDbContext _db, } // now create a websocket, we always do this whether its reconnect or not, only data is persistent - // NOTE: on the reconnect path the previous websocket may still be physically open (e.g. its receive loop - // is parked in a 30s ReceiveAsync). Overwriting the entry without closing it leaks a zombie connection - // that keeps running and sending on a socket nobody owns any more. + // close the previous socket first: it may still be open, and its send loop must stop reading the session's queue if (m_dictWebsockets[sessionType].TryRemove(ownerID, out UserWebSocketInstance? supersededSess) && supersededSess != null) { Console.WriteLine("Closing superseded websocket for {0} ({1})", ownerID, strDisplayName); await supersededSess.CloseAsync(WebSocketCloseStatus.NormalClosure, "Superseded by a newer connection"); } - UserWebSocketInstance newSess = new UserWebSocketInstance(sessionType, ownerID); + UserWebSocketInstance newSess = new UserWebSocketInstance(sessionType, ownerID, userCacheData); m_dictWebsockets[sessionType][ownerID] = newSess; // update last login and last ip @@ -526,24 +525,17 @@ public static async Task CreateSession(AppDbContext _db, outboundMsg.num_online = numOnline; outboundMsg.num_pending = numPending; byte[] bytesJSON = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(outboundMsg)); - await newSess.SendAsync(bytesJSON, WebSocketMessageType.Text); + // queued: the socket isn't attached until the handshake completes + userCacheData.QueueWebsocketSend(bytesJSON); } } return newSess; } - public static async Task Tick() + public static void Tick() { FlushLobbyListUpdates(); - - // Give the entire tick a 20 ms deadline. All users drain concurrently via - // Task.WhenAll, so a slow/stuck client cannot delay others. If the deadline - // fires, the CancellationToken propagates into each in-flight SendAsync and - // into the dequeue loop guard, so the stuck user is skipped and their unsent - // messages stay in the queue for the next tick. - using var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(20)); - await Task.WhenAll(m_dictUserSessions.Values.SelectMany(inner => inner.Values).Select(sess => sess.TickWebsocket(cts.Token))); } public static int GetNumberOfUsersOnline() @@ -559,33 +551,14 @@ public static int GetNumberOfUsersOnline() public static async Task CheckForTimeouts() { - List lstSessionsToDestroy = new(); foreach (var sessionDataByClient in m_dictWebsockets) { foreach (var sessionData in sessionDataByClient.Value) { -#if DEBUG - const int timeoutVal = 60000 * 10; -#else - const int timeoutVal = 20000; -#endif - if (sessionData.Value.GetTimeSinceLastPing() >= timeoutVal) - { - lstSessionsToDestroy.Add(sessionData.Value); - } - else - { - await sessionData.Value.SendPong(); - } + sessionData.Value.QueueKeepAlive(); } } - foreach (UserWebSocketInstance wsSess in lstSessionsToDestroy) - { - Console.WriteLine("Timing out WS session for {0}", wsSess.m_UserID); - await DeleteSession(wsSess.m_UserID, wsSess.m_SessionType, wsSess, false); - } - // do we need to clear out cache entries? List> lstCacheEntriesToDestroy = new(); foreach (var sessionDataPerClientType in m_dictUserSessions) @@ -850,7 +823,9 @@ public static async Task DisconnectUser(Int64 userID, byte[] finalMessage) UserWebSocketInstance? oldWS = GetWebSocketForSession(userSession); if (oldWS != null) { - await oldWS.SendAsync(finalMessage, WebSocketMessageType.Text); + // short bound: the socket is torn down right after, don't hold up the caller on a dead peer + userSession.QueueWebsocketSend(finalMessage); + await userSession.FlushWebsocketSendsAsync(TimeSpan.FromSeconds(2)); } await DeleteSession(userID, userSession.GetSessionType(), oldWS, true); @@ -1275,9 +1250,23 @@ public void QueueWebsocketSend(byte[] bytesJSON) return; } - // Always enqueue; the TickWebsocket drain loop is the sole sender, - // ensuring WebSocket.SendAsync is never called concurrently. - m_lstPendingWebsocketSends.Enqueue(bytesJSON); + if (!m_outbound.Writer.TryWrite(bytesJSON)) + { + // too far behind to catch up: drop the connection so the client reconnects + if (WebSocketManager.GetWebSocketForSession(this)?.AbortIfOpen() == true) + { + Console.WriteLine("[WebSocket] Outbound queue full for {0}, connection aborted", m_UserID); + } + } + } + + public async Task FlushWebsocketSendsAsync(TimeSpan timeout) + { + long deadline = Environment.TickCount64 + (long)timeout.TotalMilliseconds; + while (m_outbound.Reader.Count > 0 && Environment.TickCount64 < deadline) + { + await Task.Delay(20); + } } public async Task CloseWebsocket(WebSocketCloseStatus reason, string strReason) @@ -1291,25 +1280,15 @@ public async Task CloseWebsocket(WebSocketCloseStatus rea return websocketForUser; } - public async Task TickWebsocket(CancellationToken tickToken = default) + // outlives the socket so queued messages reach the client after a reconnect; the attached socket's send loop is the only reader + private const int c_MaxQueuedSends = 1024; + private readonly Channel m_outbound = Channel.CreateBounded(new BoundedChannelOptions(c_MaxQueuedSends) { - // Do we have a connection to send on? - UserWebSocketInstance websocketForUser = WebSocketManager.GetWebSocketForSession(this); - if (websocketForUser != null) - { - const int maxMessagesSendPerFrame = 50; - int messagesSent = 0; - // start dequeing and sending - while (!tickToken.IsCancellationRequested && messagesSent < maxMessagesSendPerFrame && m_lstPendingWebsocketSends.TryDequeue(out byte[] packetData)) - { - await websocketForUser.SendAsync(packetData, WebSocketMessageType.Text, tickToken); - ++messagesSent; - } - } - } - - // TODO_CACHE: Size limit this? - ConcurrentQueue m_lstPendingWebsocketSends = new ConcurrentQueue(); + SingleReader = true, + FullMode = BoundedChannelFullMode.Wait + }); + + public ChannelReader OutboundWebsocketSends => m_outbound.Reader; public bool NeedsCleanup() { @@ -1464,141 +1443,140 @@ public class UserWebSocketInstance public EUserSessionType m_SessionType = EUserSessionType.None; public Int64 m_UserID = -1; - public Int64 m_lastPingTime = Environment.TickCount64; // last time we pinged this user, used to detect disconnects - - + // pinged after c_KeepAliveInterval of silence, aborted if no pong within c_KeepAliveTimeout + public static readonly TimeSpan c_KeepAliveInterval = TimeSpan.FromSeconds(15); +#if DEBUG + public static readonly TimeSpan c_KeepAliveTimeout = TimeSpan.FromMinutes(10); // survive debugger breaks +#else + public static readonly TimeSpan c_KeepAliveTimeout = TimeSpan.FromSeconds(45); +#endif + + // cancelling a pending send aborts the socket, so only give up on a peer that stopped reading + private static readonly TimeSpan c_SendStallTimeout = c_KeepAliveInterval + c_KeepAliveTimeout; // TODO: Start using nullable for int values etc instead of doing 0 or -1 - public async Task SendPong() + // reply to legacy JSON PING; released clients only reset their timeout on it + public void QueuePong() { - OnPing(); - - // send pong back WebSocketMessage_PONG outboundMsg = new WebSocketMessage_PONG(); outboundMsg.msg_id = (int)EWebSocketMessageID.PONG; - byte[] bytesJSON = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(outboundMsg)); - await SendAsync(bytesJSON, WebSocketMessageType.Text); + m_OwnerSession.QueueWebsocketSend(Encoding.UTF8.GetBytes(JsonSerializer.Serialize(outboundMsg))); + } + + // legacy JSON keep-alive for clients that don't send their own PING; any queued message serves the same purpose + public void QueueKeepAlive() + { + if (m_OwnerSession.OutboundWebsocketSends.Count == 0) + { + QueuePong(); + } } - private WebSocket? m_SockInternal = null; + private readonly UserSession m_OwnerSession; + private readonly CancellationTokenSource m_sendLoopStop = new CancellationTokenSource(); + private Task m_sendLoop = Task.CompletedTask; - public UserWebSocketInstance(EUserSessionType sessionType, Int64 ownerID) : base() + public UserWebSocketInstance(EUserSessionType sessionType, Int64 ownerID, UserSession ownerSession) : base() { m_SessionType = sessionType; m_UserID = ownerID; + m_OwnerSession = ownerSession; } public void AttachWebsocket(WebSocket sock) { m_SockInternal = sock; + m_sendLoop = Task.Run(RunSendLoop); } - public void OnPing() - { - m_lastPingTime = Environment.TickCount64; - } - - public Int64 GetLastPingTime() + public bool AbortIfOpen() { - return m_lastPingTime; - } + WebSocket? sock = m_SockInternal; + if (sock == null || sock.State != WebSocketState.Open) + { + return false; + } - public Int64 GetTimeSinceLastPing() - { - return Environment.TickCount64 - m_lastPingTime; + sock.Abort(); + return true; } - public async Task SendAsync(byte[] buffer, WebSocketMessageType messageType, CancellationToken externalToken = default) + // sole writer to the socket; a message leaves the queue only once sent, so a dead socket leaves it for the next connection + private async Task RunSendLoop() { - if (m_SockInternal != null) + ChannelReader outbound = m_OwnerSession.OutboundWebsocketSends; + try { - // WebSocket.SendAsync must never be called concurrently on the same socket or the frame stream gets - // corrupted. Several paths (tick drain, pongs, direct sends) can send at the same time, so serialize here. - await m_SendLock.WaitAsync(); - try + while (await outbound.WaitToReadAsync(m_sendLoopStop.Token)) { - // should we chunked send? - /* - const int frameMax = 99999999; - if (buffer.Length > frameMax) + while (!m_sendLoopStop.IsCancellationRequested && outbound.TryPeek(out byte[]? buffer)) { - int bytresRemaining = buffer.Length; - int numFrames = (int)Math.Ceiling((float)buffer.Length / (float)frameMax); - - System.Diagnostics.Debug.WriteLine("[Websocket] sending {0} bytes in {1} chunks", bytresRemaining, numFrames); - - for (int i = 0; i < numFrames; ++i) + // a failed send on a still-open socket is dropped rather than retried forever + if (!await SendFrameAsync(buffer) && m_SockInternal?.State != WebSocketState.Open) { - int bytesToSend = Math.Min(bytresRemaining, frameMax); - bool bLastFrame = i == numFrames - 1; - - - ArraySegment arrSegment = new ArraySegment(buffer, i * frameMax, bytesToSend); - System.Diagnostics.Debug.WriteLine("[Websocket] send frame {0} with {1} bytes (last: {2})", i, bytesToSend, bLastFrame); - await m_SockInternal.SendAsync(arrSegment, messageType, bLastFrame, CancellationToken.None); - - bytresRemaining -= bytesToSend; + return; } + outbound.TryRead(out _); } - else // just send whole - { - await m_SockInternal.SendAsync(buffer, messageType, true, CancellationToken.None); - } - */ - - CancellationTokenSource cts = CancellationTokenSource.CreateLinkedTokenSource(externalToken); - try - { - cts.CancelAfter(TimeSpan.FromMilliseconds(500)); - await m_SockInternal.SendAsync(buffer, messageType, true, cts.Token); - } - finally - { - try - { - // cts is intentionally not disposed: disposing it races with the parent token's - // timer callback (ThreadPool thread), causing ObjectDisposedException. GC reclaims it. - - - - } - catch (ObjectDisposedException) - { - - - } - } - } - catch - { - - } - finally - { - m_SendLock.Release(); } } + catch (OperationCanceledException) + { + } + catch (Exception ex) + { + Console.WriteLine($"[WebSocket] Send loop failed for {m_UserID}: {ex}"); + } } - private readonly SemaphoreSlim m_SendLock = new SemaphoreSlim(1, 1); + private async Task SendFrameAsync(byte[] buffer) + { + WebSocket? sock = m_SockInternal; + if (sock == null || sock.State != WebSocketState.Open) + { + return false; + } + + try + { + using var cts = new CancellationTokenSource(c_SendStallTimeout); + await sock.SendAsync(buffer, WebSocketMessageType.Text, true, cts.Token); + return true; + } + catch (OperationCanceledException) + { + Console.WriteLine("[WebSocket] Send to {0} stalled for {1}s, connection aborted", m_UserID, c_SendStallTimeout.TotalSeconds); + return false; + } + catch + { + return false; + } + } public async Task CloseAsync(WebSocketCloseStatus closeStatus, string? statusDescription) { + // stop reading the queue first so the next connection's send loop is the only reader + m_sendLoopStop.Cancel(); + if (m_SockInternal != null) { try { // dont wait forever, certain situations can cause that in ASP.NET - var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); + using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); await m_SockInternal.CloseAsync(closeStatus, statusDescription, cts.Token); } catch { - + // a timed-out close leaves the socket open behind a stuck send; abort to end both + m_SockInternal.Abort(); } } + + await m_sendLoop; } } diff --git a/GenOnlineService/Controllers/WebSocket/WebSocketController.cs b/GenOnlineService/Controllers/WebSocket/WebSocketController.cs index 2f30ed1..7abfa74 100644 --- a/GenOnlineService/Controllers/WebSocket/WebSocketController.cs +++ b/GenOnlineService/Controllers/WebSocket/WebSocketController.cs @@ -184,8 +184,8 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) string ipAddress = IPHelpers.NormalizeIP(HttpContext.Connection.RemoteIpAddress?.ToString()); string ipContinent = "NA"; string ipCountry = "US"; - double dLongitude = 38.8977; // the whitehouse; - double dLatitude = 77.0365f; // the whitehouse; + double dLongitude = -77.0365; // the whitehouse + double dLatitude = 38.8977; // the whitehouse try { @@ -255,7 +255,19 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) } // accept WS - using var webSocket = await HttpContext.WebSockets.AcceptWebSocketAsync(); + WebSocket acceptedSocket; + try + { + acceptedSocket = await HttpContext.WebSockets.AcceptWebSocketAsync(); + } + catch + { + // handshake failed: release the session registered for this socket, or it stays online with no connection + await WebSocketManager.DeleteSession(user_id, wsSess.m_SessionType, wsSess, false); + throw; + } + + using var webSocket = acceptedSocket; // attach wsSess.AttachWebsocket(webSocket); @@ -279,15 +291,18 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) try { - using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); // timeout + // cancelling a pending receive aborts the socket receiveResult = await webSocket.ReceiveAsync( - new ArraySegment(buffer), cts.Token); + new ArraySegment(buffer), HttpContext.RequestAborted); } catch (OperationCanceledException) { - // No message received in 30s � send a keep-alive pong and continue waiting - await wsSess.SendPong(); - continue; + break; + } + catch (WebSocketException ex) when (ex.WebSocketErrorCode == WebSocketError.ConnectionClosedPrematurely) + { + // client dropped without a close handshake + break; } catch (Exception ex) { @@ -299,8 +314,7 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) if (receiveResult.MessageType == WebSocketMessageType.Close) { - using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); // timeout - await webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing", cts.Token); + await wsSess.CloseAsync(WebSocketCloseStatus.NormalClosure, "Closing"); break; } @@ -316,8 +330,7 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) fragmentBuffer.Dispose(); fragmentBuffer = null; - using var ctsTooBig = new CancellationTokenSource(TimeSpan.FromSeconds(10)); - await webSocket.CloseAsync(WebSocketCloseStatus.MessageTooBig, "Message too large", ctsTooBig.Token); + await wsSess.CloseAsync(WebSocketCloseStatus.MessageTooBig, "Message too large"); break; } @@ -343,7 +356,7 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) // if we lost session data, close WS if (sourceUserData == null) { - wsSess.CloseAsync(WebSocketCloseStatus.NormalClosure, "User signed in from another point of presence [B]"); + await wsSess.CloseAsync(WebSocketCloseStatus.NormalClosure, "User signed in from another point of presence [B]"); break; } @@ -375,8 +388,7 @@ public async Task Get([FromHeader(Name = "is-reconnect")] bool bIsReconnect) } } - using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); // timeout - await webSocket.CloseAsync(closeStatus, closeStatusDescription, cts.Token); + await wsSess.CloseAsync(closeStatus, closeStatusDescription); } } @@ -491,7 +503,7 @@ private async Task ProcessWSMessage(UserWebSocketInstance sourceWS, UserSession { if (msgID == EWebSocketMessageID.PING) { - await sourceWS.SendPong(); + sourceWS.QueuePong(); } else if (msgID == EWebSocketMessageID.SOCIAL_SUBSCRIBE_REALTIME_UPDATES) { @@ -562,7 +574,7 @@ private async Task ProcessWSMessage(UserWebSocketInstance sourceWS, UserSession // send to source byte[] bytesJSON = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(outboundMsg)); - await sourceWS.SendAsync(bytesJSON, WebSocketMessageType.Text); + sourceUserSession.QueueWebsocketSend(bytesJSON); } } } diff --git a/GenOnlineService/Program.cs b/GenOnlineService/Program.cs index f872451..7de732d 100644 --- a/GenOnlineService/Program.cs +++ b/GenOnlineService/Program.cs @@ -1551,7 +1551,8 @@ public static async Task Main(string[] args) var webSocketOptions = new WebSocketOptions { - KeepAliveInterval = TimeSpan.FromSeconds(30) + KeepAliveInterval = UserWebSocketInstance.c_KeepAliveInterval, + KeepAliveTimeout = UserWebSocketInstance.c_KeepAliveTimeout }; app.UseWebSockets(webSocketOptions); @@ -1653,7 +1654,7 @@ public static async Task Main(string[] args) { var lobbyManager = ServiceLocator.Services.GetRequiredService(); await lobbyManager.Tick(); - await WebSocketManager.Tick(); + WebSocketManager.Tick(); } catch (Exception ex) {