| | | 1 | | // Licensed to the .NET Foundation under one or more agreements. |
| | | 2 | | // The .NET Foundation licenses this file to you under the MIT license. |
| | | 3 | | |
| | | 4 | | using System; |
| | | 5 | | using System.Collections.Generic; |
| | | 6 | | using System.Threading; |
| | | 7 | | using System.Threading.Tasks; |
| | | 8 | | |
| | | 9 | | namespace CoreWCF.Queue.Common |
| | | 10 | | { |
| | | 11 | | /// <summary> |
| | | 12 | | /// DefaultQueueTransportPump model serves the purpose of pull model |
| | | 13 | | /// </summary> |
| | | 14 | | internal class DefaultQueueTransportPump : QueueTransportPump |
| | | 15 | | { |
| | | 16 | | private readonly IQueueTransport _transport; |
| | 3 | 17 | | private readonly CancellationTokenSource _cancellationTokenSource = new(); |
| | 3 | 18 | | private readonly List<Task> _tasks = new(); |
| | | 19 | | |
| | 3 | 20 | | public DefaultQueueTransportPump(IQueueTransport queueTransport) |
| | | 21 | | { |
| | 3 | 22 | | _transport = queueTransport; |
| | 3 | 23 | | } |
| | | 24 | | |
| | | 25 | | public override Task StartPumpAsync(QueueTransportContext queueTransportContext, CancellationToken token) |
| | | 26 | | { |
| | 12 | 27 | | for (int i = 0; i < _transport.ConcurrencyLevel; i++) |
| | | 28 | | { |
| | 3 | 29 | | _tasks.Add(FetchAndProcessAsync(queueTransportContext, _cancellationTokenSource.Token)); |
| | | 30 | | } |
| | | 31 | | |
| | 3 | 32 | | return Task.CompletedTask; |
| | | 33 | | } |
| | | 34 | | |
| | | 35 | | public override Task StopPumpAsync(CancellationToken token) |
| | | 36 | | { |
| | 3 | 37 | | _cancellationTokenSource.Cancel(); |
| | 3 | 38 | | if (_transport is IDisposable disposable) |
| | | 39 | | { |
| | 0 | 40 | | disposable.Dispose(); |
| | | 41 | | } |
| | | 42 | | |
| | 3 | 43 | | return Task.WhenAll(_tasks); |
| | | 44 | | } |
| | | 45 | | |
| | | 46 | | private async Task FetchAndProcessAsync(QueueTransportContext queueTransportContext, CancellationToken token) |
| | | 47 | | { |
| | 3 | 48 | | CancellationTokenSource cts = new(); |
| | 3 | 49 | | TimeSpan receiveTimeout = queueTransportContext.ServiceDispatcher.Binding.ReceiveTimeout; |
| | 300 | 50 | | while (!token.IsCancellationRequested) |
| | | 51 | | { |
| | 297 | 52 | | cts.CancelAfter(receiveTimeout); |
| | | 53 | | QueueMessageContext queueMessageContext; |
| | | 54 | | try |
| | | 55 | | { |
| | 297 | 56 | | using var linkedCts = |
| | 297 | 57 | | CancellationTokenSource.CreateLinkedTokenSource(cts.Token, token); |
| | 297 | 58 | | queueMessageContext = await _transport.ReceiveQueueMessageContextAsync(linkedCts.Token); |
| | 196 | 59 | | } |
| | 101 | 60 | | catch (OperationCanceledException) |
| | | 61 | | { |
| | 101 | 62 | | cts = new(); |
| | 101 | 63 | | continue; |
| | | 64 | | } |
| | | 65 | | |
| | 196 | 66 | | if (queueMessageContext == null) |
| | | 67 | | { |
| | 98 | 68 | | cts = new(); |
| | 98 | 69 | | continue; |
| | | 70 | | } |
| | | 71 | | |
| | 98 | 72 | | cts.CancelAfter(-1); |
| | | 73 | | |
| | 98 | 74 | | queueMessageContext.QueueTransportContext = queueTransportContext; |
| | 98 | 75 | | await queueTransportContext.QueueMessageDispatcher(queueMessageContext); |
| | | 76 | | } |
| | 3 | 77 | | } |
| | | 78 | | } |
| | | 79 | | } |