From 89e9f16c56b2089bbd59da5b1ab9966ef34a3a0b Mon Sep 17 00:00:00 2001 From: jvyden Date: Tue, 22 Jul 2025 01:27:31 -0400 Subject: [PATCH 1/2] Support storing job state in database --- .../GameDatabaseContext.Workers.cs | 21 +++++++++++ Refresh.Database/GameDatabaseContext.cs | 1 + .../20250722045529_AddJobStateTable.cs | 36 +++++++++++++++++++ .../GameDatabaseContextModelSnapshot.cs | 14 ++++++++ .../Models/Workers/PersistentJobState.cs | 7 ++++ Refresh.Workers/State/IJobStoresState.cs | 8 +++++ Refresh.Workers/WorkerManager.cs | 14 ++++++++ 7 files changed, 101 insertions(+) create mode 100644 Refresh.Database/Migrations/20250722045529_AddJobStateTable.cs create mode 100644 Refresh.Database/Models/Workers/PersistentJobState.cs create mode 100644 Refresh.Workers/State/IJobStoresState.cs diff --git a/Refresh.Database/GameDatabaseContext.Workers.cs b/Refresh.Database/GameDatabaseContext.Workers.cs index 3d8308d76..7c50af114 100644 --- a/Refresh.Database/GameDatabaseContext.Workers.cs +++ b/Refresh.Database/GameDatabaseContext.Workers.cs @@ -40,4 +40,25 @@ public bool MarkWorkerContacted(int id) return true; } + + public object? GetJobState(string jobId, Type type) + { + PersistentJobState? state = this.JobStates.FirstOrDefault(s => s.JobId == jobId); + if (state == null) + return null; + + return JsonConvert.DeserializeObject(state.State, type); + } + + public void UpdateOrCreateJobState(string jobId, object state) + { + PersistentJobState jobState = new() + { + JobId = jobId, + State = JsonConvert.SerializeObject(state, Formatting.None), + }; + + this.JobStates.Update(jobState); + this.SaveChanges(); + } } \ No newline at end of file diff --git a/Refresh.Database/GameDatabaseContext.cs b/Refresh.Database/GameDatabaseContext.cs index a35fb1f56..3e8c909ff 100644 --- a/Refresh.Database/GameDatabaseContext.cs +++ b/Refresh.Database/GameDatabaseContext.cs @@ -73,6 +73,7 @@ public partial class GameDatabaseContext : DbContext, IDatabaseContext internal DbSet ProfilePinRelations { get; set; } internal DbSet GameSkillRewards { get; set; } internal DbSet Workers { get; set; } + internal DbSet JobStates { 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/20250722045529_AddJobStateTable.cs b/Refresh.Database/Migrations/20250722045529_AddJobStateTable.cs new file mode 100644 index 000000000..57aca6e95 --- /dev/null +++ b/Refresh.Database/Migrations/20250722045529_AddJobStateTable.cs @@ -0,0 +1,36 @@ +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Refresh.Database.Migrations +{ + [DbContext(typeof(GameDatabaseContext))] + [Migration("20250722045529_AddJobStateTable")] + /// + public partial class AddJobStateTable : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "JobStates", + columns: table => new + { + JobId = table.Column(type: "text", nullable: false), + State = table.Column(type: "jsonb", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_JobStates", x => x.JobId); + }); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "JobStates"); + } + } +} diff --git a/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs b/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs index 69be41cf6..7596bb8b4 100644 --- a/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs +++ b/Refresh.Database/Migrations/GameDatabaseContextModelSnapshot.cs @@ -1437,6 +1437,20 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("QueuedRegistrations"); }); + modelBuilder.Entity("Refresh.Database.Models.Workers.PersistentJobState", b => + { + b.Property("JobId") + .HasColumnType("text"); + + b.Property("State") + .IsRequired() + .HasColumnType("jsonb"); + + b.HasKey("JobId"); + + b.ToTable("JobStates"); + }); + modelBuilder.Entity("Refresh.Database.Models.Workers.WorkerInfo", b => { b.Property("WorkerId") diff --git a/Refresh.Database/Models/Workers/PersistentJobState.cs b/Refresh.Database/Models/Workers/PersistentJobState.cs new file mode 100644 index 000000000..6830d2e9e --- /dev/null +++ b/Refresh.Database/Models/Workers/PersistentJobState.cs @@ -0,0 +1,7 @@ +namespace Refresh.Database.Models.Workers; + +public class PersistentJobState +{ + [Key, Required] public string JobId { get; set; } = null!; + [Column(TypeName = "jsonb"), Required] public string State { get; set; } = null!; +} \ No newline at end of file diff --git a/Refresh.Workers/State/IJobStoresState.cs b/Refresh.Workers/State/IJobStoresState.cs new file mode 100644 index 000000000..bfaa12a66 --- /dev/null +++ b/Refresh.Workers/State/IJobStoresState.cs @@ -0,0 +1,8 @@ +namespace Refresh.Workers.State; + +public interface IJobStoresState +{ + public string JobId { get; set; } + public object JobState { get; protected internal set; } + public Type JobStateType { get; } +} \ No newline at end of file diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index cfcab9a8c..f71f8d9e8 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -2,6 +2,7 @@ using NotEnoughLogs; using Refresh.Core; using Refresh.Database; +using Refresh.Workers.State; namespace Refresh.Workers; @@ -55,6 +56,16 @@ private void RunWorkCycle() continue; this._logger.LogTrace(RefreshContext.Worker, $"Running work cycle for {job.GetType().Name}"); + + IJobStoresState? jobWithState = job as IJobStoresState; + if (jobWithState != null) + { + object? jobState = context.Database.GetJobState(jobWithState.JobId, jobWithState.JobStateType); + jobState ??= Activator.CreateInstance(jobWithState.JobStateType); + + jobWithState.JobState = jobState!; + } + try { job.ExecuteJob(context); @@ -64,6 +75,9 @@ private void RunWorkCycle() { this._logger.LogError(RefreshContext.Worker, $"Unhandled exception while running work cycle for {job.GetType().Name}: {e}"); } + + if (jobWithState != null) + context.Database.UpdateOrCreateJobState(jobWithState.JobId, jobWithState.JobState); } long now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); From daae2eecc8ef97ad227f97807cff62593f3053d0 Mon Sep 17 00:00:00 2001 From: jvyden Date: Tue, 22 Jul 2025 02:07:41 -0400 Subject: [PATCH 2/2] Add support for simple background migration jobs --- .../GameDatabaseContext.Workers.cs | 15 ++++-- Refresh.Workers/MigrationJob.cs | 49 +++++++++++++++++++ Refresh.Workers/Refresh.Workers.csproj | 1 + Refresh.Workers/State/IJobStoresState.cs | 4 +- Refresh.Workers/State/MigrationJobState.cs | 11 +++++ Refresh.Workers/WorkerManager.cs | 9 +++- .../GameServer/TestMigrationJob.cs | 20 ++++++++ .../Tests/Workers/MigrationTests.cs | 32 ++++++++++++ 8 files changed, 132 insertions(+), 9 deletions(-) create mode 100644 Refresh.Workers/MigrationJob.cs create mode 100644 Refresh.Workers/State/MigrationJobState.cs create mode 100644 RefreshTests.GameServer/GameServer/TestMigrationJob.cs create mode 100644 RefreshTests.GameServer/Tests/Workers/MigrationTests.cs diff --git a/Refresh.Database/GameDatabaseContext.Workers.cs b/Refresh.Database/GameDatabaseContext.Workers.cs index 7c50af114..7d22b7c6a 100644 --- a/Refresh.Database/GameDatabaseContext.Workers.cs +++ b/Refresh.Database/GameDatabaseContext.Workers.cs @@ -52,13 +52,18 @@ public bool MarkWorkerContacted(int id) public void UpdateOrCreateJobState(string jobId, object state) { - PersistentJobState jobState = new() + PersistentJobState? jobState = this.JobStates.FirstOrDefault(s => s.JobId == jobId); + if (jobState == null) { - JobId = jobId, - State = JsonConvert.SerializeObject(state, Formatting.None), - }; + jobState = new PersistentJobState + { + JobId = jobId, + }; - this.JobStates.Update(jobState); + this.JobStates.Add(jobState); + } + + jobState.State = JsonConvert.SerializeObject(state, Formatting.None); this.SaveChanges(); } } \ No newline at end of file diff --git a/Refresh.Workers/MigrationJob.cs b/Refresh.Workers/MigrationJob.cs new file mode 100644 index 000000000..35bd8389a --- /dev/null +++ b/Refresh.Workers/MigrationJob.cs @@ -0,0 +1,49 @@ +using Microsoft.EntityFrameworkCore.Storage; +using Refresh.Workers.State; + +namespace Refresh.Workers; + +public abstract class MigrationJob : WorkerJob, IJobStoresState where TEntity : class +{ + public virtual string JobId => this.GetType().Name; + public object JobState { get; set; } = null!; + public Type JobStateType => typeof(MigrationJobState); + + public MigrationJobState? MigrationJobState => JobState as MigrationJobState; + + protected virtual int BatchCount => 1_000; + protected virtual IQueryable SortAndFilter(IQueryable query) => query; + + public override bool CanExecute() + { + return this.MigrationJobState == null || !this.MigrationJobState.StateInitialized || !this.MigrationJobState.Complete; + } + + public override void ExecuteJob(WorkContext context) + { + IQueryable query = context.Database.Set(); + query = this.SortAndFilter(query); + + MigrationJobState state = this.MigrationJobState!; + + if (!state.StateInitialized) + { + state.Total = query.Count(); + state.StateInitialized = true; + } + + query = query.Skip(state.Processed).Take(this.BatchCount); + + using IDbContextTransaction transaction = context.Database.Database.BeginTransaction(); + + TEntity[] batch = query.ToArray(); + + Migrate(context, batch); + context.Database.SaveChanges(); + transaction.Commit(); + + state.Processed += batch.Length; + } + + protected abstract void Migrate(WorkContext context, TEntity[] batch); +} \ No newline at end of file diff --git a/Refresh.Workers/Refresh.Workers.csproj b/Refresh.Workers/Refresh.Workers.csproj index 4a42bbe59..eb7166453 100644 --- a/Refresh.Workers/Refresh.Workers.csproj +++ b/Refresh.Workers/Refresh.Workers.csproj @@ -11,6 +11,7 @@ + diff --git a/Refresh.Workers/State/IJobStoresState.cs b/Refresh.Workers/State/IJobStoresState.cs index bfaa12a66..8f4cefed9 100644 --- a/Refresh.Workers/State/IJobStoresState.cs +++ b/Refresh.Workers/State/IJobStoresState.cs @@ -2,7 +2,7 @@ public interface IJobStoresState { - public string JobId { get; set; } - public object JobState { get; protected internal set; } + public string JobId { get; } + public object JobState { get; set; } public Type JobStateType { get; } } \ No newline at end of file diff --git a/Refresh.Workers/State/MigrationJobState.cs b/Refresh.Workers/State/MigrationJobState.cs new file mode 100644 index 000000000..a1c95fb94 --- /dev/null +++ b/Refresh.Workers/State/MigrationJobState.cs @@ -0,0 +1,11 @@ +namespace Refresh.Workers.State; + +public class MigrationJobState +{ + public bool StateInitialized; + public int Total; + public int Processed; + + public bool Complete => this.Remaining <= 0; + public int Remaining => this.Total - this.Processed; +} \ No newline at end of file diff --git a/Refresh.Workers/WorkerManager.cs b/Refresh.Workers/WorkerManager.cs index f71f8d9e8..5152f96a2 100644 --- a/Refresh.Workers/WorkerManager.cs +++ b/Refresh.Workers/WorkerManager.cs @@ -55,8 +55,6 @@ private void RunWorkCycle() if (!job.CanExecute()) continue; - this._logger.LogTrace(RefreshContext.Worker, $"Running work cycle for {job.GetType().Name}"); - IJobStoresState? jobWithState = job as IJobStoresState; if (jobWithState != null) { @@ -64,8 +62,15 @@ private void RunWorkCycle() jobState ??= Activator.CreateInstance(jobWithState.JobStateType); jobWithState.JobState = jobState!; + + // jobs that consume state may have different execution requirements when state is updated + // check again to handle this case. the check above is still retained to avoid unnecessary db lookups + if (!job.CanExecute()) + continue; } + this._logger.LogTrace(RefreshContext.Worker, $"Running work cycle for {job.GetType().Name}"); + try { job.ExecuteJob(context); diff --git a/RefreshTests.GameServer/GameServer/TestMigrationJob.cs b/RefreshTests.GameServer/GameServer/TestMigrationJob.cs new file mode 100644 index 000000000..015dc316d --- /dev/null +++ b/RefreshTests.GameServer/GameServer/TestMigrationJob.cs @@ -0,0 +1,20 @@ +using Refresh.Database.Models.Levels; +using Refresh.Workers; + +namespace RefreshTests.GameServer.GameServer; + +public class TestMigrationJob : MigrationJob +{ + protected override void Migrate(WorkContext context, GameLevel[] batch) + { + foreach (GameLevel level in batch) + { + level.Title += " test"; + } + } + + protected override IQueryable SortAndFilter(IQueryable query) + { + return query.OrderBy(l => l.LevelId); + } +} \ No newline at end of file diff --git a/RefreshTests.GameServer/Tests/Workers/MigrationTests.cs b/RefreshTests.GameServer/Tests/Workers/MigrationTests.cs new file mode 100644 index 000000000..9e7b74658 --- /dev/null +++ b/RefreshTests.GameServer/Tests/Workers/MigrationTests.cs @@ -0,0 +1,32 @@ +using Refresh.Database.Models.Authentication; +using Refresh.Database.Models.Levels; +using Refresh.Database.Models.Users; +using Refresh.Database.Query; +using Refresh.Workers.State; + +namespace RefreshTests.GameServer.Tests.Workers; + +public class MigrationTests : GameServerTest +{ + [Test] + public void MigrationJobWorks() + { + using TestContext context = this.GetServer(); + TestMigrationJob job = new(); + GameUser user = context.CreateUser(); + + for (int i = 0; i < 100; i++) + { + context.CreateLevel(user); + } + + IEnumerable allLevels = context.Database.GetNewestLevels(100, 0, null, new LevelFilterSettings(TokenGame.Website)).Items; + Assert.That(allLevels.All(l => l.Title == "Level"), Is.True); + + job.JobState = new MigrationJobState(); + job.ExecuteJob(context.GetWorkContext()); + + allLevels = context.Database.GetNewestLevels(100, 0, null, new LevelFilterSettings(TokenGame.Website)).Items; + Assert.That(allLevels.All(l => l.Title == "Level test"), Is.True); + } +} \ No newline at end of file