From 77b12626013b04e0db30d0b2e093cf477114b8e6 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 17:35:37 -0400 Subject: [PATCH 01/12] Don't require WorkerManager to use DataContext --- Refresh.GameServer/RefreshGameServer.cs | 2 +- Refresh.Interfaces.Workers/IWorker.cs | 2 +- .../RequestStatisticSubmitWorker.cs | 2 +- Refresh.Interfaces.Workers/WorkContext.cs | 12 +++++++++ Refresh.Interfaces.Workers/WorkerManager.cs | 13 +++------- .../Workers/CoolLevelsWorker.cs | 6 ++--- .../Workers/DiscordIntegrationWorker.cs | 26 +++++++++++-------- .../Workers/ExpiredObjectWorker.cs | 2 +- .../Workers/ObjectStatisticsWorker.cs | 2 +- .../Workers/PunishmentExpiryWorker.cs | 2 +- RefreshTests.GameServer/TestContext.cs | 12 +++++++++ .../Tests/Workers/PunishmentExpiryTests.cs | 8 +++--- 12 files changed, 55 insertions(+), 34 deletions(-) create mode 100644 Refresh.Interfaces.Workers/WorkContext.cs diff --git a/Refresh.GameServer/RefreshGameServer.cs b/Refresh.GameServer/RefreshGameServer.cs index caedaecb8..b98145fc3 100644 --- a/Refresh.GameServer/RefreshGameServer.cs +++ b/Refresh.GameServer/RefreshGameServer.cs @@ -206,7 +206,7 @@ protected override void SetupServices() protected virtual void SetupWorkers() { - this.WorkerManager = new WorkerManager(this.Logger, this._dataStore, this._databaseProvider, this._matchService, this._guidCheckerService); + this.WorkerManager = new WorkerManager(this.Logger, this._dataStore, this._databaseProvider); this.WorkerManager.AddWorker(); this.WorkerManager.AddWorker(); diff --git a/Refresh.Interfaces.Workers/IWorker.cs b/Refresh.Interfaces.Workers/IWorker.cs index edeb96858..ccecabc2a 100644 --- a/Refresh.Interfaces.Workers/IWorker.cs +++ b/Refresh.Interfaces.Workers/IWorker.cs @@ -16,5 +16,5 @@ public interface IWorker /// The server's data store, for workers to use. /// A database context, for workers to use. /// True if the worker did work, false if it did not. - public void DoWork(DataContext context); + public void DoWork(WorkContext context); } \ No newline at end of file diff --git a/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs b/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs index 1de24798c..e5f2c6ad8 100644 --- a/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs +++ b/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs @@ -7,7 +7,7 @@ public class RequestStatisticSubmitWorker : IWorker { public int WorkInterval => 5_000; - public void DoWork(DataContext context) + public void DoWork(WorkContext context) { (int game, int api) = RequestStatisticTrackingMiddleware.SubmitAndClearRequests(); diff --git a/Refresh.Interfaces.Workers/WorkContext.cs b/Refresh.Interfaces.Workers/WorkContext.cs new file mode 100644 index 000000000..03e4cc579 --- /dev/null +++ b/Refresh.Interfaces.Workers/WorkContext.cs @@ -0,0 +1,12 @@ +using Bunkum.Core.Storage; +using NotEnoughLogs; +using Refresh.Database; + +namespace Refresh.Interfaces.Workers; + +public class WorkContext : IDataContext +{ + public required GameDatabaseContext Database { get; init; } + public required Logger Logger { get; init; } + public required IDataStore DataStore { get; init; } +} \ No newline at end of file diff --git a/Refresh.Interfaces.Workers/WorkerManager.cs b/Refresh.Interfaces.Workers/WorkerManager.cs index 91758bf35..42da43fec 100644 --- a/Refresh.Interfaces.Workers/WorkerManager.cs +++ b/Refresh.Interfaces.Workers/WorkerManager.cs @@ -12,16 +12,12 @@ public class WorkerManager private readonly Logger _logger; private readonly IDataStore _dataStore; private readonly GameDatabaseProvider _databaseProvider; - private readonly MatchService _matchService; - private readonly GuidCheckerService _guidCheckerService; - public WorkerManager(Logger logger, IDataStore dataStore, GameDatabaseProvider databaseProvider, MatchService matchService, GuidCheckerService guidCheckerService) + public WorkerManager(Logger logger, IDataStore dataStore, GameDatabaseProvider databaseProvider) { this._dataStore = dataStore; this._databaseProvider = databaseProvider; this._logger = logger; - this._matchService = matchService; - this._guidCheckerService = guidCheckerService; } private Thread? _thread = null; @@ -42,14 +38,11 @@ public void AddWorker(IWorker worker) private void RunWorkCycle() { - Lazy dataContext = new(() => new DataContext + Lazy workContext = new(() => new WorkContext { Database = this._databaseProvider.GetContext(), Logger = this._logger, DataStore = this._dataStore, - Match = this._matchService, - Token = null, - GuidChecker = this._guidCheckerService, }); foreach (IWorker worker in this._workers) @@ -67,7 +60,7 @@ private void RunWorkCycle() } this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + worker.GetType().Name); - worker.DoWork(dataContext.Value); + worker.DoWork(workContext.Value); } } diff --git a/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs b/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs index 6da335674..572ff801a 100644 --- a/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs +++ b/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs @@ -23,7 +23,7 @@ public class CoolLevelsWorker : IWorker public int WorkInterval => 600_000; // Every 10 minutes [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] - public void DoWork(DataContext context) + public void DoWork(WorkContext context) { const int pageSize = 1000; DatabaseList levels = context.Database.GetUserLevelsChunk(0, pageSize); @@ -104,7 +104,7 @@ private static float CalculateLevelDecayMultiplier(Logger logger, long now, Game return (float)multiplier; } - private static float CalculatePositiveScore(GameLevel level, DataContext context) + private static float CalculatePositiveScore(GameLevel level, WorkContext context) { Debug.Assert(level.Statistics != null); @@ -163,7 +163,7 @@ private static float CalculatePositiveScore(GameLevel level, DataContext context return score; } - private static float CalculateNegativeScore(GameLevel level, DataContext context) + private static float CalculateNegativeScore(GameLevel level, WorkContext context) { Debug.Assert(level.Statistics != null); diff --git a/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs b/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs index 7a443cde4..4754792b4 100644 --- a/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs +++ b/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs @@ -5,6 +5,10 @@ using Refresh.Core.Types.Data; using Refresh.Database; using Refresh.Database.Models.Activity; +using Refresh.Database.Models.Levels; +using Refresh.Database.Models.Levels.Scores; +using Refresh.Database.Models.Photos; +using Refresh.Database.Models.Users; using Refresh.Database.Query; using Refresh.Interfaces.APIv3.Endpoints.DataTypes.Response.Levels; using Refresh.Interfaces.APIv3.Endpoints.DataTypes.Response.Users; @@ -43,21 +47,21 @@ private string GetAssetUrl(string hash) return $"{this._externalUrl}/api/v3/assets/{hash}/image"; } - private Embed? GenerateEmbedFromEvent(Event @event, DataContext context) + private Embed? GenerateEmbedFromEvent(Event @event, WorkContext context) { EmbedBuilder embed = new(); - ApiGameLevelResponse? level = @event.StoredDataType == EventDataType.Level ? - ApiGameLevelResponse.FromOld(context.Database.GetLevelById(@event.StoredSequentialId!.Value), context) + GameLevel? level = @event.StoredDataType == EventDataType.Level ? + context.Database.GetLevelById(@event.StoredSequentialId!.Value) : null; - ApiGameUserResponse? user = @event.StoredDataType == EventDataType.User ? - ApiGameUserResponse.FromOld(context.Database.GetUserByObjectId(@event.StoredObjectId), context) + GameUser? user = @event.StoredDataType == EventDataType.User ? + context.Database.GetUserByObjectId(@event.StoredObjectId) : null; - ApiGameScoreResponse? score = @event.StoredDataType == EventDataType.Score ? - ApiGameScoreResponse.FromOld(context.Database.GetScoreByObjectId(@event.StoredObjectId), context) + GameScore? score = @event.StoredDataType == EventDataType.Score ? + context.Database.GetScoreByObjectId(@event.StoredObjectId) : null; - ApiGamePhotoResponse? photo = @event.StoredDataType == EventDataType.Photo ? - ApiGamePhotoResponse.FromOld(context.Database.GetPhotoFromEvent(@event), context) + GamePhoto? photo = @event.StoredDataType == EventDataType.Photo ? + context.Database.GetPhotoFromEvent(@event) : null; if (photo != null) @@ -92,7 +96,7 @@ private string GetAssetUrl(string hash) embed.WithDescription($"[{@event.User.Username}]({this._externalUrl}/u/{@event.User.UserId}) {description}"); if (photo != null) - embed.WithImageUrl(this.GetAssetUrl(photo.LargeHash)); + embed.WithImageUrl(this.GetAssetUrl(photo.LargeAsset.AssetHash)); else if (level != null) embed.WithThumbnailUrl(this.GetAssetUrl(level.IconHash)); else if (user != null) @@ -104,7 +108,7 @@ private string GetAssetUrl(string hash) return embed.Build(); } - public void DoWork(DataContext context) + public void DoWork(WorkContext context) { if (this._firstCycle) { diff --git a/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs b/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs index 6152ebb4f..244e7d950 100644 --- a/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs +++ b/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs @@ -8,7 +8,7 @@ namespace Refresh.Interfaces.Workers.Workers; public class ExpiredObjectWorker : IWorker { public int WorkInterval => 60_000; // 1 minute - public void DoWork(DataContext context) + public void DoWork(WorkContext context) { List registrationsToRemove = []; List codesToRemove = []; diff --git a/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs b/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs index 4308310bb..d16247169 100644 --- a/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs +++ b/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs @@ -11,7 +11,7 @@ public class ObjectStatisticsWorker : IWorker public int WorkInterval => 60_000; [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] - public void DoWork(DataContext context) + public void DoWork(WorkContext context) { GameLevel[] levels = context.Database.GetLevelsWithStatisticsNeedingUpdates() .Take(500) diff --git a/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs b/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs index 03de4bba7..b046a77da 100644 --- a/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs +++ b/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs @@ -12,7 +12,7 @@ public class PunishmentExpiryWorker : IWorker { public int WorkInterval => 60_000; // 1 minute - public void DoWork(DataContext context) + public void DoWork(WorkContext context) { DatabaseList bannedUsers = context.Database.GetAllUsersWithRole(GameUserRole.Banned); DatabaseList restrictedUsers = context.Database.GetAllUsersWithRole(GameUserRole.Restricted); diff --git a/RefreshTests.GameServer/TestContext.cs b/RefreshTests.GameServer/TestContext.cs index e139d6494..32cdcd8ec 100644 --- a/RefreshTests.GameServer/TestContext.cs +++ b/RefreshTests.GameServer/TestContext.cs @@ -12,6 +12,7 @@ using Refresh.Database.Models.Levels.Scores; using Refresh.Database.Models.Levels; using Refresh.Interfaces.Game.Types.UserData.Leaderboard; +using Refresh.Interfaces.Workers; namespace RefreshTests.GameServer; @@ -181,6 +182,17 @@ public DataContext GetDataContext(Token? token = null) GuidChecker = this.GetService(), }; } + + public WorkContext GetWorkContext() + { + return new WorkContext + { + Database = this.Database, + Logger = this.Server.Value.Logger, + DataStore = (IDataStore)this.GetService() + .AddParameterToEndpoint(null!, new BunkumParameterInfo(typeof(IDataStore), ""), null!)!, + }; + } public void Dispose() { diff --git a/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs b/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs index 844a59f3f..655dbadaf 100644 --- a/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs +++ b/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs @@ -28,14 +28,14 @@ public void BannedUsersExpire() Assert.That(context.Database.GetAllUsersWithRole(Banned).Items, Contains.Item(user)); }); - worker.DoWork(context.GetDataContext()); + worker.DoWork(context.GetWorkContext()); Assert.Multiple(() => { Assert.That(user.Role, Is.EqualTo(Banned)); }); context.Time.TimestampMilliseconds = 2000; - worker.DoWork(context.GetDataContext()); + worker.DoWork(context.GetWorkContext()); context.Database.Refresh(); user = context.Database.GetUserByObjectId(user.UserId)!; @@ -57,14 +57,14 @@ public void RestrictedUsersExpire() context.Database.RestrictUser(user, "", DateTimeOffset.FromUnixTimeMilliseconds(1000)); Assert.That(user.Role, Is.EqualTo(Restricted)); - worker.DoWork(context.GetDataContext()); + worker.DoWork(context.GetWorkContext()); Assert.Multiple(() => { Assert.That(user.Role, Is.EqualTo(Restricted)); }); context.Time.TimestampMilliseconds = 2000; - worker.DoWork(context.GetDataContext()); + worker.DoWork(context.GetWorkContext()); context.Database.Refresh(); user = context.Database.GetUserByObjectId(user.UserId)!; From c7c883601b580fc217034d4d529a7bbdf8a84fb8 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 17:45:49 -0400 Subject: [PATCH 02/12] Move base worker implementation to project --- Refresh.GameServer/RefreshGameServer.cs | 15 +++++----- .../{Workers => }/CoolLevelsWorker.cs | 10 +++---- .../{Workers => }/DiscordIntegrationWorker.cs | 15 ++++------ .../{Workers => }/ExpiredObjectWorker.cs | 10 +++---- Refresh.Interfaces.Workers/IWorker.cs | 20 ------------- .../{Workers => }/ObjectStatisticsWorker.cs | 10 +++---- .../{Workers => }/PunishmentExpiryWorker.cs | 10 +++---- .../Refresh.Interfaces.Workers.csproj | 2 +- .../RequestStatisticSubmitWorker.cs | 16 ++++++++++ .../RequestStatisticSubmitWorker.cs | 16 ---------- Refresh.Workers/Refresh.Workers.csproj | 17 +++++++++++ .../WorkContext.cs | 2 +- Refresh.Workers/WorkerJob.cs | 18 ++++++++++++ .../WorkerManager.cs | 16 +++++----- Refresh.sln | 9 ++++++ RefreshTests.GameServer/TestContext.cs | 1 + .../Tests/Workers/PunishmentExpiryTests.cs | 29 +++++++------------ 17 files changed, 113 insertions(+), 103 deletions(-) rename Refresh.Interfaces.Workers/{Workers => }/CoolLevelsWorker.cs (97%) rename Refresh.Interfaces.Workers/{Workers => }/DiscordIntegrationWorker.cs (89%) rename Refresh.Interfaces.Workers/{Workers => }/ExpiredObjectWorker.cs (89%) delete mode 100644 Refresh.Interfaces.Workers/IWorker.cs rename Refresh.Interfaces.Workers/{Workers => }/ObjectStatisticsWorker.cs (81%) rename Refresh.Interfaces.Workers/{Workers => }/PunishmentExpiryWorker.cs (84%) create mode 100644 Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs delete mode 100644 Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs create mode 100644 Refresh.Workers/Refresh.Workers.csproj rename {Refresh.Interfaces.Workers => Refresh.Workers}/WorkContext.cs (88%) create mode 100644 Refresh.Workers/WorkerJob.cs rename {Refresh.Interfaces.Workers => Refresh.Workers}/WorkerManager.cs (86%) diff --git a/Refresh.GameServer/RefreshGameServer.cs b/Refresh.GameServer/RefreshGameServer.cs index b98145fc3..ebab8b29f 100644 --- a/Refresh.GameServer/RefreshGameServer.cs +++ b/Refresh.GameServer/RefreshGameServer.cs @@ -31,8 +31,7 @@ using Refresh.Interfaces.Game; using Refresh.Interfaces.Internal; using Refresh.Interfaces.Workers; -using Refresh.Interfaces.Workers.RequestTracking; -using Refresh.Interfaces.Workers.Workers; +using Refresh.Workers; namespace Refresh.GameServer; @@ -208,15 +207,15 @@ protected virtual void SetupWorkers() { this.WorkerManager = new WorkerManager(this.Logger, this._dataStore, this._databaseProvider); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); + this.WorkerManager.AddWorker(); + this.WorkerManager.AddWorker(); + this.WorkerManager.AddWorker(); + this.WorkerManager.AddWorker(); + this.WorkerManager.AddWorker(); if ((this._integrationConfig?.DiscordWebhookEnabled ?? false) && this._config != null && this._config.PermitShowingOnlineUsers) { - this.WorkerManager.AddWorker(new DiscordIntegrationWorker(this._integrationConfig, this._config)); + this.WorkerManager.AddWorker(new DiscordIntegrationJob(this._integrationConfig, this._config)); } } diff --git a/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs b/Refresh.Interfaces.Workers/CoolLevelsWorker.cs similarity index 97% rename from Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs rename to Refresh.Interfaces.Workers/CoolLevelsWorker.cs index 572ff801a..1bcf71917 100644 --- a/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs +++ b/Refresh.Interfaces.Workers/CoolLevelsWorker.cs @@ -10,20 +10,20 @@ using System.Diagnostics.CodeAnalysis; using NotEnoughLogs; using Refresh.Core; -using Refresh.Core.Types.Data; using Refresh.Database; using Refresh.Database.Models.Authentication; using Refresh.Database.Models.Levels; using Refresh.Database.Models.Users; +using Refresh.Workers; -namespace Refresh.Interfaces.Workers.Workers; +namespace Refresh.Interfaces.Workers; -public class CoolLevelsWorker : IWorker +public class CoolLevelsJob : WorkerJob { - public int WorkInterval => 600_000; // Every 10 minutes + public override int WorkInterval => 600_000; // Every 10 minutes [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] - public void DoWork(WorkContext context) + public override void ExecuteJob(WorkContext context) { const int pageSize = 1000; DatabaseList levels = context.Database.GetUserLevelsChunk(0, pageSize); diff --git a/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs similarity index 89% rename from Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs rename to Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs index 4754792b4..129c6e7d0 100644 --- a/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs +++ b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs @@ -2,7 +2,6 @@ using Discord.Webhook; using Refresh.Core; using Refresh.Core.Configuration; -using Refresh.Core.Types.Data; using Refresh.Database; using Refresh.Database.Models.Activity; using Refresh.Database.Models.Levels; @@ -10,13 +9,11 @@ using Refresh.Database.Models.Photos; using Refresh.Database.Models.Users; using Refresh.Database.Query; -using Refresh.Interfaces.APIv3.Endpoints.DataTypes.Response.Levels; -using Refresh.Interfaces.APIv3.Endpoints.DataTypes.Response.Users; -using Refresh.Interfaces.APIv3.Endpoints.DataTypes.Response.Users.Photos; +using Refresh.Workers; -namespace Refresh.Interfaces.Workers.Workers; +namespace Refresh.Interfaces.Workers; -public class DiscordIntegrationWorker : IWorker +public class DiscordIntegrationJob : WorkerJob { private readonly IntegrationConfig _config; private readonly string _externalUrl; @@ -26,9 +23,9 @@ public class DiscordIntegrationWorker : IWorker private long _lastTimestamp; private static long Now => DateTimeOffset.Now.ToUnixTimeMilliseconds(); - public int WorkInterval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default + public override int WorkInterval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default - public DiscordIntegrationWorker(IntegrationConfig config, GameServerConfig gameConfig) + public DiscordIntegrationJob(IntegrationConfig config, GameServerConfig gameConfig) { this._config = config; this._externalUrl = gameConfig.WebExternalUrl; @@ -108,7 +105,7 @@ private string GetAssetUrl(string hash) return embed.Build(); } - public void DoWork(WorkContext context) + public override void ExecuteJob(WorkContext context) { if (this._firstCycle) { diff --git a/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs b/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs similarity index 89% rename from Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs rename to Refresh.Interfaces.Workers/ExpiredObjectWorker.cs index 244e7d950..07073bfad 100644 --- a/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs +++ b/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs @@ -1,14 +1,14 @@ using Refresh.Core; -using Refresh.Core.Types.Data; using Refresh.Database.Models.Authentication; using Refresh.Database.Models.Users; +using Refresh.Workers; -namespace Refresh.Interfaces.Workers.Workers; +namespace Refresh.Interfaces.Workers; -public class ExpiredObjectWorker : IWorker +public class CleanupExpiredObjectsJob : WorkerJob { - public int WorkInterval => 60_000; // 1 minute - public void DoWork(WorkContext context) + public override int WorkInterval => 60_000; // 1 minute + public override void ExecuteJob(WorkContext context) { List registrationsToRemove = []; List codesToRemove = []; diff --git a/Refresh.Interfaces.Workers/IWorker.cs b/Refresh.Interfaces.Workers/IWorker.cs deleted file mode 100644 index ccecabc2a..000000000 --- a/Refresh.Interfaces.Workers/IWorker.cs +++ /dev/null @@ -1,20 +0,0 @@ -using Refresh.Core.Types.Data; - -namespace Refresh.Interfaces.Workers; - -public interface IWorker -{ - /// - /// How often to perform work, in milliseconds - /// - public int WorkInterval { get; } - - /// - /// Instructs the worker to do work. - /// - /// - /// The server's data store, for workers to use. - /// A database context, for workers to use. - /// True if the worker did work, false if it did not. - public void DoWork(WorkContext context); -} \ No newline at end of file diff --git a/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs b/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs similarity index 81% rename from Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs rename to Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs index d16247169..414c6c878 100644 --- a/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs +++ b/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs @@ -1,17 +1,17 @@ using System.Diagnostics.CodeAnalysis; using Refresh.Core; -using Refresh.Core.Types.Data; using Refresh.Database.Models.Levels; using Refresh.Database.Models.Users; +using Refresh.Workers; -namespace Refresh.Interfaces.Workers.Workers; +namespace Refresh.Interfaces.Workers; -public class ObjectStatisticsWorker : IWorker +public class ObjectStatisticsJob : WorkerJob { - public int WorkInterval => 60_000; + public override int WorkInterval => 60_000; [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] - public void DoWork(WorkContext context) + public override void ExecuteJob(WorkContext context) { GameLevel[] levels = context.Database.GetLevelsWithStatisticsNeedingUpdates() .Take(500) diff --git a/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs b/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs similarity index 84% rename from Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs rename to Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs index b046a77da..0699a560d 100644 --- a/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs +++ b/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs @@ -1,18 +1,18 @@ using Refresh.Core; -using Refresh.Core.Types.Data; using Refresh.Database; using Refresh.Database.Models.Users; +using Refresh.Workers; -namespace Refresh.Interfaces.Workers.Workers; +namespace Refresh.Interfaces.Workers; /// /// A worker that checks all users for bans/restrictions, then removes them if expired. /// -public class PunishmentExpiryWorker : IWorker +public class PunishmentExpiryJob : WorkerJob { - public int WorkInterval => 60_000; // 1 minute + public override int WorkInterval => 60_000; // 1 minute - public void DoWork(WorkContext context) + public override void ExecuteJob(WorkContext context) { DatabaseList bannedUsers = context.Database.GetAllUsersWithRole(GameUserRole.Banned); DatabaseList restrictedUsers = context.Database.GetAllUsersWithRole(GameUserRole.Restricted); diff --git a/Refresh.Interfaces.Workers/Refresh.Interfaces.Workers.csproj b/Refresh.Interfaces.Workers/Refresh.Interfaces.Workers.csproj index f64ff9419..98671073d 100644 --- a/Refresh.Interfaces.Workers/Refresh.Interfaces.Workers.csproj +++ b/Refresh.Interfaces.Workers/Refresh.Interfaces.Workers.csproj @@ -12,7 +12,7 @@ - + diff --git a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs b/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs new file mode 100644 index 000000000..e032a7363 --- /dev/null +++ b/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs @@ -0,0 +1,16 @@ +using Refresh.Core.Metrics; +using Refresh.Workers; + +namespace Refresh.Interfaces.Workers; + +public class RequestStatisticSubmitJob : WorkerJob +{ + public override int WorkInterval => 5_000; + + public override void ExecuteJob(WorkContext context) + { + (int game, int api) = RequestStatisticTrackingMiddleware.SubmitAndClearRequests(); + + context.Database.IncrementRequests(api, game); + } +} \ No newline at end of file diff --git a/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs b/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs deleted file mode 100644 index e5f2c6ad8..000000000 --- a/Refresh.Interfaces.Workers/RequestTracking/RequestStatisticSubmitWorker.cs +++ /dev/null @@ -1,16 +0,0 @@ -using Refresh.Core.Metrics; -using Refresh.Core.Types.Data; - -namespace Refresh.Interfaces.Workers.RequestTracking; - -public class RequestStatisticSubmitWorker : IWorker -{ - public int WorkInterval => 5_000; - - public void DoWork(WorkContext context) - { - (int game, int api) = RequestStatisticTrackingMiddleware.SubmitAndClearRequests(); - - context.Database.IncrementRequests(api, game); - } -} \ No newline at end of file diff --git a/Refresh.Workers/Refresh.Workers.csproj b/Refresh.Workers/Refresh.Workers.csproj new file mode 100644 index 000000000..4a42bbe59 --- /dev/null +++ b/Refresh.Workers/Refresh.Workers.csproj @@ -0,0 +1,17 @@ + + + + net9.0 + enable + enable + true + 612,618 + Debug;Release + AnyCPU + + + + + + + diff --git a/Refresh.Interfaces.Workers/WorkContext.cs b/Refresh.Workers/WorkContext.cs similarity index 88% rename from Refresh.Interfaces.Workers/WorkContext.cs rename to Refresh.Workers/WorkContext.cs index 03e4cc579..6e4bc2289 100644 --- a/Refresh.Interfaces.Workers/WorkContext.cs +++ b/Refresh.Workers/WorkContext.cs @@ -2,7 +2,7 @@ using NotEnoughLogs; using Refresh.Database; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Workers; public class WorkContext : IDataContext { diff --git a/Refresh.Workers/WorkerJob.cs b/Refresh.Workers/WorkerJob.cs new file mode 100644 index 000000000..31b4ecc09 --- /dev/null +++ b/Refresh.Workers/WorkerJob.cs @@ -0,0 +1,18 @@ +namespace Refresh.Workers; + +public abstract class WorkerJob +{ + /// + /// How often to perform work, in milliseconds + /// + public abstract int WorkInterval { get; } + + public bool FirstCycle { get; internal set; } + + + /// + /// Executes the job. + /// + /// Contains data access so workers can interact with the server. + public abstract void ExecuteJob(WorkContext context); +} \ No newline at end of file diff --git a/Refresh.Interfaces.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs similarity index 86% rename from Refresh.Interfaces.Workers/WorkerManager.cs rename to Refresh.Workers/WorkerManager.cs index 42da43fec..b8b000f34 100644 --- a/Refresh.Interfaces.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -1,11 +1,9 @@ using Bunkum.Core.Storage; using NotEnoughLogs; using Refresh.Core; -using Refresh.Core.Services; -using Refresh.Core.Types.Data; using Refresh.Database; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Workers; public class WorkerManager { @@ -23,15 +21,15 @@ public WorkerManager(Logger logger, IDataStore dataStore, GameDatabaseProvider d private Thread? _thread = null; private bool _threadShouldRun = false; - private readonly List _workers = []; - private readonly Dictionary _lastWorkTimestamps = new(); + private readonly List _workers = []; + private readonly Dictionary _lastWorkTimestamps = new(); - public void AddWorker() where TWorker : IWorker, new() + public void AddWorker() where TWorker : WorkerJob, new() { TWorker worker = new(); this._workers.Add(worker); } - public void AddWorker(IWorker worker) + public void AddWorker(WorkerJob worker) { this._workers.Add(worker); } @@ -45,7 +43,7 @@ private void RunWorkCycle() DataStore = this._dataStore, }); - foreach (IWorker worker in this._workers) + foreach (WorkerJob worker in this._workers) { long now = DateTimeOffset.Now.ToUnixTimeMilliseconds(); if (this._lastWorkTimestamps.TryGetValue(worker, out long lastWork)) @@ -60,7 +58,7 @@ private void RunWorkCycle() } this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + worker.GetType().Name); - worker.DoWork(workContext.Value); + worker.ExecuteJob(workContext.Value); } } diff --git a/Refresh.sln b/Refresh.sln index 111f272f3..cb10f4662 100644 --- a/Refresh.sln +++ b/Refresh.sln @@ -33,6 +33,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Refresh.Interfaces.Internal EndProject Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Interfaces", "Interfaces", "{2B6AE134-DE2B-46A8-925D-6D4B1908BFAA}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Refresh.Workers", "Refresh.Workers\Refresh.Workers.csproj", "{F33E3D8D-7867-41C6-8E24-5BA95642D001}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -109,6 +111,12 @@ Global {14EF66DA-391A-40A1-80B3-03A2553C71C9}.Release|Any CPU.Build.0 = Release|Any CPU {14EF66DA-391A-40A1-80B3-03A2553C71C9}.DebugLocalBunkum|Any CPU.ActiveCfg = Debug|Any CPU {14EF66DA-391A-40A1-80B3-03A2553C71C9}.DebugLocalBunkum|Any CPU.Build.0 = Debug|Any CPU + {F33E3D8D-7867-41C6-8E24-5BA95642D001}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {F33E3D8D-7867-41C6-8E24-5BA95642D001}.Debug|Any CPU.Build.0 = Debug|Any CPU + {F33E3D8D-7867-41C6-8E24-5BA95642D001}.Release|Any CPU.ActiveCfg = Release|Any CPU + {F33E3D8D-7867-41C6-8E24-5BA95642D001}.Release|Any CPU.Build.0 = Release|Any CPU + {F33E3D8D-7867-41C6-8E24-5BA95642D001}.DebugLocalBunkum|Any CPU.ActiveCfg = Debug|Any CPU + {F33E3D8D-7867-41C6-8E24-5BA95642D001}.DebugLocalBunkum|Any CPU.Build.0 = Debug|Any CPU EndGlobalSection GlobalSection(NestedProjects) = preSolution {A3D76B31-8732-4134-907B-F36773DF70D3} = {AA547E95-555C-470E-995A-D66E56660542} @@ -121,5 +129,6 @@ Global {9EACC405-0340-46C9-8437-68AC71AB05CF} = {2B6AE134-DE2B-46A8-925D-6D4B1908BFAA} {14EF66DA-391A-40A1-80B3-03A2553C71C9} = {2B6AE134-DE2B-46A8-925D-6D4B1908BFAA} {0FD65A5F-F20F-4A30-B486-B21146686FA6} = {2B6AE134-DE2B-46A8-925D-6D4B1908BFAA} + {F33E3D8D-7867-41C6-8E24-5BA95642D001} = {AA547E95-555C-470E-995A-D66E56660542} EndGlobalSection EndGlobal diff --git a/RefreshTests.GameServer/TestContext.cs b/RefreshTests.GameServer/TestContext.cs index 32cdcd8ec..3563eb873 100644 --- a/RefreshTests.GameServer/TestContext.cs +++ b/RefreshTests.GameServer/TestContext.cs @@ -13,6 +13,7 @@ using Refresh.Database.Models.Levels; using Refresh.Interfaces.Game.Types.UserData.Leaderboard; using Refresh.Interfaces.Workers; +using Refresh.Workers; namespace RefreshTests.GameServer; diff --git a/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs b/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs index 655dbadaf..4acc1fd91 100644 --- a/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs +++ b/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs @@ -1,6 +1,6 @@ using NotEnoughLogs; using Refresh.Database.Models.Users; -using Refresh.Interfaces.Workers.Workers; +using Refresh.Interfaces.Workers; using RefreshTests.GameServer.Logging; using static Refresh.Database.Models.Users.GameUserRole; @@ -12,7 +12,7 @@ public class PunishmentExpiryTests : GameServerTest public void BannedUsersExpire() { using TestContext context = this.GetServer(); - PunishmentExpiryWorker worker = new(); + PunishmentExpiryJob worker = new(); GameUser user = context.CreateUser(); Assert.Multiple(() => { @@ -28,28 +28,22 @@ public void BannedUsersExpire() Assert.That(context.Database.GetAllUsersWithRole(Banned).Items, Contains.Item(user)); }); - worker.DoWork(context.GetWorkContext()); - Assert.Multiple(() => - { - Assert.That(user.Role, Is.EqualTo(Banned)); - }); + worker.ExecuteJob(context.GetWorkContext()); + Assert.That(user.Role, Is.EqualTo(Banned)); context.Time.TimestampMilliseconds = 2000; - worker.DoWork(context.GetWorkContext()); + worker.ExecuteJob(context.GetWorkContext()); context.Database.Refresh(); user = context.Database.GetUserByObjectId(user.UserId)!; - Assert.Multiple(() => - { - Assert.That(user.Role, Is.EqualTo(User)); - }); + Assert.That(user.Role, Is.EqualTo(User)); } [Test] public void RestrictedUsersExpire() { using TestContext context = this.GetServer(); - PunishmentExpiryWorker worker = new(); + PunishmentExpiryJob worker = new(); GameUser user = context.CreateUser(); Assert.That(user.Role, Is.EqualTo(User)); @@ -57,14 +51,11 @@ public void RestrictedUsersExpire() context.Database.RestrictUser(user, "", DateTimeOffset.FromUnixTimeMilliseconds(1000)); Assert.That(user.Role, Is.EqualTo(Restricted)); - worker.DoWork(context.GetWorkContext()); - Assert.Multiple(() => - { - Assert.That(user.Role, Is.EqualTo(Restricted)); - }); + worker.ExecuteJob(context.GetWorkContext()); + Assert.That(user.Role, Is.EqualTo(Restricted)); context.Time.TimestampMilliseconds = 2000; - worker.DoWork(context.GetWorkContext()); + worker.ExecuteJob(context.GetWorkContext()); context.Database.Refresh(); user = context.Database.GetUserByObjectId(user.UserId)!; From 4851064015fe8879d3dd97c8116f15cdfbf32a43 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 17:47:20 -0400 Subject: [PATCH 03/12] Simply first cycle logic --- .../DiscordIntegrationWorker.cs | 14 ++------------ Refresh.Workers/WorkerJob.cs | 3 +-- Refresh.Workers/WorkerManager.cs | 1 + 3 files changed, 4 insertions(+), 14 deletions(-) diff --git a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs index 129c6e7d0..13c83ad4f 100644 --- a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs +++ b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs @@ -18,8 +18,6 @@ public class DiscordIntegrationJob : WorkerJob private readonly IntegrationConfig _config; private readonly string _externalUrl; private readonly DiscordWebhookClient _client; - - private bool _firstCycle = true; private long _lastTimestamp; private static long Now => DateTimeOffset.Now.ToUnixTimeMilliseconds(); @@ -33,12 +31,6 @@ public DiscordIntegrationJob(IntegrationConfig config, GameServerConfig gameConf this._client = new DiscordWebhookClient(config.DiscordWebhookUrl); } - private void DoFirstCycle() - { - this._firstCycle = false; - this._lastTimestamp = Now; - } - private string GetAssetUrl(string hash) { return $"{this._externalUrl}/api/v3/assets/{hash}/image"; @@ -107,10 +99,8 @@ private string GetAssetUrl(string hash) public override void ExecuteJob(WorkContext context) { - if (this._firstCycle) - { - this.DoFirstCycle(); - } + if (this.FirstCycle) + this._lastTimestamp = Now; DatabaseList activity = context.Database.GetGlobalRecentActivity(new ActivityQueryParameters { diff --git a/Refresh.Workers/WorkerJob.cs b/Refresh.Workers/WorkerJob.cs index 31b4ecc09..52ab3450b 100644 --- a/Refresh.Workers/WorkerJob.cs +++ b/Refresh.Workers/WorkerJob.cs @@ -7,8 +7,7 @@ public abstract class WorkerJob /// public abstract int WorkInterval { get; } - public bool FirstCycle { get; internal set; } - + public bool FirstCycle { get; internal set; } = true; /// /// Executes the job. diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index b8b000f34..48abaa19d 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -59,6 +59,7 @@ private void RunWorkCycle() this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + worker.GetType().Name); worker.ExecuteJob(workContext.Value); + worker.FirstCycle = false; } } From ff54bcbb1458d54eb12dbb28067c30b1b9651b2c Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 17:53:44 -0400 Subject: [PATCH 04/12] Rename everything to "job", add StartupJob --- Refresh.GameServer/RefreshGameServer.cs | 12 +++---- .../CoolLevelsWorker.cs | 2 +- .../DiscordIntegrationWorker.cs | 2 +- .../ExpiredObjectWorker.cs | 2 +- .../ObjectStatisticsWorker.cs | 2 +- .../PunishmentExpiryWorker.cs | 2 +- .../RequestStatisticSubmitWorker.cs | 2 +- Refresh.Workers/StartupJob.cs | 17 ++++++++++ Refresh.Workers/WorkerJob.cs | 2 +- Refresh.Workers/WorkerManager.cs | 32 +++++++++---------- 10 files changed, 46 insertions(+), 29 deletions(-) create mode 100644 Refresh.Workers/StartupJob.cs diff --git a/Refresh.GameServer/RefreshGameServer.cs b/Refresh.GameServer/RefreshGameServer.cs index ebab8b29f..120b45483 100644 --- a/Refresh.GameServer/RefreshGameServer.cs +++ b/Refresh.GameServer/RefreshGameServer.cs @@ -207,15 +207,15 @@ protected virtual void SetupWorkers() { this.WorkerManager = new WorkerManager(this.Logger, this._dataStore, this._databaseProvider); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); - this.WorkerManager.AddWorker(); + this.WorkerManager.AddJob(); + this.WorkerManager.AddJob(); + this.WorkerManager.AddJob(); + this.WorkerManager.AddJob(); + this.WorkerManager.AddJob(); if ((this._integrationConfig?.DiscordWebhookEnabled ?? false) && this._config != null && this._config.PermitShowingOnlineUsers) { - this.WorkerManager.AddWorker(new DiscordIntegrationJob(this._integrationConfig, this._config)); + this.WorkerManager.AddJob(new DiscordIntegrationJob(this._integrationConfig, this._config)); } } diff --git a/Refresh.Interfaces.Workers/CoolLevelsWorker.cs b/Refresh.Interfaces.Workers/CoolLevelsWorker.cs index 1bcf71917..cd17f1677 100644 --- a/Refresh.Interfaces.Workers/CoolLevelsWorker.cs +++ b/Refresh.Interfaces.Workers/CoolLevelsWorker.cs @@ -20,7 +20,7 @@ namespace Refresh.Interfaces.Workers; public class CoolLevelsJob : WorkerJob { - public override int WorkInterval => 600_000; // Every 10 minutes + public override int Interval => 600_000; // Every 10 minutes [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] public override void ExecuteJob(WorkContext context) diff --git a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs index 13c83ad4f..c5713ed1d 100644 --- a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs +++ b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs @@ -21,7 +21,7 @@ public class DiscordIntegrationJob : WorkerJob private long _lastTimestamp; private static long Now => DateTimeOffset.Now.ToUnixTimeMilliseconds(); - public override int WorkInterval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default + public override int Interval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default public DiscordIntegrationJob(IntegrationConfig config, GameServerConfig gameConfig) { diff --git a/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs b/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs index 07073bfad..fde6c5834 100644 --- a/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs +++ b/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs @@ -7,7 +7,7 @@ namespace Refresh.Interfaces.Workers; public class CleanupExpiredObjectsJob : WorkerJob { - public override int WorkInterval => 60_000; // 1 minute + public override int Interval => 60_000; // 1 minute public override void ExecuteJob(WorkContext context) { List registrationsToRemove = []; diff --git a/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs b/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs index 414c6c878..a64032173 100644 --- a/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs +++ b/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs @@ -8,7 +8,7 @@ namespace Refresh.Interfaces.Workers; public class ObjectStatisticsJob : WorkerJob { - public override int WorkInterval => 60_000; + public override int Interval => 60_000; [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] public override void ExecuteJob(WorkContext context) diff --git a/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs b/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs index 0699a560d..ae86556c4 100644 --- a/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs +++ b/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs @@ -10,7 +10,7 @@ namespace Refresh.Interfaces.Workers; /// public class PunishmentExpiryJob : WorkerJob { - public override int WorkInterval => 60_000; // 1 minute + public override int Interval => 60_000; // 1 minute public override void ExecuteJob(WorkContext context) { diff --git a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs b/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs index e032a7363..d55b66e43 100644 --- a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs +++ b/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs @@ -5,7 +5,7 @@ namespace Refresh.Interfaces.Workers; public class RequestStatisticSubmitJob : WorkerJob { - public override int WorkInterval => 5_000; + public override int Interval => 5_000; public override void ExecuteJob(WorkContext context) { diff --git a/Refresh.Workers/StartupJob.cs b/Refresh.Workers/StartupJob.cs new file mode 100644 index 000000000..6f78d4065 --- /dev/null +++ b/Refresh.Workers/StartupJob.cs @@ -0,0 +1,17 @@ +using System.Diagnostics; + +namespace Refresh.Workers; + +public abstract class StartupJob : WorkerJob +{ + public override int Interval => int.MaxValue; + public sealed override void ExecuteJob(WorkContext context) + { + if (!this.FirstCycle) + throw new UnreachableException("Invoked startup job twice"); + + this.ExecuteStartupJob(context); + } + + protected abstract void ExecuteStartupJob(WorkContext context); +} \ No newline at end of file diff --git a/Refresh.Workers/WorkerJob.cs b/Refresh.Workers/WorkerJob.cs index 52ab3450b..910bb6e3f 100644 --- a/Refresh.Workers/WorkerJob.cs +++ b/Refresh.Workers/WorkerJob.cs @@ -5,7 +5,7 @@ public abstract class WorkerJob /// /// How often to perform work, in milliseconds /// - public abstract int WorkInterval { get; } + public abstract int Interval { get; } public bool FirstCycle { get; internal set; } = true; diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index 48abaa19d..5f994dafb 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -21,45 +21,45 @@ public WorkerManager(Logger logger, IDataStore dataStore, GameDatabaseProvider d private Thread? _thread = null; private bool _threadShouldRun = false; - private readonly List _workers = []; - private readonly Dictionary _lastWorkTimestamps = new(); + private readonly List _jobs = []; + private readonly Dictionary _lastJobTimestamps = new(); - public void AddWorker() where TWorker : WorkerJob, new() + public void AddJob() where TJob : WorkerJob, new() { - TWorker worker = new(); - this._workers.Add(worker); + TJob worker = new(); + this._jobs.Add(worker); } - public void AddWorker(WorkerJob worker) + public void AddJob(WorkerJob worker) { - this._workers.Add(worker); + this._jobs.Add(worker); } private void RunWorkCycle() { - Lazy workContext = new(() => new WorkContext + Lazy context = new(() => new WorkContext { Database = this._databaseProvider.GetContext(), Logger = this._logger, DataStore = this._dataStore, }); - foreach (WorkerJob worker in this._workers) + foreach (WorkerJob job in this._jobs) { long now = DateTimeOffset.Now.ToUnixTimeMilliseconds(); - if (this._lastWorkTimestamps.TryGetValue(worker, out long lastWork)) + if (this._lastJobTimestamps.TryGetValue(job, out long lastJob)) { - if(now - lastWork < worker.WorkInterval) continue; + if(now - lastJob < job.Interval) continue; - this._lastWorkTimestamps[worker] = now; + this._lastJobTimestamps[job] = now; } else { - this._lastWorkTimestamps.Add(worker, now); + this._lastJobTimestamps.Add(job, now); } - this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + worker.GetType().Name); - worker.ExecuteJob(workContext.Value); - worker.FirstCycle = false; + this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + job.GetType().Name); + job.ExecuteJob(context.Value); + job.FirstCycle = false; } } From 5aa30b76b6508e3ba076a2e570fdd3a6c62e0e20 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 17:58:51 -0400 Subject: [PATCH 05/12] Make job logic more generalizable --- .../CoolLevelsWorker.cs | 4 ++-- .../DiscordIntegrationWorker.cs | 4 ++-- .../ExpiredObjectWorker.cs | 4 ++-- .../ObjectStatisticsWorker.cs | 4 ++-- .../PunishmentExpiryWorker.cs | 4 ++-- .../RequestStatisticSubmitWorker.cs | 4 ++-- Refresh.Workers/RepeatingJob.cs | 20 +++++++++++++++++++ Refresh.Workers/StartupJob.cs | 15 ++------------ Refresh.Workers/WorkerJob.cs | 7 ++----- Refresh.Workers/WorkerManager.cs | 16 +++------------ 10 files changed, 39 insertions(+), 43 deletions(-) create mode 100644 Refresh.Workers/RepeatingJob.cs diff --git a/Refresh.Interfaces.Workers/CoolLevelsWorker.cs b/Refresh.Interfaces.Workers/CoolLevelsWorker.cs index cd17f1677..ade9d0c71 100644 --- a/Refresh.Interfaces.Workers/CoolLevelsWorker.cs +++ b/Refresh.Interfaces.Workers/CoolLevelsWorker.cs @@ -18,9 +18,9 @@ namespace Refresh.Interfaces.Workers; -public class CoolLevelsJob : WorkerJob +public class CoolLevelsJob : RepeatingJob { - public override int Interval => 600_000; // Every 10 minutes + protected override int Interval => 600_000; // Every 10 minutes [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] public override void ExecuteJob(WorkContext context) diff --git a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs index c5713ed1d..f91f12352 100644 --- a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs +++ b/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs @@ -13,7 +13,7 @@ namespace Refresh.Interfaces.Workers; -public class DiscordIntegrationJob : WorkerJob +public class DiscordIntegrationJob : RepeatingJob { private readonly IntegrationConfig _config; private readonly string _externalUrl; @@ -21,7 +21,7 @@ public class DiscordIntegrationJob : WorkerJob private long _lastTimestamp; private static long Now => DateTimeOffset.Now.ToUnixTimeMilliseconds(); - public override int Interval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default + protected override int Interval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default public DiscordIntegrationJob(IntegrationConfig config, GameServerConfig gameConfig) { diff --git a/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs b/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs index fde6c5834..18861e8b8 100644 --- a/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs +++ b/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs @@ -5,9 +5,9 @@ namespace Refresh.Interfaces.Workers; -public class CleanupExpiredObjectsJob : WorkerJob +public class CleanupExpiredObjectsJob : RepeatingJob { - public override int Interval => 60_000; // 1 minute + protected override int Interval => 60_000; // 1 minute public override void ExecuteJob(WorkContext context) { List registrationsToRemove = []; diff --git a/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs b/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs index a64032173..e75fc3638 100644 --- a/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs +++ b/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs @@ -6,9 +6,9 @@ namespace Refresh.Interfaces.Workers; -public class ObjectStatisticsJob : WorkerJob +public class ObjectStatisticsJob : RepeatingJob { - public override int Interval => 60_000; + protected override int Interval => 60_000; [SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")] public override void ExecuteJob(WorkContext context) diff --git a/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs b/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs index ae86556c4..690bf8302 100644 --- a/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs +++ b/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs @@ -8,9 +8,9 @@ namespace Refresh.Interfaces.Workers; /// /// A worker that checks all users for bans/restrictions, then removes them if expired. /// -public class PunishmentExpiryJob : WorkerJob +public class PunishmentExpiryJob : RepeatingJob { - public override int Interval => 60_000; // 1 minute + protected override int Interval => 60_000; // 1 minute public override void ExecuteJob(WorkContext context) { diff --git a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs b/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs index d55b66e43..c789e43fe 100644 --- a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs +++ b/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs @@ -3,9 +3,9 @@ namespace Refresh.Interfaces.Workers; -public class RequestStatisticSubmitJob : WorkerJob +public class RequestStatisticSubmitJob : RepeatingJob { - public override int Interval => 5_000; + protected override int Interval => 5_000; public override void ExecuteJob(WorkContext context) { diff --git a/Refresh.Workers/RepeatingJob.cs b/Refresh.Workers/RepeatingJob.cs new file mode 100644 index 000000000..3284d72ae --- /dev/null +++ b/Refresh.Workers/RepeatingJob.cs @@ -0,0 +1,20 @@ +namespace Refresh.Workers; + +public abstract class RepeatingJob : WorkerJob +{ + /// + /// How often to perform work, in milliseconds + /// + protected abstract int Interval { get; } + + private long _lastExecute = 0; + + public override bool CanExecute() + { + long now = DateTimeOffset.Now.ToUnixTimeMilliseconds(); + if(now - _lastExecute < this.Interval) return false; + + this._lastExecute = now; + return true; + } +} \ No newline at end of file diff --git a/Refresh.Workers/StartupJob.cs b/Refresh.Workers/StartupJob.cs index 6f78d4065..5310cf18d 100644 --- a/Refresh.Workers/StartupJob.cs +++ b/Refresh.Workers/StartupJob.cs @@ -1,17 +1,6 @@ -using System.Diagnostics; - -namespace Refresh.Workers; +namespace Refresh.Workers; public abstract class StartupJob : WorkerJob { - public override int Interval => int.MaxValue; - public sealed override void ExecuteJob(WorkContext context) - { - if (!this.FirstCycle) - throw new UnreachableException("Invoked startup job twice"); - - this.ExecuteStartupJob(context); - } - - protected abstract void ExecuteStartupJob(WorkContext context); + public override bool CanExecute() => this.FirstCycle; } \ No newline at end of file diff --git a/Refresh.Workers/WorkerJob.cs b/Refresh.Workers/WorkerJob.cs index 910bb6e3f..b2d8ee0de 100644 --- a/Refresh.Workers/WorkerJob.cs +++ b/Refresh.Workers/WorkerJob.cs @@ -2,13 +2,10 @@ public abstract class WorkerJob { - /// - /// How often to perform work, in milliseconds - /// - public abstract int Interval { get; } - public bool FirstCycle { get; internal set; } = true; + public virtual bool CanExecute() => true; + /// /// Executes the job. /// diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index 5f994dafb..847e9c05f 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -22,7 +22,6 @@ public WorkerManager(Logger logger, IDataStore dataStore, GameDatabaseProvider d private bool _threadShouldRun = false; private readonly List _jobs = []; - private readonly Dictionary _lastJobTimestamps = new(); public void AddJob() where TJob : WorkerJob, new() { @@ -45,19 +44,10 @@ private void RunWorkCycle() foreach (WorkerJob job in this._jobs) { - long now = DateTimeOffset.Now.ToUnixTimeMilliseconds(); - if (this._lastJobTimestamps.TryGetValue(job, out long lastJob)) - { - if(now - lastJob < job.Interval) continue; - - this._lastJobTimestamps[job] = now; - } - else - { - this._lastJobTimestamps.Add(job, now); - } + if (!job.CanExecute()) + continue; - this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + job.GetType().Name); + this._logger.LogDebug(RefreshContext.Worker, "Running work cycle for " + job.GetType().Name); job.ExecuteJob(context.Value); job.FirstCycle = false; } From 5551a44285fa297ef1f39dc7dfcbbb9e1b30e844 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:01:12 -0400 Subject: [PATCH 06/12] Organize jobs into "repeating" folder --- Refresh.GameServer/RefreshGameServer.cs | 1 + .../{CoolLevelsWorker.cs => Repeating/CoolLevelsJob.cs} | 2 +- .../DiscordIntegrationJob.cs} | 2 +- .../{ExpiredObjectWorker.cs => Repeating/ExpiredObjectJob.cs} | 2 +- .../ObjectStatisticsJob.cs} | 2 +- .../PunishmentExpiryJob.cs} | 2 +- .../RequestStatisticSubmitJob.cs} | 2 +- RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs | 1 + 8 files changed, 8 insertions(+), 6 deletions(-) rename Refresh.Interfaces.Workers/{CoolLevelsWorker.cs => Repeating/CoolLevelsJob.cs} (99%) rename Refresh.Interfaces.Workers/{DiscordIntegrationWorker.cs => Repeating/DiscordIntegrationJob.cs} (99%) rename Refresh.Interfaces.Workers/{ExpiredObjectWorker.cs => Repeating/ExpiredObjectJob.cs} (97%) rename Refresh.Interfaces.Workers/{ObjectStatisticsWorker.cs => Repeating/ObjectStatisticsJob.cs} (95%) rename Refresh.Interfaces.Workers/{PunishmentExpiryWorker.cs => Repeating/PunishmentExpiryJob.cs} (96%) rename Refresh.Interfaces.Workers/{RequestStatisticSubmitWorker.cs => Repeating/RequestStatisticSubmitJob.cs} (88%) diff --git a/Refresh.GameServer/RefreshGameServer.cs b/Refresh.GameServer/RefreshGameServer.cs index 120b45483..afe627dad 100644 --- a/Refresh.GameServer/RefreshGameServer.cs +++ b/Refresh.GameServer/RefreshGameServer.cs @@ -31,6 +31,7 @@ using Refresh.Interfaces.Game; using Refresh.Interfaces.Internal; using Refresh.Interfaces.Workers; +using Refresh.Interfaces.Workers.Repeating; using Refresh.Workers; namespace Refresh.GameServer; diff --git a/Refresh.Interfaces.Workers/CoolLevelsWorker.cs b/Refresh.Interfaces.Workers/Repeating/CoolLevelsJob.cs similarity index 99% rename from Refresh.Interfaces.Workers/CoolLevelsWorker.cs rename to Refresh.Interfaces.Workers/Repeating/CoolLevelsJob.cs index ade9d0c71..fc8bb4966 100644 --- a/Refresh.Interfaces.Workers/CoolLevelsWorker.cs +++ b/Refresh.Interfaces.Workers/Repeating/CoolLevelsJob.cs @@ -16,7 +16,7 @@ using Refresh.Database.Models.Users; using Refresh.Workers; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Interfaces.Workers.Repeating; public class CoolLevelsJob : RepeatingJob { diff --git a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs b/Refresh.Interfaces.Workers/Repeating/DiscordIntegrationJob.cs similarity index 99% rename from Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs rename to Refresh.Interfaces.Workers/Repeating/DiscordIntegrationJob.cs index f91f12352..d109586c7 100644 --- a/Refresh.Interfaces.Workers/DiscordIntegrationWorker.cs +++ b/Refresh.Interfaces.Workers/Repeating/DiscordIntegrationJob.cs @@ -11,7 +11,7 @@ using Refresh.Database.Query; using Refresh.Workers; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Interfaces.Workers.Repeating; public class DiscordIntegrationJob : RepeatingJob { diff --git a/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs b/Refresh.Interfaces.Workers/Repeating/ExpiredObjectJob.cs similarity index 97% rename from Refresh.Interfaces.Workers/ExpiredObjectWorker.cs rename to Refresh.Interfaces.Workers/Repeating/ExpiredObjectJob.cs index 18861e8b8..86ef474e8 100644 --- a/Refresh.Interfaces.Workers/ExpiredObjectWorker.cs +++ b/Refresh.Interfaces.Workers/Repeating/ExpiredObjectJob.cs @@ -3,7 +3,7 @@ using Refresh.Database.Models.Users; using Refresh.Workers; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Interfaces.Workers.Repeating; public class CleanupExpiredObjectsJob : RepeatingJob { diff --git a/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs b/Refresh.Interfaces.Workers/Repeating/ObjectStatisticsJob.cs similarity index 95% rename from Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs rename to Refresh.Interfaces.Workers/Repeating/ObjectStatisticsJob.cs index e75fc3638..a8e8f56fe 100644 --- a/Refresh.Interfaces.Workers/ObjectStatisticsWorker.cs +++ b/Refresh.Interfaces.Workers/Repeating/ObjectStatisticsJob.cs @@ -4,7 +4,7 @@ using Refresh.Database.Models.Users; using Refresh.Workers; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Interfaces.Workers.Repeating; public class ObjectStatisticsJob : RepeatingJob { diff --git a/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs b/Refresh.Interfaces.Workers/Repeating/PunishmentExpiryJob.cs similarity index 96% rename from Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs rename to Refresh.Interfaces.Workers/Repeating/PunishmentExpiryJob.cs index 690bf8302..f2cc0e4d6 100644 --- a/Refresh.Interfaces.Workers/PunishmentExpiryWorker.cs +++ b/Refresh.Interfaces.Workers/Repeating/PunishmentExpiryJob.cs @@ -3,7 +3,7 @@ using Refresh.Database.Models.Users; using Refresh.Workers; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Interfaces.Workers.Repeating; /// /// A worker that checks all users for bans/restrictions, then removes them if expired. diff --git a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs b/Refresh.Interfaces.Workers/Repeating/RequestStatisticSubmitJob.cs similarity index 88% rename from Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs rename to Refresh.Interfaces.Workers/Repeating/RequestStatisticSubmitJob.cs index c789e43fe..201c1d111 100644 --- a/Refresh.Interfaces.Workers/RequestStatisticSubmitWorker.cs +++ b/Refresh.Interfaces.Workers/Repeating/RequestStatisticSubmitJob.cs @@ -1,7 +1,7 @@ using Refresh.Core.Metrics; using Refresh.Workers; -namespace Refresh.Interfaces.Workers; +namespace Refresh.Interfaces.Workers.Repeating; public class RequestStatisticSubmitJob : RepeatingJob { diff --git a/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs b/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs index 4acc1fd91..e12e34d7c 100644 --- a/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs +++ b/RefreshTests.GameServer/Tests/Workers/PunishmentExpiryTests.cs @@ -1,6 +1,7 @@ using NotEnoughLogs; using Refresh.Database.Models.Users; using Refresh.Interfaces.Workers; +using Refresh.Interfaces.Workers.Repeating; using RefreshTests.GameServer.Logging; using static Refresh.Database.Models.Users.GameUserRole; From ccdd459317117b7b969cfd2bea8241442e4d7428 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:08:57 -0400 Subject: [PATCH 07/12] Join thread instead of sleeping --- Refresh.Workers/WorkerManager.cs | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index 847e9c05f..cabd029e8 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -86,9 +86,6 @@ public void Stop() this._logger.LogDebug(RefreshContext.Worker, "Stopping the worker thread"); this._threadShouldRun = false; - while (this._thread.IsAlive) - { - Thread.Sleep(10); - } + this._thread.Join(); } } \ No newline at end of file From a1bc15d565d4d7df5e13813afa037568ba96e487 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:10:07 -0400 Subject: [PATCH 08/12] Add per-job exception handling --- Refresh.Workers/WorkerManager.cs | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index cabd029e8..609c204ed 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -47,9 +47,16 @@ private void RunWorkCycle() if (!job.CanExecute()) continue; - this._logger.LogDebug(RefreshContext.Worker, "Running work cycle for " + job.GetType().Name); - job.ExecuteJob(context.Value); - job.FirstCycle = false; + this._logger.LogDebug(RefreshContext.Worker, $"Running work cycle for {job.GetType().Name}"); + try + { + job.ExecuteJob(context.Value); + job.FirstCycle = false; + } + catch(Exception e) + { + this._logger.LogError(RefreshContext.Worker, $"Unhandled exception while running work cycle for {job.GetType().Name}: {e}"); + } } } From 1c4fa36f35c50e4de2cfa7a8b451627b5ff3b7af Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:35:26 -0400 Subject: [PATCH 09/12] Store information regarding worker status in database --- .../GameDatabaseContext.Workers.cs | 43 +++++++++++++++++++ Refresh.Database/GameDatabaseContext.cs | 2 + .../20250721221859_AddWorkerInfoTable.cs | 41 ++++++++++++++++++ .../GameDatabaseContextModelSnapshot.cs | 24 ++++++++++- .../Models/Workers/WorkerClass.cs | 13 ++++++ Refresh.Database/Models/Workers/WorkerInfo.cs | 21 +++++++++ Refresh.Workers/WorkerManager.cs | 31 +++++++++---- Refresh.sln.DotSettings | 1 + 8 files changed, 166 insertions(+), 10 deletions(-) create mode 100644 Refresh.Database/GameDatabaseContext.Workers.cs create mode 100644 Refresh.Database/Migrations/20250721221859_AddWorkerInfoTable.cs create mode 100644 Refresh.Database/Models/Workers/WorkerClass.cs create mode 100644 Refresh.Database/Models/Workers/WorkerInfo.cs diff --git a/Refresh.Database/GameDatabaseContext.Workers.cs b/Refresh.Database/GameDatabaseContext.Workers.cs new file mode 100644 index 000000000..3d8308d76 --- /dev/null +++ b/Refresh.Database/GameDatabaseContext.Workers.cs @@ -0,0 +1,43 @@ +using Refresh.Database.Models.Workers; + +namespace Refresh.Database; + +public partial class GameDatabaseContext // Workers +{ + public int CreateWorker() + { + DateTimeOffset now = this._time.Now; + + // remove workers with this class so there's only one of us + this.Workers.RemoveRange(w => w.Class == WorkerClass.Refresh); + + WorkerInfo worker = new() + { + Class = WorkerClass.Refresh, + CreatedAt = now, + LastContact = now, + }; + + this.Workers.Add(worker); + this.SaveChanges(); + + return worker.WorkerId; + } + + /// + /// Mark a worker as contacted. + /// + /// Our worker ID. + /// False if the worker doesn't exist, and the worker should shut down. + public bool MarkWorkerContacted(int id) + { + WorkerInfo? worker = this.Workers.FirstOrDefault(w => w.WorkerId == id); + if (worker == null) + return false; + + worker.LastContact = this._time.Now; + this.SaveChanges(); + + return true; + } +} \ No newline at end of file diff --git a/Refresh.Database/GameDatabaseContext.cs b/Refresh.Database/GameDatabaseContext.cs index 40761df71..a35fb1f56 100644 --- a/Refresh.Database/GameDatabaseContext.cs +++ b/Refresh.Database/GameDatabaseContext.cs @@ -19,6 +19,7 @@ using MongoDB.Bson; using NotEnoughLogs; using Refresh.Database.Models.Statistics; +using Refresh.Database.Models.Workers; using LogLevel = Microsoft.Extensions.Logging.LogLevel; namespace Refresh.Database; @@ -71,6 +72,7 @@ public partial class GameDatabaseContext : DbContext, IDatabaseContext internal DbSet PinProgressRelations { get; set; } internal DbSet ProfilePinRelations { get; set; } internal DbSet GameSkillRewards { get; set; } + internal DbSet Workers { get; set; } #pragma warning disable CS8618 // Non-nullable variable must contain a non-null value when exiting constructor. Consider declaring it as nullable. internal GameDatabaseContext(Logger logger, IDateTimeProvider time, IDatabaseConfig dbConfig) diff --git a/Refresh.Database/Migrations/20250721221859_AddWorkerInfoTable.cs b/Refresh.Database/Migrations/20250721221859_AddWorkerInfoTable.cs new file mode 100644 index 000000000..3076bd61a --- /dev/null +++ b/Refresh.Database/Migrations/20250721221859_AddWorkerInfoTable.cs @@ -0,0 +1,41 @@ +using System; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; + +#nullable disable + +namespace Refresh.Database.Migrations +{ + [DbContext(typeof(GameDatabaseContext))] + [Migration("20250721221859_AddWorkerInfoTable")] + /// + public partial class AddWorkerInfoTable : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "Workers", + columns: table => new + { + WorkerId = table.Column(type: "integer", nullable: false) + .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn), + Class = table.Column(type: "integer", nullable: false), + CreatedAt = table.Column(type: "timestamp with time zone", nullable: false), + LastContact = table.Column(type: "timestamp with time zone", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_Workers", x => x.WorkerId); + }); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "Workers"); + } + } +} diff --git a/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs b/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs index 490549e92..69be41cf6 100644 --- a/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs +++ b/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs @@ -18,7 +18,7 @@ protected override void BuildModel(ModelBuilder modelBuilder) { #pragma warning disable 612, 618 modelBuilder - .HasAnnotation("ProductVersion", "9.0.6") + .HasAnnotation("ProductVersion", "9.0.7") .HasAnnotation("Relational:MaxIdentifierLength", 63); NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); @@ -1437,6 +1437,28 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("QueuedRegistrations"); }); + modelBuilder.Entity("Refresh.Database.Models.Workers.WorkerInfo", b => + { + b.Property("WorkerId") + .ValueGeneratedOnAdd() + .HasColumnType("integer"); + + NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("WorkerId")); + + b.Property("Class") + .HasColumnType("integer"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("LastContact") + .HasColumnType("timestamp with time zone"); + + b.HasKey("WorkerId"); + + b.ToTable("Workers"); + }); + modelBuilder.Entity("Refresh.Database.Models.Activity.Event", b => { b.HasOne("Refresh.Database.Models.Users.GameUser", "User") diff --git a/Refresh.Database/Models/Workers/WorkerClass.cs b/Refresh.Database/Models/Workers/WorkerClass.cs new file mode 100644 index 000000000..234b16ce8 --- /dev/null +++ b/Refresh.Database/Models/Workers/WorkerClass.cs @@ -0,0 +1,13 @@ +namespace Refresh.Database.Models.Workers; + +public enum WorkerClass +{ + /// + /// A worker based on the Refresh codebase. + /// + Refresh, + /// + /// A worker based on the CwLib codebase. + /// + Craftworld, +} \ No newline at end of file diff --git a/Refresh.Database/Models/Workers/WorkerInfo.cs b/Refresh.Database/Models/Workers/WorkerInfo.cs new file mode 100644 index 000000000..614aee4ed --- /dev/null +++ b/Refresh.Database/Models/Workers/WorkerInfo.cs @@ -0,0 +1,21 @@ +namespace Refresh.Database.Models.Workers; + +/// +/// Information and metadata about a worker. +/// +public class WorkerInfo +{ + [Key, Required] public int WorkerId { get; set; } + /// + /// The class of worker this is. This determines what types of jobs this worker runs. + /// + public WorkerClass Class { get; set; } + /// + /// When did this worker come online? + /// + public DateTimeOffset CreatedAt { get; set; } + /// + /// When is the last time we heard from this worker? + /// + public DateTimeOffset LastContact { get; set; } +} \ No newline at end of file diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index 609c204ed..ebc3afa5e 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -11,18 +11,25 @@ public class WorkerManager private readonly IDataStore _dataStore; private readonly GameDatabaseProvider _databaseProvider; + private readonly int _workerId; + + private Thread? _thread = null; + private bool _threadShouldRun = false; + + private long _lastContactUpdate = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + + private readonly List _jobs = []; + public WorkerManager(Logger logger, IDataStore dataStore, GameDatabaseProvider databaseProvider) { this._dataStore = dataStore; this._databaseProvider = databaseProvider; this._logger = logger; + + using GameDatabaseContext context = this._databaseProvider.GetContext(); + this._workerId = context.CreateWorker(); } - private Thread? _thread = null; - private bool _threadShouldRun = false; - - private readonly List _jobs = []; - public void AddJob() where TJob : WorkerJob, new() { TJob worker = new(); @@ -35,12 +42,12 @@ public void AddJob(WorkerJob worker) private void RunWorkCycle() { - Lazy context = new(() => new WorkContext + WorkContext context = new() { Database = this._databaseProvider.GetContext(), Logger = this._logger, DataStore = this._dataStore, - }); + }; foreach (WorkerJob job in this._jobs) { @@ -50,7 +57,7 @@ private void RunWorkCycle() this._logger.LogDebug(RefreshContext.Worker, $"Running work cycle for {job.GetType().Name}"); try { - job.ExecuteJob(context.Value); + job.ExecuteJob(context); job.FirstCycle = false; } catch(Exception e) @@ -58,6 +65,12 @@ private void RunWorkCycle() this._logger.LogError(RefreshContext.Worker, $"Unhandled exception while running work cycle for {job.GetType().Name}: {e}"); } } + + long now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + if (now - this._lastContactUpdate < 5000) return; + + this._lastContactUpdate = now; + context.Database.MarkWorkerContacted(this._workerId); } public void Start() @@ -71,7 +84,7 @@ public void Start() try { this.RunWorkCycle(); - Thread.Sleep(100); + Thread.Sleep(1000); } catch(Exception e) { diff --git a/Refresh.sln.DotSettings b/Refresh.sln.DotSettings index baa09e3cf..71e08e2f8 100644 --- a/Refresh.sln.DotSettings +++ b/Refresh.sln.DotSettings @@ -31,6 +31,7 @@ True True True + True True True From baa02f8f4ac5c7684b2dc1303e1bc6878f446209 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:38:39 -0400 Subject: [PATCH 10/12] Shut down worker if worker is replaced --- Refresh.Workers/WorkerManager.cs | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index ebc3afa5e..06bc258ef 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -70,7 +70,13 @@ private void RunWorkCycle() if (now - this._lastContactUpdate < 5000) return; this._lastContactUpdate = now; - context.Database.MarkWorkerContacted(this._workerId); + bool updated = context.Database.MarkWorkerContacted(this._workerId); + + if (!updated) + { + this._logger.LogInfo(RefreshContext.Worker, "Worker is shutting down as we've been replaced."); + this.Stop(false); + } } public void Start() @@ -83,8 +89,8 @@ public void Start() { try { + Thread.Sleep(500); this.RunWorkCycle(); - Thread.Sleep(1000); } catch(Exception e) { @@ -100,12 +106,13 @@ public void Start() this._thread = thread; } - public void Stop() + public void Stop(bool join = true) { if (this._thread == null) return; this._logger.LogDebug(RefreshContext.Worker, "Stopping the worker thread"); this._threadShouldRun = false; - this._thread.Join(); + if(join) + this._thread.Join(); } } \ No newline at end of file From ca9c306c171cd99709797738c89bc57d5bb68049 Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:50:23 -0400 Subject: [PATCH 11/12] Add simple dedicated worker manager project --- Refresh.GameServer/RefreshGameServer.cs | 8 +---- .../RefreshWorkerManager.cs | 23 +++++++++++++ Refresh.WorkerManager/Program.cs | 34 +++++++++++++++++++ .../Refresh.WorkerManager.csproj | 17 ++++++++++ Refresh.Workers/WorkerManager.cs | 8 +++++ Refresh.sln | 8 +++++ 6 files changed, 91 insertions(+), 7 deletions(-) create mode 100644 Refresh.Interfaces.Workers/RefreshWorkerManager.cs create mode 100644 Refresh.WorkerManager/Program.cs create mode 100644 Refresh.WorkerManager/Refresh.WorkerManager.csproj diff --git a/Refresh.GameServer/RefreshGameServer.cs b/Refresh.GameServer/RefreshGameServer.cs index afe627dad..3af37fba7 100644 --- a/Refresh.GameServer/RefreshGameServer.cs +++ b/Refresh.GameServer/RefreshGameServer.cs @@ -206,13 +206,7 @@ protected override void SetupServices() protected virtual void SetupWorkers() { - this.WorkerManager = new WorkerManager(this.Logger, this._dataStore, this._databaseProvider); - - this.WorkerManager.AddJob(); - this.WorkerManager.AddJob(); - this.WorkerManager.AddJob(); - this.WorkerManager.AddJob(); - this.WorkerManager.AddJob(); + this.WorkerManager = RefreshWorkerManager.Create(this.Logger, this._dataStore, this._databaseProvider); if ((this._integrationConfig?.DiscordWebhookEnabled ?? false) && this._config != null && this._config.PermitShowingOnlineUsers) { diff --git a/Refresh.Interfaces.Workers/RefreshWorkerManager.cs b/Refresh.Interfaces.Workers/RefreshWorkerManager.cs new file mode 100644 index 000000000..a738f969d --- /dev/null +++ b/Refresh.Interfaces.Workers/RefreshWorkerManager.cs @@ -0,0 +1,23 @@ +using Bunkum.Core.Storage; +using NotEnoughLogs; +using Refresh.Database; +using Refresh.Interfaces.Workers.Repeating; +using Refresh.Workers; + +namespace Refresh.Interfaces.Workers; + +public static class RefreshWorkerManager +{ + public static WorkerManager Create(Logger logger, IDataStore dataStore, GameDatabaseProvider databaseProvider) + { + WorkerManager manager = new(logger, dataStore, databaseProvider); + + manager.AddJob(); + manager.AddJob(); + manager.AddJob(); + manager.AddJob(); + manager.AddJob(); + + return manager; + } +} \ No newline at end of file diff --git a/Refresh.WorkerManager/Program.cs b/Refresh.WorkerManager/Program.cs new file mode 100644 index 000000000..091403fde --- /dev/null +++ b/Refresh.WorkerManager/Program.cs @@ -0,0 +1,34 @@ +using Bunkum.Core; +using Bunkum.Core.Storage; +using NotEnoughLogs; +using NotEnoughLogs.Behaviour; +using Refresh.Database; +using Refresh.Database.Configuration; +using Refresh.Interfaces.Workers; +using Refresh.Workers; + +LoggerConfiguration loggerConfiguration = new() +{ + Behaviour = new QueueLoggingBehaviour(), +#if DEBUG + MaxLevel = LogLevel.Trace, +#else + MaxLevel = LogLevel.Info, +#endif +}; + +using Logger logger = new(loggerConfiguration); + +logger.LogInfo(BunkumCategory.Startup, "Starting up worker manager..."); + +using GameDatabaseProvider database = new(logger, new EmptyDatabaseConfig()); +logger.LogInfo(BunkumCategory.Startup, "Initializing database..."); +database.Initialize(); +logger.LogInfo(BunkumCategory.Startup, "Warming up database..."); +database.Warmup(); + +logger.LogInfo(BunkumCategory.Startup, "Starting worker manager!"); +WorkerManager manager = RefreshWorkerManager.Create(logger, new FileSystemDataStore(), database); + +manager.Start(); +manager.WaitForExit(); \ No newline at end of file diff --git a/Refresh.WorkerManager/Refresh.WorkerManager.csproj b/Refresh.WorkerManager/Refresh.WorkerManager.csproj new file mode 100644 index 000000000..d5f6e982f --- /dev/null +++ b/Refresh.WorkerManager/Refresh.WorkerManager.csproj @@ -0,0 +1,17 @@ + + + + Exe + net9.0 + enable + enable + true + 612,618 + + + + + + + + diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index 06bc258ef..5e79d614a 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -115,4 +115,12 @@ public void Stop(bool join = true) if(join) this._thread.Join(); } + + public void WaitForExit() + { + while(this._thread == null) + Thread.Sleep(20); + + this._thread.Join(); + } } \ No newline at end of file diff --git a/Refresh.sln b/Refresh.sln index cb10f4662..edffd1690 100644 --- a/Refresh.sln +++ b/Refresh.sln @@ -35,6 +35,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Interfaces", "Interfaces", EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Refresh.Workers", "Refresh.Workers\Refresh.Workers.csproj", "{F33E3D8D-7867-41C6-8E24-5BA95642D001}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Refresh.WorkerManager", "Refresh.WorkerManager\Refresh.WorkerManager.csproj", "{DFB85FD4-CB9D-48B4-A2B7-41A592136568}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -117,6 +119,12 @@ Global {F33E3D8D-7867-41C6-8E24-5BA95642D001}.Release|Any CPU.Build.0 = Release|Any CPU {F33E3D8D-7867-41C6-8E24-5BA95642D001}.DebugLocalBunkum|Any CPU.ActiveCfg = Debug|Any CPU {F33E3D8D-7867-41C6-8E24-5BA95642D001}.DebugLocalBunkum|Any CPU.Build.0 = Debug|Any CPU + {DFB85FD4-CB9D-48B4-A2B7-41A592136568}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {DFB85FD4-CB9D-48B4-A2B7-41A592136568}.Debug|Any CPU.Build.0 = Debug|Any CPU + {DFB85FD4-CB9D-48B4-A2B7-41A592136568}.Release|Any CPU.ActiveCfg = Release|Any CPU + {DFB85FD4-CB9D-48B4-A2B7-41A592136568}.Release|Any CPU.Build.0 = Release|Any CPU + {DFB85FD4-CB9D-48B4-A2B7-41A592136568}.DebugLocalBunkum|Any CPU.ActiveCfg = Debug|Any CPU + {DFB85FD4-CB9D-48B4-A2B7-41A592136568}.DebugLocalBunkum|Any CPU.Build.0 = Debug|Any CPU EndGlobalSection GlobalSection(NestedProjects) = preSolution {A3D76B31-8732-4134-907B-F36773DF70D3} = {AA547E95-555C-470E-995A-D66E56660542} From 633b398edb69143843cd49ea7e5310029395cb7e Mon Sep 17 00:00:00 2001 From: jvyden Date: Mon, 21 Jul 2025 18:59:25 -0400 Subject: [PATCH 12/12] Add warning regarding worker manager project --- Refresh.WorkerManager/Program.cs | 1 + 1 file changed, 1 insertion(+) diff --git a/Refresh.WorkerManager/Program.cs b/Refresh.WorkerManager/Program.cs index 091403fde..bcf158956 100644 --- a/Refresh.WorkerManager/Program.cs +++ b/Refresh.WorkerManager/Program.cs @@ -20,6 +20,7 @@ using Logger logger = new(loggerConfiguration); logger.LogInfo(BunkumCategory.Startup, "Starting up worker manager..."); +logger.LogCritical(BunkumCategory.Startup, "The dedicated worker manager isn't complete yet! It will work in a debug setting, but do not use this in production yet."); using GameDatabaseProvider database = new(logger, new EmptyDatabaseConfig()); logger.LogInfo(BunkumCategory.Startup, "Initializing database...");