From faf81f906ddf99dd0c480bed72770ab775825276 Mon Sep 17 00:00:00 2001 From: Jordan Brown Date: Tue, 15 Dec 2020 13:18:51 -0500 Subject: [PATCH] Calls into IChatManager now queue messages - Better timeout handling and hopefully kills the last spurious test failure --- .../Components/Chat/ChatManager.cs | 82 ++++++++++++++--- .../Components/Chat/IChatManager.cs | 28 +++--- .../Components/Deployment/DreamMaker.cs | 33 ++++--- .../Components/Session/SessionController.cs | 91 ++++++++++--------- .../Components/Watchdog/BasicWatchdog.cs | 8 +- .../Components/Watchdog/WatchdogBase.cs | 49 +++++----- tools/ReleaseNotes/Program.cs | 1 - 7 files changed, 174 insertions(+), 118 deletions(-) diff --git a/src/Tgstation.Server.Host/Components/Chat/ChatManager.cs b/src/Tgstation.Server.Host/Components/Chat/ChatManager.cs index 25807b4d2f..93e8971229 100644 --- a/src/Tgstation.Server.Host/Components/Chat/ChatManager.cs +++ b/src/Tgstation.Server.Host/Components/Chat/ChatManager.cs @@ -111,6 +111,11 @@ namespace Tgstation.Server.Host.Components.Chat /// Task initialProviderConnectionsTask; + /// + /// A that represents all sent messages. + /// + Task messageSendTask; + /// /// The that completes when s change /// @@ -159,6 +164,8 @@ namespace Tgstation.Server.Host.Components.Chat trackingContexts = new List(); handlerCts = new CancellationTokenSource(); connectionsUpdated = new TaskCompletionSource(); + + messageSendTask = Task.CompletedTask; channelIdCounter = 1; } @@ -170,6 +177,8 @@ namespace Tgstation.Server.Host.Components.Chat handlerCts.Dispose(); foreach (var I in providers) await I.Value.DisposeAsync().ConfigureAwait(false); + + await messageSendTask.ConfigureAwait(false); } /// @@ -614,16 +623,29 @@ namespace Tgstation.Server.Host.Components.Chat } /// - public async Task SendMessage(string message, IEnumerable channelIds, CancellationToken cancellationToken) + public void QueueMessage(string message, IEnumerable channelIds) { if (message == null) throw new ArgumentNullException(nameof(message)); if (channelIds == null) throw new ArgumentNullException(nameof(channelIds)); + var task = SendMessage(message, channelIds, handlerCts.Token); + AddMessageTask(task); + } + + /// + /// Asynchronously send a given to a set of . + /// + /// The message to send. + /// The s of the s to send to. + /// The for the operation. + /// A representing the running operation. + Task SendMessage(string message, IEnumerable channelIds, CancellationToken cancellationToken) + { logger.LogTrace("Chat send \"{0}\" to channels: {1}", message, String.Join(", ", channelIds)); - await Task.WhenAll( + return Task.WhenAll( channelIds.Select(x => { ChannelMapping channelMapping; @@ -635,12 +657,11 @@ namespace Tgstation.Server.Host.Components.Chat if (!providers.TryGetValue(channelMapping.ProviderId, out provider)) return Task.CompletedTask; return provider.SendMessage(channelMapping.ProviderChannelId, message, cancellationToken); - })) - .ConfigureAwait(false); + })); } /// - public async Task SendWatchdogMessage(string message, CancellationToken cancellationToken) + public async Task QueueWatchdogMessage(string message, CancellationToken cancellationToken) { List wdChannels = null; message = String.Format(CultureInfo.InvariantCulture, "WD: {0}", message); @@ -654,18 +675,17 @@ namespace Tgstation.Server.Host.Components.Chat lock (mappedChannels) wdChannels = mappedChannels.Where(x => x.Value.IsWatchdogChannel).Select(x => x.Key).ToList(); - await SendMessage(message, wdChannels, cancellationToken).ConfigureAwait(false); + QueueMessage(message, wdChannels); } /// - public async Task> SendDeploymentMessage( + public Action QueueDeploymentMessage( Models.RevisionInformation revisionInformation, Version byondVersion, DateTimeOffset? estimatedCompletionTime, string gitHubOwner, string gitHubRepo, - bool localCommitPushed, - CancellationToken cancellationToken) + bool localCommitPushed) { List wdChannels; lock (mappedChannels) // so it doesn't change while we're using it @@ -675,7 +695,7 @@ namespace Tgstation.Server.Host.Components.Chat var callbacks = new List>(); - await Task.WhenAll( + var task = Task.WhenAll( wdChannels.Select( async x => { @@ -697,7 +717,7 @@ namespace Tgstation.Server.Host.Components.Chat gitHubRepo, channelMapping.ProviderChannelId, localCommitPushed, - cancellationToken) + handlerCts.Token) .ConfigureAwait(false); callbacks.Add(callback); @@ -709,10 +729,16 @@ namespace Tgstation.Server.Host.Components.Chat "Error sending deploy message to provider {0}!", channelMapping.ProviderId); } - })) - .ConfigureAwait(false); + })); - return (errorMessage, dreamMakerOutput) => Task.WhenAll(callbacks.Select(x => x(errorMessage, dreamMakerOutput))); + AddMessageTask(task); + + return (errorMessage, dreamMakerOutput) => AddMessageTask( + Task.WhenAll( + callbacks.Select( + x => x( + errorMessage, + dreamMakerOutput)))); } /// @@ -744,6 +770,7 @@ namespace Tgstation.Server.Host.Components.Chat if (chatHandler != null) await chatHandler.ConfigureAwait(false); await Task.WhenAll(providers.Select(x => x.Value).Select(x => x.Disconnect(cancellationToken))).ConfigureAwait(false); + await messageSendTask.ConfigureAwait(false); } /// @@ -800,7 +827,32 @@ namespace Tgstation.Server.Host.Components.Chat List wdChannels; lock (mappedChannels) // so it doesn't change while we're using it wdChannels = mappedChannels.Select(x => x.Key).ToList(); - return SendMessage(message, wdChannels, cancellationToken); + + QueueMessage(message, wdChannels); + return Task.CompletedTask; + } + + /// + /// Adds a given to . + /// + /// The to add. + void AddMessageTask(Task task) + { + async Task Wrap(Task originalTask) + { + await originalTask.ConfigureAwait(false); + try + { + await task.ConfigureAwait(false); + } + catch (Exception ex) + { + logger.LogWarning(ex, "Error in asynchronous chat message!"); + } + } + + lock (handlerCts) + messageSendTask = Wrap(messageSendTask); } } } diff --git a/src/Tgstation.Server.Host/Components/Chat/IChatManager.cs b/src/Tgstation.Server.Host/Components/Chat/IChatManager.cs index 3cfe424951..367274d8b2 100644 --- a/src/Tgstation.Server.Host/Components/Chat/IChatManager.cs +++ b/src/Tgstation.Server.Host/Components/Chat/IChatManager.cs @@ -44,21 +44,19 @@ namespace Tgstation.Server.Host.Components.Chat Task ChangeChannels(long connectionId, IEnumerable newChannels, CancellationToken cancellationToken); /// - /// Send a chat to a given set of + /// Queue a chat to a given set of . /// - /// The message being sent - /// The s of the s to send to - /// The for the operation - /// A representing the running operation - Task SendMessage(string message, IEnumerable channelIds, CancellationToken cancellationToken); + /// The message being sent. + /// The s of the s to send to. + void QueueMessage(string message, IEnumerable channelIds); /// - /// Send a chat to configured watchdog channels + /// Queue a chat to configured watchdog channels. /// - /// The message being sent - /// The for the operation - /// A representing the running operation - Task SendWatchdogMessage(string message, CancellationToken cancellationToken); + /// The message being sent. + /// The for the operation. + /// A representing the running operation. + Task QueueWatchdogMessage(string message, CancellationToken cancellationToken); /// /// Send the message for a deployment to configured deployment channels. @@ -69,16 +67,14 @@ namespace Tgstation.Server.Host.Components.Chat /// The repository GitHub owner, if any. /// The repository GitHub name, if any. /// if the local deployment commit was pushed to the remote repository. - /// The for the operation. - /// A resulting in a to call to update the message at the deployment's conclusion. Parameters: Error message if any, DreamMaker output if any. - Task> SendDeploymentMessage( + /// An to call to update the message at the deployment's conclusion. Parameters: Error message if any, DreamMaker output if any. + Action QueueDeploymentMessage( Models.RevisionInformation revisionInformation, Version byondVersion, DateTimeOffset? estimatedCompletionTime, string gitHubOwner, string gitHubRepo, - bool localCommitPushed, - CancellationToken cancellationToken); + bool localCommitPushed); /// /// Start tracking s and s. diff --git a/src/Tgstation.Server.Host/Components/Deployment/DreamMaker.cs b/src/Tgstation.Server.Host/Components/Deployment/DreamMaker.cs index a7fb801b9e..c031f55a4d 100644 --- a/src/Tgstation.Server.Host/Components/Deployment/DreamMaker.cs +++ b/src/Tgstation.Server.Host/Components/Deployment/DreamMaker.cs @@ -101,8 +101,14 @@ namespace Tgstation.Server.Host.Components.Deployment /// readonly object deploymentLock; - Func currentChatCallback; + /// + /// The active callback from . + /// + Action currentChatCallback; + /// + /// Cached for . + /// string currentDreamMakerOutput; /// @@ -716,25 +722,26 @@ namespace Tgstation.Server.Host.Components.Deployment var eventTask = eventConsumer.HandleEvent(EventType.DeploymentComplete, null, cancellationToken); - var chatTask = currentChatCallback(null, compileJob.Output); - currentChatCallback = null; - try { - await Task.WhenAll(commentsTask, eventTask, chatTask).ConfigureAwait(false); + currentChatCallback(null, compileJob.Output); + + await Task.WhenAll(commentsTask, eventTask).ConfigureAwait(false); } catch (Exception ex) { throw new JobException(ErrorCode.PostDeployFailure, ex); } + finally + { + currentChatCallback = null; + } } catch (Exception ex) { - if (currentChatCallback != null) - await currentChatCallback( - FormatExceptionForUsers(ex), - currentDreamMakerOutput) - .ConfigureAwait(false); + currentChatCallback?.Invoke( + FormatExceptionForUsers(ex), + currentDreamMakerOutput); throw; } @@ -797,15 +804,13 @@ namespace Tgstation.Server.Host.Components.Deployment try { using var byondLock = await byond.UseExecutables(null, cancellationToken).ConfigureAwait(false); - currentChatCallback = await chatManager.SendDeploymentMessage( + currentChatCallback = chatManager.QueueDeploymentMessage( revisionInformation, byondLock.Version, DateTimeOffset.Now + estimatedDuration, repository.RemoteRepositoryOwner, repository.RemoteRepositoryName, - localCommitExistsOnRemote, - cancellationToken) - .ConfigureAwait(false); + localCommitExistsOnRemote); var job = new Models.CompileJob { diff --git a/src/Tgstation.Server.Host/Components/Session/SessionController.cs b/src/Tgstation.Server.Host/Components/Session/SessionController.cs index 79b17c930b..1e068f0188 100644 --- a/src/Tgstation.Server.Host/Components/Session/SessionController.cs +++ b/src/Tgstation.Server.Host/Components/Session/SessionController.cs @@ -343,7 +343,7 @@ namespace Tgstation.Server.Host.Components.Session } /// - public async Task ProcessBridgeRequest(BridgeParameters parameters, CancellationToken cancellationToken) + public Task ProcessBridgeRequest(BridgeParameters parameters, CancellationToken cancellationToken) { if (parameters == null) throw new ArgumentNullException(nameof(parameters)); @@ -358,34 +358,36 @@ namespace Tgstation.Server.Host.Components.Session { case BridgeCommandType.ChatSend: if (parameters.ChatMessage == null) - return new BridgeResponse - { - ErrorMessage = "Missing chatMessage field!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Missing chatMessage field!" + }); if (parameters.ChatMessage.ChannelIds == null) - return new BridgeResponse - { - ErrorMessage = "Missing channelIds field in chatMessage!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Missing channelIds field in chatMessage!" + }); if (parameters.ChatMessage.ChannelIds.Any(channelIdString => !UInt64.TryParse(channelIdString, out var _))) - return new BridgeResponse - { - ErrorMessage = "Invalid channelIds in chatMessage!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Invalid channelIds in chatMessage!" + }); if (parameters.ChatMessage.Text == null) - return new BridgeResponse - { - ErrorMessage = "Missing message field in chatMessage!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Missing message field in chatMessage!" + }); - await chat.SendMessage( + chat.QueueMessage( parameters.ChatMessage.Text, - parameters.ChatMessage.ChannelIds.Select(UInt64.Parse), - cancellationToken) - .ConfigureAwait(false); + parameters.ChatMessage.ChannelIds.Select(UInt64.Parse)); break; case BridgeCommandType.Prime: var oldPrimeTcs = primeTcs; @@ -404,10 +406,11 @@ namespace Tgstation.Server.Host.Components.Session { /////UHHHH logger.LogWarning("DreamDaemon sent new port command without providing it's own!"); - return new BridgeResponse - { - ErrorMessage = "Missing stringified port as data parameter!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Missing stringified port as data parameter!" + }); } var currentPort = parameters.CurrentPort.Value; @@ -434,19 +437,21 @@ namespace Tgstation.Server.Host.Components.Session case BridgeCommandType.Startup: apiValidationStatus = ApiValidationStatus.BadValidationRequest; if (parameters.Version == null) - return new BridgeResponse - { - ErrorMessage = "Missing dmApiVersion field!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Missing dmApiVersion field!" + }); DMApiVersion = parameters.Version; if (DMApiVersion.Major != DMApiConstants.Version.Major) { apiValidationStatus = ApiValidationStatus.Incompatible; - return new BridgeResponse - { - ErrorMessage = "Incompatible dmApiVersion!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Incompatible dmApiVersion!" + }); } switch (parameters.MinimumSecurityLevel) @@ -461,15 +466,17 @@ namespace Tgstation.Server.Host.Components.Session apiValidationStatus = ApiValidationStatus.RequiresTrusted; break; case null: - return new BridgeResponse - { - ErrorMessage = "Missing minimumSecurityLevel field!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Missing minimumSecurityLevel field!" + }); default: - return new BridgeResponse - { - ErrorMessage = "Invalid minimumSecurityLevel!" - }; + return Task.FromResult( + new BridgeResponse + { + ErrorMessage = "Invalid minimumSecurityLevel!" + }); } response.RuntimeInformation = new RuntimeInformation( @@ -504,7 +511,7 @@ namespace Tgstation.Server.Host.Components.Session break; } - return response; + return Task.FromResult(response); } } diff --git a/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs b/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs index 48098d98de..c364a8cf65 100644 --- a/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs +++ b/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs @@ -95,7 +95,7 @@ namespace Tgstation.Server.Host.Components.Watchdog if (Server.RebootState == Session.RebootState.Shutdown) { // the time for graceful shutdown is now - await Chat.SendWatchdogMessage( + await Chat.QueueWatchdogMessage( String.Format( CultureInfo.InvariantCulture, "Server {0}! Shutting down due to graceful termination request...", @@ -105,7 +105,7 @@ namespace Tgstation.Server.Host.Components.Watchdog return MonitorAction.Exit; } - await Chat.SendWatchdogMessage( + await Chat.QueueWatchdogMessage( String.Format( CultureInfo.InvariantCulture, "Server {0}! Rebooting...", @@ -132,7 +132,7 @@ namespace Tgstation.Server.Host.Components.Watchdog return MonitorAction.Restart; case Session.RebootState.Shutdown: // graceful shutdown time - await Chat.SendWatchdogMessage( + await Chat.QueueWatchdogMessage( "Active server rebooted! Shutting down due to graceful termination request...", cancellationToken) .ConfigureAwait(false); @@ -265,7 +265,7 @@ namespace Tgstation.Server.Host.Components.Watchdog { gracefulRebootRequired = true; if (Server.CompileJob.DMApiVersion == null) - return Chat.SendWatchdogMessage( + return Chat.QueueWatchdogMessage( "A new deployment has been made but cannot be applied automatically as the currently running server has no DMAPI. Please manually reboot the server to apply the update.", cancellationToken); return Server.SetRebootState(Session.RebootState.Restart, cancellationToken); diff --git a/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs b/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs index 8275f6df83..cf873601f4 100644 --- a/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs +++ b/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs @@ -269,7 +269,7 @@ namespace Tgstation.Server.Host.Components.Watchdog { var eventTask = eventConsumer.HandleEvent(releaseServers ? EventType.WatchdogDetach : EventType.WatchdogShutdown, null, cancellationToken); - var chatTask = announce ? Chat.SendWatchdogMessage("Shutting down...", cancellationToken) : Task.CompletedTask; + var chatTask = announce ? Chat.QueueWatchdogMessage("Shutting down...", cancellationToken) : Task.CompletedTask; await eventTask.ConfigureAwait(false); @@ -314,7 +314,7 @@ namespace Tgstation.Server.Host.Components.Watchdog case 2: var message2 = "DEFCON 3: DreamDaemon has missed 2 heartbeats!"; Logger.LogInformation(message2); - await Chat.SendWatchdogMessage(message2, cancellationToken).ConfigureAwait(false); + await Chat.QueueWatchdogMessage(message2, cancellationToken).ConfigureAwait(false); break; case 3: var actionToTake = shouldShutdown @@ -322,7 +322,7 @@ namespace Tgstation.Server.Host.Components.Watchdog : "be restarted"; var message3 = $"DEFCON 2: DreamDaemon has missed 3 heartbeats! If it does not respond to the next one, the watchdog will {actionToTake}!"; Logger.LogWarning(message3); - await Chat.SendWatchdogMessage(message3, cancellationToken).ConfigureAwait(false); + await Chat.QueueWatchdogMessage(message3, cancellationToken).ConfigureAwait(false); break; case 4: var actionTaken = shouldShutdown @@ -330,7 +330,7 @@ namespace Tgstation.Server.Host.Components.Watchdog : "Restarting"; var message4 = $"DEFCON 1: Four heartbeats have been missed! {actionTaken}..."; Logger.LogWarning(message4); - await Chat.SendWatchdogMessage(message4, cancellationToken).ConfigureAwait(false); + await Chat.QueueWatchdogMessage(message4, cancellationToken).ConfigureAwait(false); await DisposeAndNullControllers(cancellationToken).ConfigureAwait(false); return shouldShutdown ? MonitorAction.Exit : MonitorAction.Restart; default: @@ -371,7 +371,7 @@ namespace Tgstation.Server.Host.Components.Watchdog Task announceTask; if (announce) { - announceTask = Chat.SendWatchdogMessage( + announceTask = Chat.QueueWatchdogMessage( reattachInfo == null ? "Launching..." : "Reattaching...", @@ -405,7 +405,7 @@ namespace Tgstation.Server.Host.Components.Watchdog { await originalChatTask.ConfigureAwait(false); if (announceFailure) - await Chat.SendWatchdogMessage("Startup failed!", cancellationToken).ConfigureAwait(false); + await Chat.QueueWatchdogMessage("Startup failed!", cancellationToken).ConfigureAwait(false); } announceTask = ChainChatTaskWithErrorMessage(); @@ -487,7 +487,7 @@ namespace Tgstation.Server.Host.Components.Watchdog const string FailReattachMessage = "Unable to properly reattach to server! Restarting watchdog..."; Logger.LogWarning(FailReattachMessage); - var chatTask = Chat.SendWatchdogMessage(FailReattachMessage, cancellationToken); + var chatTask = Chat.QueueWatchdogMessage(FailReattachMessage, cancellationToken); await InitControllers(chatTask, null, cancellationToken).ConfigureAwait(false); } @@ -591,7 +591,7 @@ namespace Tgstation.Server.Host.Components.Watchdog Math.Pow(2, retryAttempts)), TimeSpan.FromHours(1).TotalSeconds); // max of one hour, increasing by a power of 2 each time - chatTask = Chat.SendWatchdogMessage( + chatTask = Chat.QueueWatchdogMessage( $"Failed to restart (Attempt: {retryAttempts}), retrying in {retryDelay}s...", cancellationToken); @@ -733,7 +733,7 @@ namespace Tgstation.Server.Host.Components.Watchdog var nextActionMessage = nextAction != MonitorAction.Exit ? "Recovering" : "Shutting down"; - var chatTask = Chat.SendWatchdogMessage( + var chatTask = Chat.QueueWatchdogMessage( $"Monitor crashed, this should NEVER happen! Please report this, full details in logs! {nextActionMessage}. Error: {e.Message}", cancellationToken); @@ -811,22 +811,19 @@ namespace Tgstation.Server.Host.Components.Watchdog .ConfigureAwait(false); if (result?.InteropResponse?.ChatResponses != null) - await Task.WhenAll( - result.InteropResponse.ChatResponses.Select( - x => Chat.SendMessage( - x.Text, - x.ChannelIds - .Select(channelIdString => - { - if (UInt64.TryParse(channelIdString, out var channelId)) - return (ulong?)channelId; + foreach (var response in result.InteropResponse.ChatResponses) + Chat.QueueMessage( + response.Text, + response.ChannelIds + .Select(channelIdString => + { + if (UInt64.TryParse(channelIdString, out var channelId)) + return (ulong?)channelId; - return null; - }) - .Where(nullableChannelId => nullableChannelId.HasValue) - .Select(nullableChannelId => nullableChannelId.Value), - cancellationToken))) - .ConfigureAwait(false); + return null; + }) + .Where(nullableChannelId => nullableChannelId.HasValue) + .Select(nullableChannelId => nullableChannelId.Value)); } /// @@ -888,7 +885,7 @@ namespace Tgstation.Server.Host.Components.Watchdog { if (!graceful) { - var chatTask = Chat.SendWatchdogMessage("Manual restart triggered...", cancellationToken); + var chatTask = Chat.QueueWatchdogMessage("Manual restart triggered...", cancellationToken); await TerminateNoLock(false, false, cancellationToken).ConfigureAwait(false); await LaunchNoLock(true, false, true, null, cancellationToken).ConfigureAwait(false); await chatTask.ConfigureAwait(false); @@ -966,7 +963,7 @@ namespace Tgstation.Server.Host.Components.Watchdog { releaseServers = true; if (Status == WatchdogStatus.Online) - await Chat.SendWatchdogMessage("Detaching...", cancellationToken).ConfigureAwait(false); + await Chat.QueueWatchdogMessage("Detaching...", cancellationToken).ConfigureAwait(false); else Logger.LogTrace("Not sending detach chat message as status is: {0}", Status); } diff --git a/tools/ReleaseNotes/Program.cs b/tools/ReleaseNotes/Program.cs index 28cafa17a7..65760f4b8d 100644 --- a/tools/ReleaseNotes/Program.cs +++ b/tools/ReleaseNotes/Program.cs @@ -60,7 +60,6 @@ namespace ReleaseNotes { Milestone = $"v{versionString}", Type = IssueTypeQualifier.PullRequest, - State = ItemState.Closed, Repos = { { RepoOwner, RepoName } } }).ConfigureAwait(false);