diff --git a/ECMJobRunner.WebCron/Extensions/DtoExtensions.cs b/ECMJobRunner.WebCron/Extensions/DtoExtensions.cs new file mode 100644 index 0000000..0441f38 --- /dev/null +++ b/ECMJobRunner.WebCron/Extensions/DtoExtensions.cs @@ -0,0 +1,21 @@ +using ECMJobRunner.Application.Common.Dtos; +using ECMJobRunner.Application.DEXJob.Commands; +using MediatR; + +namespace ECMJobRunner.WebCron.Extensions; + +public static class DtoExtensions +{ + public static string JobId(this CfgProfileDto profile) + { + return $"profile-{profile.Id}-{profile.ProfileName.Replace(' ', '_').ToLowerInvariant()}"; + } + + public static TriggeringDEXJobBatchCommand ToJob(this CfgProfileDto profile) + { + return new TriggeringDEXJobBatchCommand + { + ProfileId = profile.Id, + }; + } +} diff --git a/ECMJobRunner.WebCron/ProfileManager.cs b/ECMJobRunner.WebCron/ProfileManager.cs new file mode 100644 index 0000000..cd5731e --- /dev/null +++ b/ECMJobRunner.WebCron/ProfileManager.cs @@ -0,0 +1,59 @@ +using ECMJobRunner.Application.Common.Dtos; +using ECMJobRunner.Application.DEXJob.Commands; +using ECMJobRunner.Application.DEXJob.Queries; +using ECMJobRunner.WebCron.Extensions; +using Hangfire; +using MediatR; +using System.Collections.Concurrent; + +namespace ECMJobRunner.WebCron +{ + public class ProfileManager(ILogger Logger, IMediator Mediator, IRecurringJobManager JobManager) : BackgroundService + { + private readonly ConcurrentDictionary Profiles = new(); + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + + while (!stoppingToken.IsCancellationRequested) + { + if (Logger.IsEnabled(LogLevel.Information)) + { + Logger.LogInformation("Worker running at: {time}", DateTimeOffset.Now); + } + + var profiles = await Mediator.Send(new GetProfileQuery() + { + Active = true, + IncludeSqlJobs = true + }, stoppingToken); + + foreach (var profile in profiles) + { + if (Profiles.TryGetValue(profile.JobId(), out var currentProfile) + && currentProfile.Schedule == profile.Schedule) + continue; + + // Add or update recurring job using MediatR command + JobManager.AddOrUpdate( + profile.JobId(), + mediator => mediator.Send(profile.ToJob(), CancellationToken.None), + profile.Schedule, + new RecurringJobOptions + { + TimeZone = TimeZoneInfo.Local + } + ); + + // Store/update in local cache + Profiles[profile.JobId()] = profile; + + Logger.LogInformation("Job {JobId} registered with schedule: {Schedule}", + profile.JobId(), profile.Schedule); + } + + await Task.Delay(1000, stoppingToken); + } + } + } +}