From c9ae5e3bd2db78cf5180cc2a013aedf10acd1da5 Mon Sep 17 00:00:00 2001 From: Leath Cooper Date: Sat, 19 Oct 2024 16:10:11 -0400 Subject: [PATCH 1/2] NATS --- Projects/Server/Events/EventSink.cs | 3 + Projects/Server/Events/MessageBus.cs | 97 +++++++++++++++++++++++++ Projects/Server/Items/Item.cs | 6 ++ Projects/Server/Main.cs | 21 ++++-- Projects/Server/Server.csproj | 1 + Projects/UOContent/Commands/Handlers.cs | 6 +- Projects/UOContent/Misc/Broadcasts.cs | 9 +++ 7 files changed, 135 insertions(+), 8 deletions(-) create mode 100644 Projects/Server/Events/MessageBus.cs diff --git a/Projects/Server/Events/EventSink.cs b/Projects/Server/Events/EventSink.cs index 6a72185f4..8b9f83533 100644 --- a/Projects/Server/Events/EventSink.cs +++ b/Projects/Server/Events/EventSink.cs @@ -41,4 +41,7 @@ public static partial class EventSink public static event Action ServerStarted; public static void InvokeServerStarted() => ServerStarted?.Invoke(); + + public static event Action IncomingMessage; + public static void InvokeIncomingMessage(string message, int hue) => IncomingMessage?.Invoke(message, hue); } diff --git a/Projects/Server/Events/MessageBus.cs b/Projects/Server/Events/MessageBus.cs new file mode 100644 index 000000000..5a30d1b49 --- /dev/null +++ b/Projects/Server/Events/MessageBus.cs @@ -0,0 +1,97 @@ +using System; +using System.Text.Json.Serialization; +using System.Threading; +using NATS.Client; +using NATS.Client.Core; +using NATS.Client.JetStream; +using NATS.Client.JetStream.Models; +using NATS.Net; + +namespace Server; + +[JsonSerializable(typeof(BroadcastMessage))] +internal partial class BroadcastJsonContext : JsonSerializerContext; + +public record BroadcastMessage +{ + [JsonPropertyName("hue")] + public int Hue { get; set; } = 0x482; + + [JsonPropertyName("message")] + public string? Message { get; set; } +} + +public static class MessageBus +{ + private static NatsClient natsClient; + private static INatsJSContext jetstreamContext; + private static INatsJSConsumer incomingBroadcastsConsumer; + private static CancellationTokenSource cancellationTokenSource; + + private static NatsJsonContextSerializer broadcastSerializer = + new (BroadcastJsonContext.Default); + + public static void Start() + { + natsClient = new NatsClient(); + jetstreamContext = natsClient.CreateJetStreamContext(); + + StartBroadcasts(); + } + + private static async void StartBroadcasts() + { + await jetstreamContext.CreateStreamAsync( + new StreamConfig(name: "INCOMING_NOTIFICATIONS", subjects: ["incoming.notifications.>"]) + ); + + incomingBroadcastsConsumer = await jetstreamContext.CreateOrUpdateConsumerAsync( + "INCOMING_NOTIFICATIONS", + new ConsumerConfig("incoming_broadcasts_consumer") + { + FilterSubject = "incoming.notifications.broadcast" + } + ); + + cancellationTokenSource = new CancellationTokenSource(); + + await foreach (var msg in incomingBroadcastsConsumer.ConsumeAsync( + cancellationToken: cancellationTokenSource.Token, + serializer: broadcastSerializer + )) + { + if (msg.Data == null) + { + Console.WriteLine($"Message Received w/ null Data"); + continue; + } + + Console.WriteLine($"Message Received {msg.Data}"); + EventSink.InvokeIncomingMessage(msg.Data.Message, msg.Data.Hue); + await msg.AckAsync(cancellationToken: cancellationTokenSource.Token); + } + } + + public static async void SendBroadcast(string message, int hue) + { + await natsClient.PublishAsync( + subject: "incoming.notifications.broadcast", + data: new BroadcastMessage() + { + Message = message, + Hue = hue, + }, + serializer: broadcastSerializer + ); + } + + public static async void Publish(string subject, T data, INatsSerializer serializer = null) + { + await natsClient.PublishAsync(subject, data); + } + + public static void Kill() + { + cancellationTokenSource.Cancel(); + } +} diff --git a/Projects/Server/Items/Item.cs b/Projects/Server/Items/Item.cs index aeab60811..df8048955 100644 --- a/Projects/Server/Items/Item.cs +++ b/Projects/Server/Items/Item.cs @@ -233,6 +233,8 @@ public class Item : IHued, IComparable, ISpawnable, IObjectPropertyListEnt SetLastMoved(); World.AddEntity(this); + + MessageBus.Publish("items.created", $"{Serial}:{itemID}"); } public Item(Serial serial) => Serial = serial; @@ -1072,6 +1074,8 @@ public class Item : IHued, IComparable, ISpawnable, IObjectPropertyListEnt { writer.WriteEncodedInt(info.m_SavedFlags); } + + MessageBus.Publish("items.serialize", $"{Serial}:{ItemID}"); } public void MoveToWorld(WorldLocation worldLocation) @@ -3036,6 +3040,8 @@ public class Item : IHued, IComparable, ISpawnable, IObjectPropertyListEnt // if (version < 9) VerifyCompactInfo(); + + MessageBus.Publish("items.deserialize", $"{Serial}:{ItemID}"); } private void FixHolding_Sandbox() diff --git a/Projects/Server/Main.cs b/Projects/Server/Main.cs index 13577440d..ab19073ff 100644 --- a/Projects/Server/Main.cs +++ b/Projects/Server/Main.cs @@ -189,9 +189,11 @@ public static class Core if (fi.Directory != null && Directory.Exists(fi.Directory.FullName)) { fullPath = fi.Directory.EnumerateFiles( - fi.Name, - new EnumerationOptions { MatchCasing = MatchCasing.CaseInsensitive } - ).FirstOrDefault()?.FullName; + fi.Name, + new EnumerationOptions { MatchCasing = MatchCasing.CaseInsensitive } + ) + .FirstOrDefault() + ?.FullName; } } @@ -317,6 +319,7 @@ public static class Core process.Start(); } + logger.Information("Restart done"); } catch (Exception e) @@ -380,13 +383,15 @@ public static class Core Utility.PopColor(); Utility.PushColor(ConsoleColor.DarkGray); - Console.WriteLine(@"Copyright 2019-2023 ModernUO Development Team + Console.WriteLine( + @"Copyright 2019-2023 ModernUO Development Team This program comes with ABSOLUTELY NO WARRANTY; This is free software, and you are welcome to redistribute it under certain conditions. You should have received a copy of the GNU General Public License along with this program. If not, see . - ".TrimMultiline()); + ".TrimMultiline() + ); Utility.PopColor(); Console.CancelKeyPress += Console_CancelKeyPressed; @@ -426,6 +431,8 @@ public static class Core TileMatrixLoader.LoadTileMatrix(); + MessageBus.Start(); + RegionJsonSerializer.LoadRegions(); World.Load(); @@ -433,6 +440,7 @@ public static class Core TcpServer.Start(); PingServer.Start(); + EventSink.InvokeServerStarted(); RunEventLoop(); } @@ -561,7 +569,8 @@ public static class Core if (World.DirtyTrackingEnabled) { var manualDirtyCheckingAttribute = type.GetCustomAttribute(false); - var codeGennedAttribute = type.GetCustomAttribute(false); + var codeGennedAttribute = + type.GetCustomAttribute(false); if (manualDirtyCheckingAttribute == null && codeGennedAttribute == null) { diff --git a/Projects/Server/Server.csproj b/Projects/Server/Server.csproj index a47b3fb1a..f05fa2864 100644 --- a/Projects/Server/Server.csproj +++ b/Projects/Server/Server.csproj @@ -35,6 +35,7 @@ + diff --git a/Projects/UOContent/Commands/Handlers.cs b/Projects/UOContent/Commands/Handlers.cs index de1e831e2..f3e6cd86a 100644 --- a/Projects/UOContent/Commands/Handlers.cs +++ b/Projects/UOContent/Commands/Handlers.cs @@ -638,8 +638,10 @@ namespace Server.Commands [Description("Broadcasts a message to everyone online.")] public static void BroadcastMessage_OnCommand(CommandEventArgs e) { - BroadcastMessage(AccessLevel.Player, 0x482, $"Staff message from {e.Mobile.Name}:"); - BroadcastMessage(AccessLevel.Player, 0x482, e.ArgString); + MessageBus.SendBroadcast($"Staff message from {e.Mobile.Name}:", 0x482); + MessageBus.SendBroadcast($"{e.ArgString}", 0x21); + // BroadcastMessage(AccessLevel.Player, 0x482, $"Staff message from {e.Mobile.Name}:"); + // BroadcastMessage(AccessLevel.Player, 0x482, e.ArgString); } public static void BroadcastMessage(AccessLevel ac, int hue, string message) diff --git a/Projects/UOContent/Misc/Broadcasts.cs b/Projects/UOContent/Misc/Broadcasts.cs index 74eb2ad4f..90e0c2cbc 100644 --- a/Projects/UOContent/Misc/Broadcasts.cs +++ b/Projects/UOContent/Misc/Broadcasts.cs @@ -1,3 +1,5 @@ +using System; + namespace Server.Misc { public static class Broadcasts @@ -6,6 +8,13 @@ namespace Server.Misc { EventSink.ServerCrashed += EventSink_Crashed; EventSink.Shutdown += EventSink_Shutdown; + EventSink.IncomingMessage += EventSink_IncomingMessage; + } + + public static void EventSink_IncomingMessage(string message, int hue) + { + Console.WriteLine("Incoming Message"); + World.Broadcast(hue, true, message); } public static void EventSink_Crashed(ServerCrashedEventArgs e) From c06c84ca9ccac6dc92041973cb7b2c2a9584dca6 Mon Sep 17 00:00:00 2001 From: Leath Cooper Date: Sat, 19 Oct 2024 17:52:53 -0400 Subject: [PATCH 2/2] NATS driven broadcast --- Projects/Server/Events/MessageBus.cs | 83 +--------------- Projects/UOContent/Commands/Broadcast.cs | 119 +++++++++++++++++++++++ Projects/UOContent/Commands/Handlers.cs | 34 ------- Projects/UOContent/UOContent.csproj | 1 + 4 files changed, 123 insertions(+), 114 deletions(-) create mode 100644 Projects/UOContent/Commands/Broadcast.cs diff --git a/Projects/Server/Events/MessageBus.cs b/Projects/Server/Events/MessageBus.cs index 5a30d1b49..18a1cf0a7 100644 --- a/Projects/Server/Events/MessageBus.cs +++ b/Projects/Server/Events/MessageBus.cs @@ -1,97 +1,20 @@ -using System; -using System.Text.Json.Serialization; -using System.Threading; -using NATS.Client; using NATS.Client.Core; using NATS.Client.JetStream; -using NATS.Client.JetStream.Models; using NATS.Net; namespace Server; -[JsonSerializable(typeof(BroadcastMessage))] -internal partial class BroadcastJsonContext : JsonSerializerContext; - -public record BroadcastMessage -{ - [JsonPropertyName("hue")] - public int Hue { get; set; } = 0x482; - - [JsonPropertyName("message")] - public string? Message { get; set; } -} - public static class MessageBus { private static NatsClient natsClient; - private static INatsJSContext jetstreamContext; - private static INatsJSConsumer incomingBroadcastsConsumer; - private static CancellationTokenSource cancellationTokenSource; + public static NatsClient Client => natsClient; - private static NatsJsonContextSerializer broadcastSerializer = - new (BroadcastJsonContext.Default); + private static INatsJSContext jetstreamContext; + public static INatsJSContext Context => jetstreamContext; public static void Start() { natsClient = new NatsClient(); jetstreamContext = natsClient.CreateJetStreamContext(); - - StartBroadcasts(); - } - - private static async void StartBroadcasts() - { - await jetstreamContext.CreateStreamAsync( - new StreamConfig(name: "INCOMING_NOTIFICATIONS", subjects: ["incoming.notifications.>"]) - ); - - incomingBroadcastsConsumer = await jetstreamContext.CreateOrUpdateConsumerAsync( - "INCOMING_NOTIFICATIONS", - new ConsumerConfig("incoming_broadcasts_consumer") - { - FilterSubject = "incoming.notifications.broadcast" - } - ); - - cancellationTokenSource = new CancellationTokenSource(); - - await foreach (var msg in incomingBroadcastsConsumer.ConsumeAsync( - cancellationToken: cancellationTokenSource.Token, - serializer: broadcastSerializer - )) - { - if (msg.Data == null) - { - Console.WriteLine($"Message Received w/ null Data"); - continue; - } - - Console.WriteLine($"Message Received {msg.Data}"); - EventSink.InvokeIncomingMessage(msg.Data.Message, msg.Data.Hue); - await msg.AckAsync(cancellationToken: cancellationTokenSource.Token); - } - } - - public static async void SendBroadcast(string message, int hue) - { - await natsClient.PublishAsync( - subject: "incoming.notifications.broadcast", - data: new BroadcastMessage() - { - Message = message, - Hue = hue, - }, - serializer: broadcastSerializer - ); - } - - public static async void Publish(string subject, T data, INatsSerializer serializer = null) - { - await natsClient.PublishAsync(subject, data); - } - - public static void Kill() - { - cancellationTokenSource.Cancel(); } } diff --git a/Projects/UOContent/Commands/Broadcast.cs b/Projects/UOContent/Commands/Broadcast.cs new file mode 100644 index 000000000..7f8fff4bc --- /dev/null +++ b/Projects/UOContent/Commands/Broadcast.cs @@ -0,0 +1,119 @@ +using System; +using System.Text.Json.Serialization; +using System.Threading; +using Server.Network; +using NATS.Client; +using NATS.Client.Core; +using NATS.Client.JetStream; +using NATS.Client.JetStream.Models; +using NATS.Net; + +namespace Server.Commands; + +[JsonSerializable(typeof(BroadcastPayload))] +internal partial class BroadcastJsonContext : JsonSerializerContext; + +public record BroadcastPayload +{ + [JsonPropertyName("hue")] + public int Hue { get; set; } = 0x482; + + [JsonPropertyName("message")] + public string? Message { get; set; } +} + +public static class Broadcast +{ + private static INatsJSConsumer incomingBroadcastsConsumer; + + private static readonly NatsJsonContextSerializer broadcastSerializer = + new(BroadcastJsonContext.Default); + + public static void Configure() + { + CommandSystem.Register("BCast", AccessLevel.GameMaster, BroadcastMessage_OnCommand); + CommandSystem.Register("SMsg", AccessLevel.Counselor, StaffMessage_OnCommand); + } + + public static void Initialize() + { + StartBroadcasts(); + } + + [Usage("BCast ")] + [Aliases("B", "BC")] + [Description("Broadcasts a message to everyone online.")] + public static void BroadcastMessage_OnCommand(CommandEventArgs e) + { + EmitBroadcast($"Staff message from {e.Mobile.Name}:", 0x482); + EmitBroadcast($"{e.ArgString}", 0x21); + // BroadcastMessage(AccessLevel.Player, 0x482, $"Staff message from {e.Mobile.Name}:"); + // BroadcastMessage(AccessLevel.Player, 0x482, e.ArgString); + } + + [Usage("SMsg ")] + [Aliases("S", "SM")] + [Description("Broadcasts a message to all online staff.")] + public static void StaffMessage_OnCommand(CommandEventArgs e) + { + BroadcastMessage(AccessLevel.Counselor, e.Mobile.SpeechHue, $"[{e.Mobile.Name}] {e.ArgString}"); + } + + public static async void EmitBroadcast(string message, int hue) + { + await MessageBus.Client.PublishAsync( + subject: "broadcast.all", + data: new BroadcastPayload() + { + Message = message, + Hue = hue, + }, + serializer: broadcastSerializer + ); + } + + public static void BroadcastMessage(AccessLevel ac, int hue, string message) + { + foreach (var state in NetState.Instances) + { + var m = state.Mobile; + + if (m?.AccessLevel >= ac) + { + m.SendMessage(hue, message); + } + } + } + + private static async void StartBroadcasts() + { + await MessageBus.Context.CreateStreamAsync( + new StreamConfig(name: "BROADCASTS", subjects: ["broadcast.>"]) + ); + + incomingBroadcastsConsumer = await MessageBus.Context.CreateOrUpdateConsumerAsync( + "BROADCASTS", + new ConsumerConfig("broadcast_incoming_consumer") + ); + + CancellationTokenSource cancellationTokenSource = new(); + + await foreach (var msg in incomingBroadcastsConsumer.ConsumeAsync( + serializer: broadcastSerializer, + cancellationToken: cancellationTokenSource.Token + )) + { + if (msg.Data == null) + { + Console.WriteLine($"Message Received w/ null Data"); + continue; + } + + Console.WriteLine($"Message Received {msg.Data}"); + EventSink.InvokeIncomingMessage(msg.Data.Message, msg.Data.Hue); + await msg.AckAsync( + cancellationToken: cancellationTokenSource.Token + ); + } + } +} diff --git a/Projects/UOContent/Commands/Handlers.cs b/Projects/UOContent/Commands/Handlers.cs index f3e6cd86a..0cd63aa95 100644 --- a/Projects/UOContent/Commands/Handlers.cs +++ b/Projects/UOContent/Commands/Handlers.cs @@ -33,8 +33,6 @@ namespace Server.Commands Register("Help", AccessLevel.Player, Help_OnCommand); Register("Move", AccessLevel.GameMaster, Move_OnCommand); Register("Client", AccessLevel.Counselor, Client_OnCommand); - Register("SMsg", AccessLevel.Counselor, StaffMessage_OnCommand); - Register("BCast", AccessLevel.GameMaster, BroadcastMessage_OnCommand); Register("Bank", AccessLevel.GameMaster, Bank_OnCommand); Register("Echo", AccessLevel.Counselor, Echo_OnCommand); Register("Sound", AccessLevel.GameMaster, Sound_OnCommand); @@ -625,38 +623,6 @@ namespace Server.Commands } } - [Usage("SMsg ")] - [Aliases("S", "SM")] - [Description("Broadcasts a message to all online staff.")] - public static void StaffMessage_OnCommand(CommandEventArgs e) - { - BroadcastMessage(AccessLevel.Counselor, e.Mobile.SpeechHue, $"[{e.Mobile.Name}] {e.ArgString}"); - } - - [Usage("BCast ")] - [Aliases("B", "BC")] - [Description("Broadcasts a message to everyone online.")] - public static void BroadcastMessage_OnCommand(CommandEventArgs e) - { - MessageBus.SendBroadcast($"Staff message from {e.Mobile.Name}:", 0x482); - MessageBus.SendBroadcast($"{e.ArgString}", 0x21); - // BroadcastMessage(AccessLevel.Player, 0x482, $"Staff message from {e.Mobile.Name}:"); - // BroadcastMessage(AccessLevel.Player, 0x482, e.ArgString); - } - - public static void BroadcastMessage(AccessLevel ac, int hue, string message) - { - foreach (var state in NetState.Instances) - { - var m = state.Mobile; - - if (m?.AccessLevel >= ac) - { - m.SendMessage(hue, message); - } - } - } - [Usage("AutoPageNotify")] [Aliases("APN")] [Description("Toggles your auto-page-notify status.")] diff --git a/Projects/UOContent/UOContent.csproj b/Projects/UOContent/UOContent.csproj index 5f07339eb..c1f692f38 100644 --- a/Projects/UOContent/UOContent.csproj +++ b/Projects/UOContent/UOContent.csproj @@ -34,6 +34,7 @@ + false