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.GameServer/RefreshGameServer.cs b/Refresh.GameServer/RefreshGameServer.cs
index caedaecb8..3af37fba7 100644
--- a/Refresh.GameServer/RefreshGameServer.cs
+++ b/Refresh.GameServer/RefreshGameServer.cs
@@ -31,8 +31,8 @@
using Refresh.Interfaces.Game;
using Refresh.Interfaces.Internal;
using Refresh.Interfaces.Workers;
-using Refresh.Interfaces.Workers.RequestTracking;
-using Refresh.Interfaces.Workers.Workers;
+using Refresh.Interfaces.Workers.Repeating;
+using Refresh.Workers;
namespace Refresh.GameServer;
@@ -206,17 +206,11 @@ protected override void SetupServices()
protected virtual void SetupWorkers()
{
- this.WorkerManager = new WorkerManager(this.Logger, this._dataStore, this._databaseProvider, this._matchService, this._guidCheckerService);
-
- this.WorkerManager.AddWorker();
- this.WorkerManager.AddWorker();
- this.WorkerManager.AddWorker();
- this.WorkerManager.AddWorker();
- this.WorkerManager.AddWorker();
+ this.WorkerManager = RefreshWorkerManager.Create(this.Logger, this._dataStore, this._databaseProvider);
if ((this._integrationConfig?.DiscordWebhookEnabled ?? false) && this._config != null && this._config.PermitShowingOnlineUsers)
{
- this.WorkerManager.AddWorker(new DiscordIntegrationWorker(this._integrationConfig, this._config));
+ this.WorkerManager.AddJob(new DiscordIntegrationJob(this._integrationConfig, this._config));
}
}
diff --git a/Refresh.Interfaces.Workers/IWorker.cs b/Refresh.Interfaces.Workers/IWorker.cs
deleted file mode 100644
index edeb96858..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(DataContext context);
-}
\ No newline at end of file
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/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.Interfaces.Workers/Workers/CoolLevelsWorker.cs b/Refresh.Interfaces.Workers/Repeating/CoolLevelsJob.cs
similarity index 96%
rename from Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs
rename to Refresh.Interfaces.Workers/Repeating/CoolLevelsJob.cs
index 6da335674..fc8bb4966 100644
--- a/Refresh.Interfaces.Workers/Workers/CoolLevelsWorker.cs
+++ b/Refresh.Interfaces.Workers/Repeating/CoolLevelsJob.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.Repeating;
-public class CoolLevelsWorker : IWorker
+public class CoolLevelsJob : RepeatingJob
{
- public int WorkInterval => 600_000; // Every 10 minutes
+ protected override int Interval => 600_000; // Every 10 minutes
[SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")]
- public void DoWork(DataContext context)
+ public override void ExecuteJob(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/Repeating/DiscordIntegrationJob.cs
similarity index 69%
rename from Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs
rename to Refresh.Interfaces.Workers/Repeating/DiscordIntegrationJob.cs
index 7a443cde4..d109586c7 100644
--- a/Refresh.Interfaces.Workers/Workers/DiscordIntegrationWorker.cs
+++ b/Refresh.Interfaces.Workers/Repeating/DiscordIntegrationJob.cs
@@ -2,29 +2,28 @@
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;
+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;
-using Refresh.Interfaces.APIv3.Endpoints.DataTypes.Response.Users.Photos;
+using Refresh.Workers;
-namespace Refresh.Interfaces.Workers.Workers;
+namespace Refresh.Interfaces.Workers.Repeating;
-public class DiscordIntegrationWorker : IWorker
+public class DiscordIntegrationJob : RepeatingJob
{
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();
- public int WorkInterval => this._config.DiscordWorkerFrequencySeconds * 1000; // 60 seconds by default
+ protected override int Interval => 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;
@@ -32,32 +31,26 @@ public DiscordIntegrationWorker(IntegrationConfig config, GameServerConfig gameC
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";
}
- 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 +85,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,12 +97,10 @@ private string GetAssetUrl(string hash)
return embed.Build();
}
- public void DoWork(DataContext context)
+ 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.Interfaces.Workers/Workers/ExpiredObjectWorker.cs b/Refresh.Interfaces.Workers/Repeating/ExpiredObjectJob.cs
similarity index 89%
rename from Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs
rename to Refresh.Interfaces.Workers/Repeating/ExpiredObjectJob.cs
index 6152ebb4f..86ef474e8 100644
--- a/Refresh.Interfaces.Workers/Workers/ExpiredObjectWorker.cs
+++ b/Refresh.Interfaces.Workers/Repeating/ExpiredObjectJob.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.Repeating;
-public class ExpiredObjectWorker : IWorker
+public class CleanupExpiredObjectsJob : RepeatingJob
{
- public int WorkInterval => 60_000; // 1 minute
- public void DoWork(DataContext context)
+ protected override int Interval => 60_000; // 1 minute
+ public override void ExecuteJob(WorkContext context)
{
List registrationsToRemove = [];
List codesToRemove = [];
diff --git a/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs b/Refresh.Interfaces.Workers/Repeating/ObjectStatisticsJob.cs
similarity index 80%
rename from Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs
rename to Refresh.Interfaces.Workers/Repeating/ObjectStatisticsJob.cs
index 4308310bb..a8e8f56fe 100644
--- a/Refresh.Interfaces.Workers/Workers/ObjectStatisticsWorker.cs
+++ b/Refresh.Interfaces.Workers/Repeating/ObjectStatisticsJob.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.Repeating;
-public class ObjectStatisticsWorker : IWorker
+public class ObjectStatisticsJob : RepeatingJob
{
- public int WorkInterval => 60_000;
+ protected override int Interval => 60_000;
[SuppressMessage("ReSharper.DPA", "DPA0005: Database issues")]
- public void DoWork(DataContext 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/Repeating/PunishmentExpiryJob.cs
similarity index 84%
rename from Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs
rename to Refresh.Interfaces.Workers/Repeating/PunishmentExpiryJob.cs
index 03de4bba7..f2cc0e4d6 100644
--- a/Refresh.Interfaces.Workers/Workers/PunishmentExpiryWorker.cs
+++ b/Refresh.Interfaces.Workers/Repeating/PunishmentExpiryJob.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.Repeating;
///
/// A worker that checks all users for bans/restrictions, then removes them if expired.
///
-public class PunishmentExpiryWorker : IWorker
+public class PunishmentExpiryJob : RepeatingJob
{
- public int WorkInterval => 60_000; // 1 minute
+ protected override int Interval => 60_000; // 1 minute
- public void DoWork(DataContext 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/Repeating/RequestStatisticSubmitJob.cs b/Refresh.Interfaces.Workers/Repeating/RequestStatisticSubmitJob.cs
new file mode 100644
index 000000000..201c1d111
--- /dev/null
+++ b/Refresh.Interfaces.Workers/Repeating/RequestStatisticSubmitJob.cs
@@ -0,0 +1,16 @@
+using Refresh.Core.Metrics;
+using Refresh.Workers;
+
+namespace Refresh.Interfaces.Workers.Repeating;
+
+public class RequestStatisticSubmitJob : RepeatingJob
+{
+ protected override int Interval => 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 1de24798c..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(DataContext context)
- {
- (int game, int api) = RequestStatisticTrackingMiddleware.SubmitAndClearRequests();
-
- context.Database.IncrementRequests(api, game);
- }
-}
\ No newline at end of file
diff --git a/Refresh.WorkerManager/Program.cs b/Refresh.WorkerManager/Program.cs
new file mode 100644
index 000000000..bcf158956
--- /dev/null
+++ b/Refresh.WorkerManager/Program.cs
@@ -0,0 +1,35 @@
+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...");
+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...");
+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/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.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
new file mode 100644
index 000000000..5310cf18d
--- /dev/null
+++ b/Refresh.Workers/StartupJob.cs
@@ -0,0 +1,6 @@
+namespace Refresh.Workers;
+
+public abstract class StartupJob : WorkerJob
+{
+ public override bool CanExecute() => this.FirstCycle;
+}
\ No newline at end of file
diff --git a/Refresh.Workers/WorkContext.cs b/Refresh.Workers/WorkContext.cs
new file mode 100644
index 000000000..6e4bc2289
--- /dev/null
+++ b/Refresh.Workers/WorkContext.cs
@@ -0,0 +1,12 @@
+using Bunkum.Core.Storage;
+using NotEnoughLogs;
+using Refresh.Database;
+
+namespace Refresh.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.Workers/WorkerJob.cs b/Refresh.Workers/WorkerJob.cs
new file mode 100644
index 000000000..b2d8ee0de
--- /dev/null
+++ b/Refresh.Workers/WorkerJob.cs
@@ -0,0 +1,14 @@
+namespace Refresh.Workers;
+
+public abstract class WorkerJob
+{
+ public bool FirstCycle { get; internal set; } = true;
+
+ public virtual bool CanExecute() => true;
+
+ ///
+ /// 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 51%
rename from Refresh.Interfaces.Workers/WorkerManager.cs
rename to Refresh.Workers/WorkerManager.cs
index 91758bf35..5e79d614a 100644
--- a/Refresh.Interfaces.Workers/WorkerManager.cs
+++ b/Refresh.Workers/WorkerManager.cs
@@ -1,73 +1,81 @@
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
{
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)
+ 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;
- this._matchService = matchService;
- this._guidCheckerService = guidCheckerService;
+
+ using GameDatabaseContext context = this._databaseProvider.GetContext();
+ this._workerId = context.CreateWorker();
}
- private Thread? _thread = null;
- private bool _threadShouldRun = false;
-
- private readonly List _workers = [];
- private readonly Dictionary _lastWorkTimestamps = new();
-
- public void AddWorker() where TWorker : IWorker, 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(IWorker worker)
+ public void AddJob(WorkerJob worker)
{
- this._workers.Add(worker);
+ this._jobs.Add(worker);
}
private void RunWorkCycle()
{
- Lazy dataContext = new(() => new DataContext
+ WorkContext context = new()
{
Database = this._databaseProvider.GetContext(),
Logger = this._logger,
DataStore = this._dataStore,
- Match = this._matchService,
- Token = null,
- GuidChecker = this._guidCheckerService,
- });
+ };
- foreach (IWorker worker in this._workers)
+ foreach (WorkerJob job in this._jobs)
{
- long now = DateTimeOffset.Now.ToUnixTimeMilliseconds();
- if (this._lastWorkTimestamps.TryGetValue(worker, out long lastWork))
+ if (!job.CanExecute())
+ continue;
+
+ this._logger.LogDebug(RefreshContext.Worker, $"Running work cycle for {job.GetType().Name}");
+ try
{
- if(now - lastWork < worker.WorkInterval) continue;
-
- this._lastWorkTimestamps[worker] = now;
+ job.ExecuteJob(context);
+ job.FirstCycle = false;
}
- else
+ catch(Exception e)
{
- this._lastWorkTimestamps.Add(worker, now);
+ this._logger.LogError(RefreshContext.Worker, $"Unhandled exception while running work cycle for {job.GetType().Name}: {e}");
}
-
- this._logger.LogTrace(RefreshContext.Worker, "Running work cycle for " + worker.GetType().Name);
- worker.DoWork(dataContext.Value);
+ }
+
+ long now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds();
+ if (now - this._lastContactUpdate < 5000) return;
+
+ this._lastContactUpdate = now;
+ 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);
}
}
@@ -81,8 +89,8 @@ public void Start()
{
try
{
+ Thread.Sleep(500);
this.RunWorkCycle();
- Thread.Sleep(100);
}
catch(Exception e)
{
@@ -98,15 +106,21 @@ 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;
- while (this._thread.IsAlive)
- {
- Thread.Sleep(10);
- }
+ 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 111f272f3..edffd1690 100644
--- a/Refresh.sln
+++ b/Refresh.sln
@@ -33,6 +33,10 @@ 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
+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
@@ -109,6 +113,18 @@ 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
+ {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}
@@ -121,5 +137,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/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
diff --git a/RefreshTests.GameServer/TestContext.cs b/RefreshTests.GameServer/TestContext.cs
index e139d6494..3563eb873 100644
--- a/RefreshTests.GameServer/TestContext.cs
+++ b/RefreshTests.GameServer/TestContext.cs
@@ -12,6 +12,8 @@
using Refresh.Database.Models.Levels.Scores;
using Refresh.Database.Models.Levels;
using Refresh.Interfaces.Game.Types.UserData.Leaderboard;
+using Refresh.Interfaces.Workers;
+using Refresh.Workers;
namespace RefreshTests.GameServer;
@@ -181,6 +183,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..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.Workers;
+using Refresh.Interfaces.Workers;
+using Refresh.Interfaces.Workers.Repeating;
using RefreshTests.GameServer.Logging;
using static Refresh.Database.Models.Users.GameUserRole;
@@ -12,7 +13,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 +29,22 @@ public void BannedUsersExpire()
Assert.That(context.Database.GetAllUsersWithRole(Banned).Items, Contains.Item(user));
});
- worker.DoWork(context.GetDataContext());
- 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.GetDataContext());
+ 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 +52,11 @@ public void RestrictedUsersExpire()
context.Database.RestrictUser(user, "", DateTimeOffset.FromUnixTimeMilliseconds(1000));
Assert.That(user.Role, Is.EqualTo(Restricted));
- worker.DoWork(context.GetDataContext());
- 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.GetDataContext());
+ worker.ExecuteJob(context.GetWorkContext());
context.Database.Refresh();
user = context.Database.GetUserByObjectId(user.UserId)!;