/
vshmidt
/
masstransit
Обзор
Документация
Войти
/
vshmidt
/
masstransit
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
src/Transports/MassTransit.Azure.ServiceBus.Core/Scheduling/ServiceBusScheduleMessageProvider.cs
59 строк
2 KB
Chris Patterson
AddSqlMessageScheduler() method created so that scheduled messages can be canceled
17 май 2024, 17:42
17 май 2024, 17:42
40775e9
Код
Авторство
О чём код?
namespace MassTransit.Scheduling { using System; using System.Threading; using System.Threading.Tasks; using Middleware; public class ServiceBusScheduleMessageProvider : IScheduleMessageProvider { readonly ISendEndpointProvider _sendEndpointProvider; public ServiceBusScheduleMessageProvider(ISendEndpointProvider sendEndpointProvider) { _sendEndpointProvider = sendEndpointProvider; } public ServiceBusScheduleMessageProvider(ConsumeContext consumeContext) { var context = InternalOutboxExtensions.SkipOutbox(consumeContext); _sendEndpointProvider = context; } public async Task<ScheduledMessage<T>> ScheduleSend<T>(Uri destinationAddress, DateTime scheduledTime, T message, IPipe<SendContext<T>> pipe, CancellationToken cancellationToken) where T : class { if (!MessageTypeCache<T>.IsValidMessageType) throw new ArgumentException(MessageTypeCache<T>.InvalidMessageTypeReason, nameof(T)); var scheduleMessagePipe = new ScheduleSendPipe<T>(pipe, scheduledTime); var endpoint = await _sendEndpointProvider.GetSendEndpoint(destinationAddress).ConfigureAwait(false); await endpoint.Send(message, scheduleMessagePipe, cancellationToken).ConfigureAwait(false); return new ScheduledMessageHandle<T>(scheduleMessagePipe.ScheduledMessageId ?? NewId.NextGuid(), scheduledTime, destinationAddress, message); } public Task CancelScheduledSend(Guid tokenId, CancellationToken cancellationToken) { return Task.CompletedTask; } public async Task CancelScheduledSend(Uri destinationAddress, Guid tokenId, CancellationToken cancellationToken) { var endpoint = await _sendEndpointProvider.GetSendEndpoint(destinationAddress).ConfigureAwait(false); await endpoint.Send<CancelScheduledMessage>(new { InVar.CorrelationId, InVar.Timestamp, TokenId = tokenId }, cancellationToken).ConfigureAwait(false); } } }