tgstation-server 6.11.1
The /tg/station 13 server suite
Loading...
Searching...
No Matches
SessionController.cs
Go to the documentation of this file.
1using System;
2using System.Collections.Generic;
3using System.Globalization;
4using System.Linq;
5using System.Text;
6using System.Threading;
7using System.Threading.Tasks;
8
9using Microsoft.Extensions.Logging;
10using Microsoft.Extensions.Logging.Abstractions;
11
12using Newtonsoft.Json;
13
14using Serilog.Context;
15
30
32{
35 {
39 internal static bool LogTopicRequests { get; set; } = true;
40
43
46 {
47 get
48 {
49 if (!Lifetime.IsCompleted)
50 throw new InvalidOperationException("ApiValidated cannot be checked while Lifetime is incomplete!");
52 }
53 }
54
57
60
63
65 public Version? DMApiVersion { get; private set; }
66
68 public bool TerminationWasIntentional => terminationWasIntentional || (Lifetime.IsCompleted && Lifetime.Result == 0);
69
71 public Task<LaunchResult> LaunchResult { get; }
72
74 public Task<int?> Lifetime { get; }
75
77 public Task OnStartup => startupTcs.Task;
78
80 public Task OnReboot => rebootTcs.Task;
81
83 public Task RebootGate
84 {
85 get => rebootGate;
86 set
87 {
88 var tcs = new TaskCompletionSource<Task>();
89 async Task Wrap()
90 {
91 var toAwait = await tcs.Task;
92 await toAwait;
93 await value;
94 }
95
96 tcs.SetResult(Interlocked.Exchange(ref rebootGate, Wrap()));
97 }
98 }
99
101 public Task OnPrime => primeTcs.Task;
102
105
108
110 public string DumpFileExtension => engineLock.UseDotnetDump
111 ? ".net.dmp"
112 : ".dmp";
113
118
123
126
128 public DateTimeOffset? LaunchTime => process.LaunchTime;
129
133 readonly Byond.TopicSender.ITopicClient byondTopicSender;
134
139
144
149
154
159
164
169
174
178 readonly TaskCompletionSource initialBridgeRequestTcs;
179
184
188 readonly CancellationTokenSource sessionDurationCts;
189
193 readonly object synchronizationLock;
194
198 readonly bool apiValidationSession;
199
203 volatile TaskCompletionSource startupTcs;
204
208 volatile TaskCompletionSource rebootTcs;
209
213 volatile TaskCompletionSource primeTcs;
214
218 volatile Task rebootGate;
219
224
229
234
239
244
249
254
276 ReattachInformation reattachInformation,
277 Api.Models.Instance metadata,
280 Byond.TopicSender.ITopicClient byondTopicSender,
282 IBridgeRegistrar bridgeRegistrar,
284 IAssemblyInformationProvider assemblyInformationProvider,
288 ILogger<SessionController> logger,
289 Func<ValueTask> postLifetimeCallback,
290 uint? startupTimeout,
291 bool reattached,
292 bool apiValidate)
293 : base(logger)
294 {
295 ReattachInformation = reattachInformation ?? throw new ArgumentNullException(nameof(reattachInformation));
296 this.metadata = metadata ?? throw new ArgumentNullException(nameof(metadata));
297 this.process = process ?? throw new ArgumentNullException(nameof(process));
298 this.engineLock = engineLock ?? throw new ArgumentNullException(nameof(engineLock));
299 this.byondTopicSender = byondTopicSender ?? throw new ArgumentNullException(nameof(byondTopicSender));
300 this.chatTrackingContext = chatTrackingContext ?? throw new ArgumentNullException(nameof(chatTrackingContext));
301 ArgumentNullException.ThrowIfNull(bridgeRegistrar);
302
303 this.chat = chat ?? throw new ArgumentNullException(nameof(chat));
304 ArgumentNullException.ThrowIfNull(assemblyInformationProvider);
305
306 this.asyncDelayer = asyncDelayer ?? throw new ArgumentNullException(nameof(asyncDelayer));
307 this.dotnetDumpService = dotnetDumpService ?? throw new ArgumentNullException(nameof(dotnetDumpService));
308 this.eventConsumer = eventConsumer ?? throw new ArgumentNullException(nameof(eventConsumer));
309
310 apiValidationSession = apiValidate;
311
312 disposed = false;
314 released = false;
315
316 startupTcs = new TaskCompletionSource();
317 rebootTcs = new TaskCompletionSource();
318 primeTcs = new TaskCompletionSource();
319
320 rebootGate = Task.CompletedTask;
321 customEventProcessingTask = Task.CompletedTask;
322
323 // Run this asynchronously because we want to try to avoid any effects sending topics to the server while the initial bridge request is processing
324 // It MAY be the source of a DD crash. See this gist https://gist.github.com/Cyberboss/7776bbeff3a957d76affe0eae95c9f14
325 // Worth further investigation as to if that sequence of events is a reliable crash vector and opening a BYOND bug if it is
326 initialBridgeRequestTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
327 sessionDurationCts = new CancellationTokenSource();
328
330 synchronizationLock = new object();
331
333 {
334 bridgeRegistration = bridgeRegistrar.RegisterHandler(this);
335 this.chatTrackingContext.SetChannelSink(this);
336 }
337 else
338 logger.LogTrace(
339 "Not registering session with {reasonWhyDmApiIsBad} DMAPI version for interop!",
340 reattachInformation.Dmb.CompileJob.DMApiVersion == null
341 ? "no"
342 : $"incompatible ({reattachInformation.Dmb.CompileJob.DMApiVersion})");
343
344 async Task<int?> WrapLifetime()
345 {
346 var exitCode = await process.Lifetime;
347 await postLifetimeCallback();
348 if (postValidationShutdownTask != null)
350
351 return exitCode;
352 }
353
354 Lifetime = WrapLifetime();
355
357 assemblyInformationProvider,
359 startupTimeout,
360 reattached,
361 apiValidate);
362
363 logger.LogDebug(
364 "Created session controller. CommsKey: {accessIdentifier}, Port: {port}",
365 reattachInformation.AccessIdentifier,
366 reattachInformation.Port);
367 }
368
370 public async ValueTask DisposeAsync()
371 {
373 {
374 if (disposed)
375 return;
376 disposed = true;
377 }
378
379 Logger.LogTrace("Disposing...");
380
381 sessionDurationCts.Cancel();
382 var cancellationToken = CancellationToken.None; // DCT: None available
383 var semaphoreLockTask = TopicSendSemaphore.Lock(cancellationToken);
384
385 if (!released)
386 {
388 Logger,
389 process,
392 cancellationToken);
393 }
394
395 await process.DisposeAsync();
396 engineLock.Dispose();
397 bridgeRegistration?.Dispose();
398 var regularDmbDisposeTask = ReattachInformation.Dmb.DisposeAsync();
399 var initialDmb = ReattachInformation.InitialDmb;
400 if (initialDmb != null)
401 await initialDmb.DisposeAsync();
402
403 await regularDmbDisposeTask;
404
405 chatTrackingContext.Dispose();
406 sessionDurationCts.Dispose();
407
408 if (!released)
409 await Lifetime; // finish the async callback
410
411 (await semaphoreLockTask).Dispose();
413
415 }
416
418 public async ValueTask<BridgeResponse?> ProcessBridgeRequest(BridgeParameters parameters, CancellationToken cancellationToken)
419 {
420 ArgumentNullException.ThrowIfNull(parameters);
421
422 using (LogContext.PushProperty(SerilogContextHelper.InstanceIdContextProperty, metadata.Id))
423 {
424 Logger.LogTrace("Handling bridge request...");
425
426 try
427 {
428 return await ProcessBridgeCommand(parameters, cancellationToken);
429 }
430 finally
431 {
432 initialBridgeRequestTcs.TrySetResult();
433 }
434 }
435 }
436
438 public ValueTask Release()
439 {
441
445 released = true;
446 return DisposeAsync();
447 }
448
450 public ValueTask<TopicResponse?> SendCommand(TopicParameters parameters, CancellationToken cancellationToken)
451 => SendCommand(parameters, false, cancellationToken);
452
454 public async ValueTask<bool> SetRebootState(RebootState newRebootState, CancellationToken cancellationToken)
455 {
456 if (RebootState == newRebootState)
457 return true;
458
459 Logger.LogTrace("Changing reboot state to {newRebootState}", newRebootState);
460
461 ReattachInformation.RebootState = newRebootState;
462 var result = await SendCommand(
463 new TopicParameters(newRebootState),
464 cancellationToken);
465
466 return result != null && result.ErrorMessage == null;
467 }
468
470 public void ResetRebootState()
471 {
473 Logger.LogTrace("Resetting reboot state...");
474 ReattachInformation.RebootState = RebootState.Normal;
475 }
476
478 public void AdjustPriority(bool higher) => process.AdjustPriority(higher);
479
482
485
488 {
489 var oldDmb = ReattachInformation.Dmb;
490 ReattachInformation.Dmb = dmbProvider ?? throw new ArgumentNullException(nameof(dmbProvider));
491 return oldDmb;
492 }
493
495 public async ValueTask InstanceRenamed(string newInstanceName, CancellationToken cancellationToken)
496 {
497 var runtimeInformation = ReattachInformation.RuntimeInformation;
498 if (runtimeInformation != null)
499 runtimeInformation.InstanceName = newInstanceName;
500
501 await SendCommand(
503 cancellationToken);
504 }
505
507 public async ValueTask UpdateChannels(IEnumerable<ChannelRepresentation> newChannels, CancellationToken cancellationToken)
508 => await SendCommand(
509 new TopicParameters(
510 new ChatUpdate(newChannels)),
511 cancellationToken);
512
514 public ValueTask CreateDump(string outputFile, bool minidump, CancellationToken cancellationToken)
515 {
517 return dotnetDumpService.Dump(process, outputFile, minidump, cancellationToken);
518
519 return process.CreateDump(outputFile, minidump, cancellationToken);
520 }
521
531 async Task<LaunchResult> GetLaunchResult(
532 IAssemblyInformationProvider assemblyInformationProvider,
534 uint? startupTimeout,
535 bool reattached,
536 bool apiValidate)
537 {
538 var startTime = DateTimeOffset.UtcNow;
539 var useBridgeRequestForLaunchResult = !reattached && (apiValidate || DMApiAvailable);
540 var startupTask = useBridgeRequestForLaunchResult
541 ? initialBridgeRequestTcs.Task
543 var toAwait = Task.WhenAny(startupTask, process.Lifetime);
544
545 if (startupTimeout.HasValue)
546 toAwait = Task.WhenAny(
547 toAwait,
549 TimeSpan.FromSeconds(startupTimeout.Value),
550 CancellationToken.None)); // DCT: None available, task will clean up after delay
551
552 Logger.LogTrace(
553 "Waiting for LaunchResult based on {launchResultCompletionCause}{possibleTimeout}...",
554 useBridgeRequestForLaunchResult ? "initial bridge request" : "process startup",
555 startupTimeout.HasValue ? $" with a timeout of {startupTimeout.Value}s" : String.Empty);
556
557 await toAwait;
558
559 var result = new LaunchResult
560 {
561 ExitCode = process.Lifetime.IsCompleted ? await process.Lifetime : null,
562 StartupTime = startupTask.IsCompleted ? (DateTimeOffset.UtcNow - startTime) : null,
563 };
564
565 Logger.LogTrace("Launch result: {launchResult}", result);
566
567 if (!result.ExitCode.HasValue && reattached && !disposed)
568 {
569 var reattachResponse = await SendCommand(
570 new TopicParameters(
571 assemblyInformationProvider.Version,
573 true,
574 sessionDurationCts.Token);
575
576 if (reattachResponse != null)
577 {
578 if (reattachResponse?.CustomCommands != null)
579 chatTrackingContext.CustomCommands = reattachResponse.CustomCommands;
580 else if (reattachResponse != null)
581 Logger.Log(
582 CompileJob.DMApiVersion >= new Version(5, 2, 0)
583 ? LogLevel.Warning
584 : LogLevel.Debug,
585 "DMAPI Interop v{interopVersion} isn't returning the TGS custom commands list. Functionality added in v5.2.0.",
586 CompileJob.DMApiVersion!.Semver());
587 }
588 }
589
590 return result;
591 }
592
596 void CheckDisposed() => ObjectDisposedException.ThrowIf(disposed, this);
597
603 async Task PostValidationShutdown(Task<bool> proceedTask)
604 {
605 Logger.LogTrace("Entered post validation terminate task.");
606 if (!await proceedTask)
607 {
608 Logger.LogTrace("Not running post validation terminate task for repeated bridge request.");
609 return;
610 }
611
612 const int GracePeriodSeconds = 30;
613 Logger.LogDebug("Server will terminated in {gracePeriodSeconds}s if it does not exit...", GracePeriodSeconds);
614 var delayTask = asyncDelayer.Delay(TimeSpan.FromSeconds(GracePeriodSeconds), CancellationToken.None); // DCT: None available
615 await Task.WhenAny(process.Lifetime, delayTask);
616
617 if (!process.Lifetime.IsCompleted)
618 {
619 Logger.LogWarning("DMAPI took too long to shutdown server after validation request!");
621 apiValidationStatus = ApiValidationStatus.BadValidationRequest;
622 }
623 else
624 Logger.LogTrace("Server exited properly post validation.");
625 }
626
633#pragma warning disable CA1502 // TODO: Decomplexify
634 async ValueTask<BridgeResponse?> ProcessBridgeCommand(BridgeParameters parameters, CancellationToken cancellationToken)
635 {
636 var response = new BridgeResponse();
637 switch (parameters.CommandType)
638 {
639 case BridgeCommandType.ChatSend:
640 if (parameters.ChatMessage == null)
641 return BridgeError("Missing chatMessage field!");
642
643 if (parameters.ChatMessage.ChannelIds == null)
644 return BridgeError("Missing channelIds field in chatMessage!");
645
646 if (parameters.ChatMessage.ChannelIds.Any(channelIdString => !UInt64.TryParse(channelIdString, out var _)))
647 return BridgeError("Invalid channelIds in chatMessage!");
648
649 if (parameters.ChatMessage.Text == null)
650 return BridgeError("Missing message field in chatMessage!");
651
652 var anyFailed = false;
653 var parsedChannels = parameters.ChatMessage.ChannelIds.Select(
654 channelString =>
655 {
656 anyFailed |= !UInt64.TryParse(channelString, out var channelId);
657 return channelId;
658 });
659
660 if (anyFailed)
661 return BridgeError("Failed to parse channelIds as U64!");
662
664 parameters.ChatMessage,
665 parsedChannels);
666 break;
667 case BridgeCommandType.Prime:
668 Interlocked.Exchange(ref primeTcs, new TaskCompletionSource()).SetResult();
669 break;
670 case BridgeCommandType.Kill:
671 Logger.LogInformation("Bridge requested process termination!");
672 chatTrackingContext.Active = false;
675 break;
676 case BridgeCommandType.DeprecatedPortUpdate:
677 return BridgeError("Port switching is no longer supported!");
678 case BridgeCommandType.Startup:
679 apiValidationStatus = ApiValidationStatus.BadValidationRequest;
680
682 {
683 var proceedTcs = new TaskCompletionSource<bool>();
684 var firstValidationRequest = Interlocked.CompareExchange(ref postValidationShutdownTask, PostValidationShutdown(proceedTcs.Task), null) == null;
685 proceedTcs.SetResult(firstValidationRequest);
686
687 if (!firstValidationRequest)
688 return BridgeError("Startup bridge request was repeated!");
689 }
690
691 if (parameters.Version == null)
692 return BridgeError("Missing dmApiVersion field!");
693
694 DMApiVersion = parameters.Version;
695
696 // TODO: When OD figures out how to unite port and topic_port, set an upper version bound on OD for this check
698 || (EngineVersion.Engine == EngineType.OpenDream && DMApiVersion < new Version(5, 7)))
699 {
701 return BridgeError("Incompatible dmApiVersion!");
702 }
703
704 switch (parameters.MinimumSecurityLevel)
705 {
706 case DreamDaemonSecurity.Ultrasafe:
707 apiValidationStatus = ApiValidationStatus.RequiresUltrasafe;
708 break;
709 case DreamDaemonSecurity.Safe:
711 break;
712 case DreamDaemonSecurity.Trusted:
714 break;
715 case null:
716 return BridgeError("Missing minimumSecurityLevel field!");
717 default:
718 return BridgeError("Invalid minimumSecurityLevel!");
719 }
720
721 Logger.LogTrace("ApiValidationStatus set to {apiValidationStatus}", apiValidationStatus);
722
723 // we create new runtime info here because of potential .Dmb changes (i think. i forget...)
724 response.RuntimeInformation = new RuntimeInformation(
733
734 if (parameters.TopicPort.HasValue)
735 {
736 var newTopicPort = parameters.TopicPort.Value;
737 Logger.LogInformation("Server is requesting use of port {topicPort} for topic communications", newTopicPort);
738 ReattachInformation.TopicPort = newTopicPort;
739 }
740
741 // Load custom commands
742 chatTrackingContext.CustomCommands = parameters.CustomCommands ?? Array.Empty<CustomCommand>();
743 chatTrackingContext.Active = true;
744 Interlocked.Exchange(ref startupTcs, new TaskCompletionSource()).SetResult();
745 break;
746 case BridgeCommandType.Reboot:
747 Interlocked.Increment(ref rebootBridgeRequestsProcessing);
748 try
749 {
750 chatTrackingContext.Active = false;
751 Interlocked.Exchange(ref rebootTcs, new TaskCompletionSource()).SetResult();
752 await RebootGate.WaitAsync(cancellationToken);
753 }
754 finally
755 {
756 Interlocked.Decrement(ref rebootBridgeRequestsProcessing);
757 }
758
759 break;
760 case BridgeCommandType.Chunk:
761 return await ProcessChunk<BridgeParameters, BridgeResponse>(ProcessBridgeCommand, BridgeError, parameters.Chunk, cancellationToken);
762 case BridgeCommandType.Event:
763 return TriggerCustomEvent(parameters.EventInvocation);
764 case null:
765 return BridgeError("Missing commandType!");
766 default:
767 return BridgeError($"commandType {parameters.CommandType} not supported!");
768 }
769
770 return response;
771 }
772#pragma warning restore CA1502
773
780 {
781 Logger.LogWarning("Bridge request error: {message}", message);
782 return new BridgeResponse
783 {
784 ErrorMessage = message,
785 };
786 }
787
794 async ValueTask<CombinedTopicResponse?> SendTopicRequest(TopicParameters parameters, CancellationToken cancellationToken)
795 {
796 parameters.AccessIdentifier = ReattachInformation.AccessIdentifier;
797
798 var fullCommandString = GenerateQueryString(parameters, out var json);
799 if (LogTopicRequests)
800 Logger.LogTrace("Topic request: {json}", json);
801 var fullCommandByteCount = Encoding.UTF8.GetByteCount(fullCommandString);
802 var topicPriority = parameters.IsPriority;
803 if (fullCommandByteCount <= DMApiConstants.MaximumTopicRequestLength)
804 return await SendRawTopic(fullCommandString, topicPriority, cancellationToken);
805
806 var interopChunkingVersion = new Version(5, 6, 0);
807 if (ReattachInformation.Dmb.CompileJob.DMApiVersion < interopChunkingVersion)
808 {
809 Logger.LogWarning(
810 "Cannot send topic request as it is exceeds the single request limit of {limitBytes}B ({actualBytes}B) and requires chunking and the current compile job's interop version must be at least {chunkingVersionRequired}!",
812 fullCommandByteCount,
813 interopChunkingVersion);
814 return null;
815 }
816
817 var payloadId = NextPayloadId;
818
819 // AccessIdentifer is just noise in a chunked request
820 parameters.AccessIdentifier = null!;
821 GenerateQueryString(parameters, out json);
822
823 // yes, this straight up ignores unicode, precalculating it is useless when we don't
824 // even know if the UTF8 bytes of the url encoded chunk will fit the window until we do said encoding
825 var fullPayloadSize = (uint)json.Length;
826
827 List<string>? chunkQueryStrings = null;
828 for (var chunkCount = 2; chunkQueryStrings == null; ++chunkCount)
829 {
830 var standardChunkSize = fullPayloadSize / chunkCount;
831 var bigChunkSize = standardChunkSize + (fullPayloadSize % chunkCount);
832 if (bigChunkSize > DMApiConstants.MaximumTopicRequestLength)
833 continue;
834
835 chunkQueryStrings = new List<string>();
836 for (var i = 0U; i < chunkCount; ++i)
837 {
838 var startIndex = i * standardChunkSize;
839 var subStringLength = Math.Min(
840 fullPayloadSize - startIndex,
841 i == chunkCount - 1
842 ? bigChunkSize
843 : standardChunkSize);
844 var chunkPayload = json.Substring((int)startIndex, (int)subStringLength);
845
846 var chunk = new ChunkData
847 {
848 Payload = chunkPayload,
849 PayloadId = payloadId,
850 SequenceId = i,
851 TotalChunks = (uint)chunkCount,
852 };
853
854 var chunkParameters = new TopicParameters(chunk)
855 {
856 AccessIdentifier = ReattachInformation.AccessIdentifier,
857 };
858
859 var chunkCommandString = GenerateQueryString(chunkParameters, out _);
860 if (Encoding.UTF8.GetByteCount(chunkCommandString) > DMApiConstants.MaximumTopicRequestLength)
861 {
862 // too long when encoded, need more chunks
863 chunkQueryStrings = null;
864 break;
865 }
866
867 chunkQueryStrings.Add(chunkCommandString);
868 }
869 }
870
871 Logger.LogTrace("Chunking topic request ({totalChunks} total)...", chunkQueryStrings.Count);
872
873 CombinedTopicResponse? combinedResponse = null;
874 bool LogRequestIssue(bool possiblyFromCompletedRequest)
875 {
876 if (combinedResponse?.InteropResponse == null || combinedResponse.InteropResponse.ErrorMessage != null)
877 {
878 Logger.LogWarning(
879 "Topic request {chunkingStatus} failed!{potentialRequestError}",
880 possiblyFromCompletedRequest ? "final chunk" : "chunking",
881 combinedResponse?.InteropResponse?.ErrorMessage != null
882 ? $" Request error: {combinedResponse.InteropResponse.ErrorMessage}"
883 : String.Empty);
884 return true;
885 }
886
887 return false;
888 }
889
890 foreach (var chunkCommandString in chunkQueryStrings)
891 {
892 combinedResponse = await SendRawTopic(chunkCommandString, topicPriority, cancellationToken);
893 if (LogRequestIssue(chunkCommandString == chunkQueryStrings.Last()))
894 return null;
895 }
896
897 while ((combinedResponse?.InteropResponse?.MissingChunks?.Count ?? 0) > 0)
898 {
899 Logger.LogWarning("DD is still missing some chunks of topic request P{payloadId}! Sending missing chunks...", payloadId);
900 var missingChunks = combinedResponse!.InteropResponse!.MissingChunks!;
901 var lastIndex = missingChunks.Last();
902 foreach (var missingChunkIndex in missingChunks)
903 {
904 var chunkCommandString = chunkQueryStrings[(int)missingChunkIndex];
905 combinedResponse = await SendRawTopic(chunkCommandString, topicPriority, cancellationToken);
906 if (LogRequestIssue(missingChunkIndex == lastIndex))
907 return null;
908 }
909 }
910
911 return combinedResponse;
912 }
913
920 string GenerateQueryString(TopicParameters parameters, out string json)
921 {
922 json = JsonConvert.SerializeObject(parameters, DMApiConstants.SerializerSettings);
923 var commandString = String.Format(
924 CultureInfo.InvariantCulture,
925 "?{0}={1}",
927 byondTopicSender.SanitizeString(json));
928 return commandString;
929 }
930
938 async ValueTask<CombinedTopicResponse?> SendRawTopic(string queryString, bool priority, CancellationToken cancellationToken)
939 {
940 if (disposed)
941 {
942 Logger.LogWarning(
943 "Attempted to send a topic on a disposed SessionController");
944 return null;
945 }
946
947 var targetPort = ReattachInformation.TopicPort ?? ReattachInformation.Port;
948 Byond.TopicSender.TopicResponse? byondResponse;
949 using (await TopicSendSemaphore.Lock(cancellationToken))
950 byondResponse = await byondTopicSender.SendWithOptionalPriority(
952 LogTopicRequests
953 ? Logger
954 : NullLogger.Instance,
955 queryString,
956 targetPort,
957 priority,
958 cancellationToken);
959
960 if (byondResponse == null)
961 {
962 if (priority)
963 Logger.LogError(
964 "Unable to send priority topic \"{queryString}\"!",
965 queryString);
966
967 return null;
968 }
969
970 var topicReturn = byondResponse.StringData;
971
972 TopicResponse? interopResponse = null;
973 if (topicReturn != null)
974 try
975 {
976 interopResponse = JsonConvert.DeserializeObject<TopicResponse>(topicReturn, DMApiConstants.SerializerSettings);
977 }
978 catch (Exception ex)
979 {
980 Logger.LogWarning(ex, "Invalid interop response: {topicReturnString}", topicReturn);
981 }
982
983 return new CombinedTopicResponse(byondResponse, interopResponse);
984 }
985
993 async ValueTask<TopicResponse?> SendCommand(TopicParameters parameters, bool bypassLaunchResult, CancellationToken cancellationToken)
994 {
995 ArgumentNullException.ThrowIfNull(parameters);
996
997 if (Lifetime.IsCompleted || disposed)
998 {
999 Logger.LogWarning(
1000 "Attempted to send a command to an inactive SessionController: {commandType}",
1001 parameters.CommandType);
1002 return null;
1003 }
1004
1005 if (!DMApiAvailable)
1006 {
1007 Logger.LogTrace("Not sending topic request {commandType} to server without/with incompatible DMAPI!", parameters.CommandType);
1008 return null;
1009 }
1010
1011 var reboot = OnReboot;
1012 if (!bypassLaunchResult)
1013 {
1014 var launchResult = await LaunchResult.WaitAsync(cancellationToken);
1015 if (launchResult.ExitCode.HasValue)
1016 {
1017 Logger.LogDebug("Not sending topic request {commandType} to server that failed to launch!", parameters.CommandType);
1018 return null;
1019 }
1020 }
1021
1022 // meh, this is kind of a hack, but it works
1024 {
1025 Logger.LogDebug("Not sending topic request {commandType} to server that is rebooting/starting.", parameters.CommandType);
1026 return null;
1027 }
1028
1029 using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
1030 var combinedCancellationToken = cts.Token;
1031 async ValueTask CancelIfLifetimeElapses()
1032 {
1033 try
1034 {
1035 var completed = await Task.WhenAny(Lifetime, reboot).WaitAsync(combinedCancellationToken);
1036
1037 Logger.LogDebug(
1038 "Server {action}, cancelling pending command: {commandType}",
1039 completed != reboot
1040 ? "process ended"
1041 : "rebooting",
1042 parameters.CommandType);
1043 cts.Cancel();
1044 }
1045 catch (OperationCanceledException)
1046 {
1047 // expected, not even worth tracing
1048 }
1049 catch (Exception ex)
1050 {
1051 Logger.LogError(ex, "Error in CancelIfLifetimeElapses!");
1052 }
1053 }
1054
1055 TopicResponse? fullResponse = null;
1056 var lifetimeWatchingTask = CancelIfLifetimeElapses();
1057 try
1058 {
1059 var combinedResponse = await SendTopicRequest(parameters, combinedCancellationToken);
1060
1061 void LogCombinedResponse()
1062 {
1063 if (LogTopicRequests && combinedResponse != null)
1064 Logger.LogTrace("Topic response: {topicString}", combinedResponse.ByondTopicResponse.StringData ?? "(NO STRING DATA)");
1065 }
1066
1067 LogCombinedResponse();
1068
1069 if (combinedResponse?.InteropResponse?.Chunk != null)
1070 {
1071 Logger.LogTrace("Topic response is chunked...");
1072
1073 ChunkData? nextChunk = combinedResponse.InteropResponse.Chunk;
1074 do
1075 {
1076 var nextRequest = await ProcessChunk<TopicResponse, ChunkedTopicParameters>(
1077 (completedResponse, _) =>
1078 {
1079 fullResponse = completedResponse;
1080 return ValueTask.FromResult<ChunkedTopicParameters?>(null);
1081 },
1082 error =>
1083 {
1084 Logger.LogWarning("Topic response chunking error: {message}", error);
1085 return null;
1086 },
1087 combinedResponse?.InteropResponse?.Chunk,
1088 combinedCancellationToken);
1089
1090 if (nextRequest != null)
1091 {
1092 nextRequest.PayloadId = nextChunk.PayloadId;
1093 combinedResponse = await SendTopicRequest(nextRequest, combinedCancellationToken);
1094 LogCombinedResponse();
1095 nextChunk = combinedResponse?.InteropResponse?.Chunk;
1096 }
1097 else
1098 nextChunk = null;
1099 }
1100 while (nextChunk != null);
1101 }
1102 else
1103 fullResponse = combinedResponse?.InteropResponse;
1104 }
1105 catch (OperationCanceledException ex)
1106 {
1107 Logger.LogDebug(
1108 ex,
1109 "Topic request {cancellationType}!",
1110 combinedCancellationToken.IsCancellationRequested
1111 ? cancellationToken.IsCancellationRequested
1112 ? "cancelled"
1113 : "aborted"
1114 : "timed out");
1115
1116 // throw only if the original token was the trigger
1117 cancellationToken.ThrowIfCancellationRequested();
1118 }
1119 finally
1120 {
1121 cts.Cancel();
1122 await lifetimeWatchingTask;
1123 }
1124
1125 if (fullResponse?.ErrorMessage != null)
1126 Logger.LogWarning(
1127 "Errored topic response for command {commandType}: {errorMessage}",
1128 parameters.CommandType,
1129 fullResponse.ErrorMessage);
1130
1131 return fullResponse;
1132 }
1133
1140 {
1141 if (invocation == null)
1142 return BridgeError("Missing eventInvocation!");
1143
1144 var eventName = invocation.EventName;
1145 if (eventName == null)
1146 return BridgeError("Missing eventName!");
1147
1148 var notifyCompletion = invocation.NotifyCompletion;
1149 if (!notifyCompletion.HasValue)
1150 return BridgeError("Missing notifyCompletion!");
1151
1152 var eventParams = new List<string>
1153 {
1155 };
1156
1157 eventParams.AddRange(invocation
1158 .Parameters?
1159 .Where(param => param != null)
1160 .Cast<string>()
1161 ?? Enumerable.Empty<string>());
1162
1163 var eventId = Guid.NewGuid();
1164 Logger.LogInformation("Triggering custom event \"{eventName}\": {eventId}", eventName, eventId);
1165
1166 var cancellationToken = sessionDurationCts.Token;
1167 ValueTask? eventTask = eventConsumer.HandleCustomEvent(eventName, eventParams, cancellationToken);
1168
1169 async Task ProcessEvent()
1170 {
1171 try
1172 {
1173 await eventTask.Value;
1174
1175 if (notifyCompletion.Value)
1176 await SendCommand(
1177 new TopicParameters(eventId),
1178 cancellationToken);
1179 else
1180 Logger.LogTrace("Finished custom event {eventId}, not sending notification.", eventId);
1181 }
1182 catch (OperationCanceledException ex)
1183 {
1184 Logger.LogDebug(ex, "Custom event invocation {eventId} aborted!", eventId);
1185 }
1186 catch (Exception ex)
1187 {
1188 Logger.LogWarning(ex, "Custom event invocation {eventId} errored!", eventId);
1189 }
1190 }
1191
1192 if (!eventTask.HasValue)
1193 return BridgeError("Event refused to execute due to matching a TGS event!");
1194
1195 lock (sessionDurationCts)
1196 {
1197 var previousEventProcessingTask = customEventProcessingTask;
1198 var eventProcessingTask = ProcessEvent();
1199 customEventProcessingTask = Task.WhenAll(customEventProcessingTask, eventProcessingTask);
1200 }
1201
1202 return new BridgeResponse
1203 {
1204 EventId = notifyCompletion.Value
1205 ? eventId.ToString()
1206 : null,
1207 };
1208 }
1209 }
1210}
Information about an engine installation.
Metadata about a server instance.
Definition Instance.cs:9
virtual ? Version DMApiVersion
The DMAPI Version.
Definition CompileJob.cs:41
ChatMessage? ChatMessage
The Interop.ChatMessage for BridgeCommandType.ChatSend requests.
CustomEventInvocation? EventInvocation
The Bridge.CustomEventInvocation being triggered.
ushort? TopicPort
The port that should be used to send world topics, if not the default.
DreamDaemonSecurity? MinimumSecurityLevel
The minimum required DreamDaemonSecurity level for BridgeCommandType.Startup requests.
ChunkData? Chunk
The ChunkData for BridgeCommandType.Chunk requests.
Version? Version
The DMAPI global::System.Version for BridgeCommandType.Startup requests.
bool? NotifyCompletion
If the DMAPI should be notified when the event compeletes.
Representation of the initial data passed as part of a BridgeCommandType.Startup request.
Version ServerVersion
The IAssemblyInformationProvider.Version.
DreamDaemonVisibility Visibility
The DreamDaemonSecurity level of the launch.
string InstanceName
The NamedEntity.Name of the owner at the time of launch.
bool ApiValidateOnly
If DD should just respond if it's API is working and then exit.
DreamDaemonSecurity SecurityLevel
The DreamDaemonSecurity level of the launch.
ICollection< string >? ChannelIds
The ICollection<T> of Chat.ChannelRepresentation.Ids to sent the MessageContent to....
Represents an update of ChannelRepresentations.
Definition ChatUpdate.cs:13
A packet of a split serialized set of data.
Definition ChunkData.cs:7
uint? PayloadId
The ID of the full request to differentiate different chunkings.Nullable to prevent default value omi...
Class that deserializes chunked interop payloads.
Definition Chunker.cs:16
ILogger< Chunker > Logger
The ILogger for the Chunker.
Definition Chunker.cs:20
uint NextPayloadId
Gets a payload ID for use in a new ChunkSetInfo.
Definition Chunker.cs:26
Constants used for communication with the DMAPI.
static readonly JsonSerializerSettings SerializerSettings
JsonSerializerSettings for use when communicating with the DMAPI.
const uint MaximumTopicRequestLength
The maximum length in bytes of a Byond.TopicSender.ITopicClient payload.
static readonly Version InteropVersion
The DMAPI InteropVersion being used.
const string TopicData
Parameter json is encoded in for topic requests.
string AccessIdentifier
Used to identify and authenticate the DreamDaemon instance.
string? ErrorMessage
Any errors in the client's parameters.
static TopicParameters CreateInstanceRenamedTopicParameters(string newInstanceName)
Initializes a new instance of the TopicParameters class.
ChunkData? Chunk
The ChunkData for a partial request.
bool IsPriority
Whether or not the TopicParameters constitute a priority request.
Combines a Byond.TopicSender.TopicResponse with a TopicResponse.
TopicResponse? InteropResponse
The interop TopicResponse, if any.
Represents the result of trying to start a DD process.
Parameters necessary for duplicating a ISessionController session.
RuntimeInformation? RuntimeInformation
The Interop.Bridge.RuntimeInformation for the DMAPI.
IDmbProvider? InitialDmb
The IDmbProvider initially used to launch DreamDaemon. Should be a different IDmbProvider than Dmb....
IDmbProvider Dmb
The IDmbProvider used by DreamDaemon.
async ValueTask InstanceRenamed(string newInstanceName, CancellationToken cancellationToken)
Called when the owning Instance is renamed.A ValueTask representing the running operation.
ValueTask< TopicResponse?> SendCommand(TopicParameters parameters, CancellationToken cancellationToken)
Sends a command to DreamDaemon through /world/Topic().A ValueTask<TResult> resulting in the TopicResp...
async ValueTask< BridgeResponse?> ProcessBridgeCommand(BridgeParameters parameters, CancellationToken cancellationToken)
Handle a set of bridge parameters .
readonly Byond.TopicSender.ITopicClient byondTopicSender
The Byond.TopicSender.ITopicClient for the SessionController.
string DumpFileExtension
The file extension to use for process dumps created from this session.
ApiValidationStatus apiValidationStatus
The ApiValidationStatus for the SessionController.
readonly object synchronizationLock
lock object for port updates and disposed.
DateTimeOffset? LaunchTime
When the process was started.
void AdjustPriority(bool higher)
Set's the owned global::System.Diagnostics.Process.PriorityClass to a non-normal value.
async ValueTask< TopicResponse?> SendCommand(TopicParameters parameters, bool bypassLaunchResult, CancellationToken cancellationToken)
Sends a command to DreamDaemon through /world/Topic().
volatile uint rebootBridgeRequestsProcessing
The number of currently active calls to ProcessBridgeRequest(BridgeParameters, CancellationToken) fro...
readonly IDotnetDumpService dotnetDumpService
The IDotnetDumpService for the SessionController.
readonly Api.Models.Instance metadata
The Instance metadata.
Task OnReboot
A Task that completes when the server calls /world/TgsReboot().
bool DMApiAvailable
If the DMAPI may be used this session.
readonly TaskCompletionSource initialBridgeRequestTcs
The TaskCompletionSource that completes when DD makes it's first bridge request.
bool terminationWasIntentional
Backing field for overriding TerminationWasIntentional.
bool disposed
If the SessionController has been disposed.
void ResetRebootState()
Changes RebootState to RebootState.Normal without telling the DMAPI.
readonly IEngineExecutableLock engineLock
The IEngineExecutableLock for the SessionController.
readonly IAsyncDelayer asyncDelayer
The IAsyncDelayer for the SessionController.
async ValueTask< BridgeResponse?> ProcessBridgeRequest(BridgeParameters parameters, CancellationToken cancellationToken)
Handle a set of bridge parameters .A ValueTask<TResult> resulting in the BridgeResponse for the reque...
ReattachInformation ReattachInformation
The up to date Session.ReattachInformation.
volatile Task rebootGate
Backing field for RebootGate.
bool TerminationWasIntentional
If the DreamDaemon instance sent a.
readonly IEventConsumer eventConsumer
The IEventConsumer for the SessionController.
BridgeResponse BridgeError(string message)
Log and return a BridgeResponse for a given message .
bool released
If process should be kept alive instead.
readonly IChatManager chat
The IChatManager for the SessionController.
readonly IChatTrackingContext chatTrackingContext
The IChatTrackingContext for the SessionController.
Task< int?> Lifetime
The Task<TResult> resulting in the exit code of the process or null if the process was detached.
readonly CancellationTokenSource sessionDurationCts
A CancellationTokenSource used for tasks that should not exceed the lifetime of the session.
ValueTask CreateDump(string outputFile, bool minidump, CancellationToken cancellationToken)
Create a dump file of the process.A ValueTask representing the running operation.
volatile TaskCompletionSource startupTcs
The TaskCompletionSource that completes when DD sends a valid startup bridge request.
async Task PostValidationShutdown(Task< bool > proceedTask)
Terminates the server after ten seconds if it does not exit.
SessionController(ReattachInformation reattachInformation, Api.Models.Instance metadata, IProcess process, IEngineExecutableLock engineLock, Byond.TopicSender.ITopicClient byondTopicSender, IChatTrackingContext chatTrackingContext, IBridgeRegistrar bridgeRegistrar, IChatManager chat, IAssemblyInformationProvider assemblyInformationProvider, IAsyncDelayer asyncDelayer, IDotnetDumpService dotnetDumpService, IEventConsumer eventConsumer, ILogger< SessionController > logger, Func< ValueTask > postLifetimeCallback, uint? startupTimeout, bool reattached, bool apiValidate)
Initializes a new instance of the SessionController class.
readonly bool apiValidationSession
If this session is meant to validate the presence of the DMAPI.
FifoSemaphore TopicSendSemaphore
The FifoSemaphore used to prevent concurrent calls into /world/Topic().
async ValueTask UpdateChannels(IEnumerable< ChannelRepresentation > newChannels, CancellationToken cancellationToken)
Called when newChannels are set.A ValueTask representing the running operation.
string GenerateQueryString(TopicParameters parameters, out string json)
Generates a Byond.TopicSender.ITopicClient query string for a given set of parameters .
void CheckDisposed()
Throws an ObjectDisposedException if DisposeAsync has been called.
readonly? IBridgeRegistration bridgeRegistration
The IBridgeRegistration for the SessionController.
long? MemoryUsage
Gets the process' memory usage in bytes.
async Task< LaunchResult > GetLaunchResult(IAssemblyInformationProvider assemblyInformationProvider, IAsyncDelayer asyncDelayer, uint? startupTimeout, bool reattached, bool apiValidate)
The Task<TResult> for LaunchResult.
IAsyncDisposable ReplaceDmbProvider(IDmbProvider dmbProvider)
Replace the IDmbProvider in use with a given newProvider , disposing the old one.An IAsyncDisposable ...
Task OnStartup
A Task that completes when the server calls /world/TgsNew().
bool ProcessingRebootBridgeRequest
If the ISessionController is currently processing a bridge request from TgsReboot().
ValueTask Release()
Releases the IProcess without terminating it. Also calls IDisposable.Dispose.A ValueTask representing...
volatile? Task postValidationShutdownTask
Task for shutting down the server if it is taking too long after validation.
volatile Task customEventProcessingTask
The Task representing calls to TriggerCustomEvent(CustomEventInvocation?).
Task OnPrime
A Task that completes when the server calls /world/TgsInitializationComplete().
async ValueTask< CombinedTopicResponse?> SendTopicRequest(TopicParameters parameters, CancellationToken cancellationToken)
Send a topic request for given parameters to DreamDaemon, chunking it if necessary.
volatile TaskCompletionSource rebootTcs
The TaskCompletionSource that completes when DD tells us about a reboot.
async ValueTask< CombinedTopicResponse?> SendRawTopic(string queryString, bool priority, CancellationToken cancellationToken)
Send a given queryString to DreamDaemon's /world/Topic.
volatile TaskCompletionSource primeTcs
The TaskCompletionSource that completes when DD tells us it's primed.
async ValueTask< bool > SetRebootState(RebootState newRebootState, CancellationToken cancellationToken)
Attempts to change the current RebootState to newRebootState .A ValueTask<TResult> resulting in true ...
Task RebootGate
A Task that must complete before a TgsReboot() bridge request can complete.
readonly IProcess process
The IProcess for the SessionController.
BridgeResponse TriggerCustomEvent(CustomEventInvocation? invocation)
Trigger a custom event from a given invocation .
RebootState RebootState
The current DreamDaemon reboot state.
ushort Port
The port the game server was last listening on.
A first-in first-out async semaphore.
async ValueTask< SemaphoreSlimContext > Lock(CancellationToken cancellationToken)
Locks the FifoSemaphore.
Helpers for manipulating the Serilog.Context.LogContext.
const string InstanceIdContextProperty
The Serilog.Context.LogContext property name for Models.Instance Api.Models.EntityId....
Notifyee of when ChannelRepresentations in a IChatTrackingContext are updated.
For managing connected chat services.
void QueueMessage(MessageContent message, IEnumerable< ulong > channelIds)
Queue a chat message to a given set of channelIds .
Represents a tracking of dynamic chat json files.
void SetChannelSink(IChannelSink channelSink)
Sets the channelSink for the IChatTrackingContext.
Provides absolute paths to the latest compiled .dmbs.
void KeepAlive()
Disposing the IDmbProvider won't cause a cleanup of the working directory.
EngineVersion EngineVersion
The Api.Models.EngineVersion used to build the .dmb.
Models.CompileJob CompileJob
The CompileJob of the .dmb.
Represents usage of the two primary BYOND server executables.
void DoNotDeleteThisSession()
Call if, during a detach, this version should not be deleted.
bool UseDotnetDump
If dotnet-dump should be used to create process dumps for this installation.
ValueTask StopServerProcess(ILogger logger, IProcess process, string accessIdentifier, ushort port, CancellationToken cancellationToken)
Kills a given engine server process .
Consumes EventTypes and takes the appropriate actions.
ValueTask? HandleCustomEvent(string eventName, IEnumerable< string?> parameters, CancellationToken cancellationToken)
Handles a given custom event.
IBridgeRegistration RegisterHandler(IBridgeHandler bridgeHandler)
Register a given bridgeHandler .
Handles communication with a DreamDaemon IProcess.
Service for managing the dotnet-dump installation.
ValueTask Dump(IProcess process, string outputFile, bool minidump, CancellationToken cancellationToken)
Attempt to dump a given process .
void SuspendProcess()
Suspends the process.
void ResumeProcess()
Resumes the process.
void AdjustPriority(bool higher)
Set's the owned global::System.Diagnostics.Process.PriorityClass to a non-normal value.
long? MemoryUsage
Gets the process' memory usage in bytes.
DateTimeOffset? LaunchTime
When the process was started.
ValueTask CreateDump(string outputFile, bool minidump, CancellationToken cancellationToken)
Create a dump file of the process.
Task< int?> Lifetime
The Task<TResult> resulting in the exit code of the process or null if the process was detached.
Abstraction over a global::System.Diagnostics.Process.
Definition IProcess.cs:11
Task Startup
The Task representing the time until the IProcess becomes "idle".
Definition IProcess.cs:20
void Terminate()
Asycnhronously terminates the process.
Task Delay(TimeSpan timeSpan, CancellationToken cancellationToken)
Create a Task that completes after a given timeSpan .
DreamDaemonSecurity
DreamDaemon's security level.
EngineType
The type of engine the codebase is using.
Definition EngineType.cs:7
@ Byond
Build your own net dream.
BridgeCommandType
Represents the BridgeParameters.CommandType.
@ Chunk
DreamDaemon attempting to send a longer bridge message.
RebootState
Represents the action to take when /world/Reboot() is called.
Definition RebootState.cs:7
ApiValidationStatus
Status of DMAPI validation.