From fce684f407f8b36cffee2d58e3925f72c018b699 Mon Sep 17 00:00:00 2001 From: Dominion Date: Sun, 16 Apr 2023 21:56:28 -0400 Subject: [PATCH] Closes #1457 - Delay bridge requests while detached up to one minute before failing - Add test coverage --- src/DMAPI/tgs/v5/api.dm | 2 ++ src/DMAPI/tgs/v5/bridge.dm | 11 ++++++++++ src/DMAPI/tgs/v5/topic.dm | 4 ++++ .../Components/Watchdog/BasicWatchdog.cs | 6 +++--- .../Components/Watchdog/WatchdogBase.cs | 13 ++++++++---- tests/DMAPI/LongRunning/Test.dm | 21 +++++++++++++++++++ .../Instance/WatchdogTest.cs | 18 ++++++++-------- .../Tgstation.Server.Tests/IntegrationTest.cs | 15 ++++++++++++- 8 files changed, 73 insertions(+), 17 deletions(-) diff --git a/src/DMAPI/tgs/v5/api.dm b/src/DMAPI/tgs/v5/api.dm index 91f4d73990..517240f12f 100644 --- a/src/DMAPI/tgs/v5/api.dm +++ b/src/DMAPI/tgs/v5/api.dm @@ -20,6 +20,8 @@ var/chunked_requests = 0 var/list/chunked_topics = list() + var/detached = FALSE + /datum/tgs_api/v5/ApiVersion() return new /datum/tgs_version( #include "__interop_version.dm" diff --git a/src/DMAPI/tgs/v5/bridge.dm b/src/DMAPI/tgs/v5/bridge.dm index df4825f89e..b3cf775939 100644 --- a/src/DMAPI/tgs/v5/bridge.dm +++ b/src/DMAPI/tgs/v5/bridge.dm @@ -60,6 +60,17 @@ return json /datum/tgs_api/v5/proc/PerformBridgeRequest(bridge_request) + if(detached) + // Wait up to one minute + for(var/i in 1 to 600) + sleep(1) + if(!detached) + break + + // dad went out for milk cigarettes 20 years ago... + if(i == 600) + detached = FALSE + // This is an infinite sleep until we get a response var/export_response = world.Export(bridge_request) if(!export_response) diff --git a/src/DMAPI/tgs/v5/topic.dm b/src/DMAPI/tgs/v5/topic.dm index c54ee810f7..28fcc14aef 100644 --- a/src/DMAPI/tgs/v5/topic.dm +++ b/src/DMAPI/tgs/v5/topic.dm @@ -69,6 +69,9 @@ return TopicResponse("Invalid or missing [DMAPI5_EVENT_NOTIFICATION_PARAMETERS]!") var/list/event_call = list(event_type) + if (event_type == TGS_EVENT_WATCHDOG_DETACH) + detached = TRUE + if(event_parameters) event_call += event_parameters @@ -137,6 +140,7 @@ return TopicResponse() if(DMAPI5_TOPIC_COMMAND_WATCHDOG_REATTACH) + detached = FALSE var/new_port = topic_parameters[DMAPI5_TOPIC_PARAMETER_NEW_PORT] var/error_message = null if (new_port != null) diff --git a/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs b/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs index 6c7be93c1c..a44adb23da 100644 --- a/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs +++ b/src/Tgstation.Server.Host/Components/Watchdog/BasicWatchdog.cs @@ -111,7 +111,7 @@ namespace Tgstation.Server.Host.Components.Watchdog var eventType = Server.TerminationWasRequested ? EventType.WorldEndProcess : EventType.WatchdogCrash; - await HandleNonRelayedEvent(eventType, Enumerable.Empty(), cancellationToken); + await HandleEvent(eventType, Enumerable.Empty(), false, cancellationToken); var exitWord = Server.TerminationWasRequested ? "exited" : "crashed"; if (Server.RebootState == Session.RebootState.Shutdown) @@ -146,7 +146,7 @@ namespace Tgstation.Server.Host.Components.Watchdog gracefulRebootRequired = false; Server.ResetRebootState(); - await HandleNonRelayedEvent(EventType.WorldReboot, Enumerable.Empty(), cancellationToken); + await HandleEvent(EventType.WorldReboot, Enumerable.Empty(), false, cancellationToken); switch (rebootState) { @@ -173,7 +173,7 @@ namespace Tgstation.Server.Host.Components.Watchdog await HandleNewDmbAvailable(cancellationToken); break; case MonitorActivationReason.ActiveServerPrimed: - await HandleNonRelayedEvent(EventType.WorldPrime, Enumerable.Empty(), cancellationToken); + await HandleEvent(EventType.WorldPrime, Enumerable.Empty(), false, cancellationToken); break; case MonitorActivationReason.Heartbeat: default: diff --git a/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs b/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs index a708244fe1..811cd907ed 100644 --- a/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs +++ b/src/Tgstation.Server.Host/Components/Watchdog/WatchdogBase.cs @@ -513,7 +513,7 @@ namespace Tgstation.Server.Host.Components.Watchdog cancellationToken); // simple announce if (reattachInfo == null) announceTask = Task.WhenAll( - HandleNonRelayedEvent(EventType.WatchdogLaunch, Enumerable.Empty(), cancellationToken), + HandleEvent(EventType.WatchdogLaunch, Enumerable.Empty(), false, cancellationToken), announceTask); } else @@ -697,13 +697,17 @@ namespace Tgstation.Server.Host.Components.Watchdog /// /// The . /// An of parameters for . + /// If the event should be sent to DreamDaemon. /// The for the operation. /// A representing the running operation. - protected async Task HandleNonRelayedEvent(EventType eventType, IEnumerable parameters, CancellationToken cancellationToken) + protected async Task HandleEvent(EventType eventType, IEnumerable parameters, bool relayToSession, CancellationToken cancellationToken) { try { - await eventConsumer.HandleEvent(eventType, parameters, cancellationToken); + var sessionEventTask = relayToSession ? ((IEventConsumer)this).HandleEvent(eventType, parameters, cancellationToken) : Task.CompletedTask; + await Task.WhenAll( + eventConsumer.HandleEvent(eventType, parameters, cancellationToken), + sessionEventTask); } catch (JobException ex) { @@ -1021,11 +1025,12 @@ namespace Tgstation.Server.Host.Components.Watchdog return; if (!graceful) { - var eventTask = HandleNonRelayedEvent( + var eventTask = HandleEvent( releaseServers ? EventType.WatchdogDetach : EventType.WatchdogShutdown, Enumerable.Empty(), + releaseServers, cancellationToken); var chatTask = announce ? Chat.QueueWatchdogMessage("Shutting down...", cancellationToken) : Task.CompletedTask; diff --git a/tests/DMAPI/LongRunning/Test.dm b/tests/DMAPI/LongRunning/Test.dm index 2fd8048d58..ca3164e217 100644 --- a/tests/DMAPI/LongRunning/Test.dm +++ b/tests/DMAPI/LongRunning/Test.dm @@ -130,10 +130,31 @@ var/run_bridge_test TgsChatBroadcast(new /datum/tgs_message_content(create_payload(3000))) return "sent" + // Bridge response queuing + var/tactics6 = data["tgs_integration_test_tactics6"] + if(tactics6) + DetachedChatMessageQueuing() + return "queued" + TgsChatBroadcast(new /datum/tgs_message_content("Recieved non-tgs topic: `[T]`")) return "feck" +// Look I always forget how waitfor = FALSE works +/proc/DetachedChatMessageQueuing() + set waitfor = FALSE + DetachedChatMessageQueuingP2() + +/proc/DetachedChatMessageQueuingP2() + sleep(1) + DetachedChatMessageQueuingP3() + +/proc/DetachedChatMessageQueuingP3() + set waitfor = FALSE + world.TgsChatBroadcast(new /datum/tgs_message_content("1/3 queued detached chat messages")) + world.TgsChatBroadcast(new /datum/tgs_message_content("2/3 queued detached chat messages")) + world.TgsChatBroadcast(new /datum/tgs_message_content("3/3 queued detached chat messages")) + /world/Reboot(reason) TgsChatBroadcast("World Rebooting") TgsReboot() diff --git a/tests/Tgstation.Server.Tests/Instance/WatchdogTest.cs b/tests/Tgstation.Server.Tests/Instance/WatchdogTest.cs index 9ccd87b122..97076d702f 100644 --- a/tests/Tgstation.Server.Tests/Instance/WatchdogTest.cs +++ b/tests/Tgstation.Server.Tests/Instance/WatchdogTest.cs @@ -108,7 +108,7 @@ namespace Tgstation.Server.Tests.Instance async Task SendChatOverloadCommand(CancellationToken cancellationToken) { // for the code coverage really... - var topicRequestResult = await topicClient.SendTopic( + var topicRequestResult = await TopicClient.SendTopic( IPAddress.Loopback, $"tgs_integration_test_tactics5=1", IntegrationTest.DDPort, @@ -379,7 +379,7 @@ namespace Tgstation.Server.Tests.Instance System.Console.WriteLine("TEST: Sending Bridge tests topic..."); - var bridgeTestTopicResult = await topicClient.SendTopic(IPAddress.Loopback, "tgs_integration_test_tactics2=1", IntegrationTest.DDPort, cancellationToken); + var bridgeTestTopicResult = await TopicClient.SendTopic(IPAddress.Loopback, "tgs_integration_test_tactics2=1", IntegrationTest.DDPort, cancellationToken); Assert.AreEqual("ack2", bridgeTestTopicResult.StringData); await bridgeTestsTcs.Task.WithToken(cancellationToken); @@ -405,7 +405,7 @@ namespace Tgstation.Server.Tests.Instance }; var json = JsonConvert.SerializeObject(baseTopic, DMApiConstants.SerializerSettings); - var topicString = $"tgs_integration_test_tactics3={topicClient.SanitizeString(json)}"; + var topicString = $"tgs_integration_test_tactics3={TopicClient.SanitizeString(json)}"; var baseSize = topicString.Length; var wrappingSize = baseSize; @@ -424,9 +424,9 @@ namespace Tgstation.Server.Tests.Instance TopicResponse topicRequestResult = null; try { - topicRequestResult = await topicClient.SendTopic( + topicRequestResult = await TopicClient.SendTopic( IPAddress.Loopback, - $"tgs_integration_test_tactics3={topicClient.SanitizeString(JsonConvert.SerializeObject(topic, DMApiConstants.SerializerSettings))}", + $"tgs_integration_test_tactics3={TopicClient.SanitizeString(JsonConvert.SerializeObject(topic, DMApiConstants.SerializerSettings))}", IntegrationTest.DDPort, cancellationToken); } @@ -462,9 +462,9 @@ namespace Tgstation.Server.Tests.Instance while (!cancellationToken.IsCancellationRequested) { var currentSize = baseSize + (int)Math.Pow(2, nextPow); - var topicRequestResult = await topicClient.SendTopic( + var topicRequestResult = await TopicClient.SendTopic( IPAddress.Loopback, - $"tgs_integration_test_tactics4={topicClient.SanitizeString(currentSize.ToString())}", + $"tgs_integration_test_tactics4={TopicClient.SanitizeString(currentSize.ToString())}", IntegrationTest.DDPort, cancellationToken); @@ -775,7 +775,7 @@ namespace Tgstation.Server.Tests.Instance return ddProc != null; } - readonly TopicClient topicClient = new (new SocketParameters + public static readonly TopicClient TopicClient = new (new SocketParameters { SendTimeout = TimeSpan.FromSeconds(30), ReceiveTimeout = TimeSpan.FromSeconds(30), @@ -791,7 +791,7 @@ namespace Tgstation.Server.Tests.Instance try { System.Console.WriteLine("TEST: Sending world reboot topic..."); - var result = await topicClient.SendTopic(IPAddress.Loopback, "tgs_integration_test_special_tactics=1", IntegrationTest.DDPort, cancellationToken); + var result = await TopicClient.SendTopic(IPAddress.Loopback, "tgs_integration_test_special_tactics=1", IntegrationTest.DDPort, cancellationToken); Assert.AreEqual("ack", result.StringData); using (var tempCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken)) diff --git a/tests/Tgstation.Server.Tests/IntegrationTest.cs b/tests/Tgstation.Server.Tests/IntegrationTest.cs index 282b4a1482..f78ae71c02 100644 --- a/tests/Tgstation.Server.Tests/IntegrationTest.cs +++ b/tests/Tgstation.Server.Tests/IntegrationTest.cs @@ -1,4 +1,6 @@ -using Microsoft.EntityFrameworkCore; +using Byond.TopicSender; + +using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.Infrastructure; using Microsoft.EntityFrameworkCore.Migrations; using Microsoft.Extensions.DependencyInjection; @@ -888,6 +890,17 @@ namespace Tgstation.Server.Tests await Task.WhenAny(serverTask, Task.Delay(TimeSpan.FromMinutes(1), cancellationToken)); Assert.IsTrue(serverTask.IsCompleted); + // test the reattach message queueing + // for the code coverage really... + var topicRequestResult = await WatchdogTest.TopicClient.SendTopic( + IPAddress.Loopback, + $"tgs_integration_test_tactics6=1", + DDPort, + cancellationToken); + + Assert.IsNotNull(topicRequestResult); + Assert.AreEqual("queued", topicRequestResult.StringData); + // http bind test https://github.com/tgstation/tgstation-server/issues/1065 if (new PlatformIdentifier().IsWindows) {