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, IServiceScopeFactory ScopeFactory, IRecurringJobManager JobManager) : BackgroundService { private readonly ConcurrentDictionary Profiles = new(); protected override async Task ExecuteAsync(CancellationToken stoppingToken) { try { while (!stoppingToken.IsCancellationRequested) { if (Logger.IsEnabled(LogLevel.Information)) { Logger.LogInformation("Worker running at: {time}", DateTimeOffset.Now); } // Create a scope to resolve scoped services (ISQLExecutor used by MediatR pipeline) using var scope = ScopeFactory.CreateScope(); var mediator = scope.ServiceProvider.GetRequiredService(); 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); } } catch (Exception ex) { Logger.LogError(ex, "An error occurred in ProfileManager."); } } } }