/
vshmidt
/
masstransit
Обзор
Документация
Войти
/
vshmidt
/
masstransit
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
src/Scheduling/MassTransit.HangfireIntegration/Configuration/HangfireIntegrationExtensions.cs
66 строк
3 KB
ushenkodmitry
Pause/Resume scheduled recurring message, Quartz, Hangfire (cleaned up by PBG)
03 дек 2022, 20:30
03 дек 2022, 20:30
5ee7c49
Код
Авторство
О чём код?
namespace MassTransit { using System; using Hangfire; using HangfireIntegration; using Scheduling; public static class HangfireIntegrationExtensions { [Obsolete("Use the new .AddHangfireConsumers() method, combined with AddHangfire(), to configure the Quartz scheduler")] public static void UseHangfireScheduler(this IBusFactoryConfigurator configurator, IBusRegistrationContext context, string queueName = "hangfire", Action<BackgroundJobServerOptions>? configureServer = null) { UseHangfireScheduler(configurator, new BusRegistrationContextComponentResolver(context, DefaultHangfireComponentResolver.Instance), queueName, configureServer); } public static void UseHangfireScheduler(this IBusFactoryConfigurator configurator, string queueName = "hangfire", Action<BackgroundJobServerOptions>? configureServer = null) { UseHangfireScheduler(configurator, DefaultHangfireComponentResolver.Instance, queueName, configureServer); } public static void UseHangfireScheduler(this IBusFactoryConfigurator configurator, IHangfireComponentResolver hangfireComponentResolver, string queueName = "hangfire", Action<BackgroundJobServerOptions>? configureServer = null) { UseHangfireScheduler(configurator, options => { options.QueueName = queueName; options.ComponentResolver = hangfireComponentResolver; options.ConfigureServer = configureServer; }); } public static void UseHangfireScheduler(this IBusFactoryConfigurator configurator, Action<HangfireSchedulerOptions>? configure) { if (configurator == null) throw new ArgumentNullException(nameof(configurator)); var options = new HangfireSchedulerOptions(); configure?.Invoke(options); configurator.ReceiveEndpoint(options.QueueName, e => { var partitioner = configurator.CreatePartitioner(Environment.ProcessorCount); e.Consumer(() => new ScheduleMessageConsumer(options.ComponentResolver.BackgroundJobClient, options.ComponentResolver.JobStorage), x => { x.Message<ScheduleMessage>(m => m.UsePartitioner(partitioner, p => p.Message.CorrelationId)); x.Message<CancelScheduledMessage>(m => m.UsePartitioner(partitioner, p => p.Message.TokenId)); }); e.Consumer(() => new ScheduleRecurringMessageConsumer(options.ComponentResolver.RecurringJobManager, options.ComponentResolver.TimeZoneResolver)); e.Consumer(() => new PauseScheduledRecurringMessageConsumer(options.ComponentResolver.RecurringJobManager, options.ComponentResolver.JobStorage)); e.Consumer(() => new ResumeScheduledRecurringMessageConsumer(options.ComponentResolver.RecurringJobManager, options.ComponentResolver.JobStorage)); var observer = new SchedulerBusObserver(options); configurator.ConnectBusObserver(observer); configurator.UseMessageScheduler(e.InputAddress); }); } } }