< Summary - CoreWCF Coverage — PR #1766

Information
Class: CoreWCF.Channels.MsmqQueueTransport
Assembly: CoreWCF.MSMQ
File(s): /home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.MSMQ/src/CoreWCF/Channels/MsmqQueueTransport.cs
Line coverage
0%
Covered lines: 0
Uncovered lines: 39
Coverable lines: 39
Total lines: 100
Line coverage: 0%
Branch coverage
0%
Covered branches: 0
Total branches: 2
Branch coverage: 0%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Cyclomatic complexity NPath complexity Sequence coverage
.ctor(...)100%110%
ReceiveQueueMessageContextAsync()0%220%
MessageQueueEndReceive(...)100%110%
GetContext(...)100%110%
MessageQueueBeginReceive(...)100%110%
Dispose()100%110%

File(s)

/home/runner/work/CoreWCF/CoreWCF/src/CoreWCF.MSMQ/src/CoreWCF/Channels/MsmqQueueTransport.cs

#LineLine coverage
 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
 4using System;
 5using System.IO.Pipelines;
 6using System.Threading;
 7using System.Threading.Tasks;
 8using CoreWCF.Configuration;
 9using CoreWCF.Queue.Common;
 10using Microsoft.Extensions.DependencyInjection;
 11using Microsoft.Extensions.Logging;
 12using MSMQ.Messaging;
 13using MSMQM = MSMQ.Messaging;
 14
 15namespace CoreWCF.Channels
 16{
 17    public class MsmqQueueTransport : IQueueTransport, IDisposable
 18    {
 19        private readonly Uri _baseAddress;
 20        private readonly MessageQueue _messageQueue;
 21        private readonly TimeSpan _queueReceiveTimeOut;
 22        private readonly DeadLetterQueueSender _deadLetterQueueSender;
 23        private readonly ILogger<MsmqQueueTransport> _logger;
 24
 025        public MsmqQueueTransport(IServiceDispatcher serviceDispatcher, IServiceProvider serviceProvider)
 26        {
 027            _deadLetterQueueSender = new DeadLetterQueueSender();
 028            _baseAddress = serviceDispatcher.BaseAddress;
 029            string nativeQueueName = MsmqQueueNameConverter.GetMsmqFormatQueueName(_baseAddress);
 030            _messageQueue = new MessageQueue(nativeQueueName);
 031            _queueReceiveTimeOut = serviceDispatcher.Binding.ReceiveTimeout;
 032            _logger = serviceProvider.GetRequiredService<ILogger<MsmqQueueTransport>>();
 033        }
 34
 035        public int ConcurrencyLevel => 1;
 36
 37        public async ValueTask<QueueMessageContext> ReceiveQueueMessageContextAsync(CancellationToken cancellationToken)
 38        {
 039            cancellationToken.ThrowIfCancellationRequested();
 040            _logger.LogInformation("Receiving message from msmq");
 041            var message = await Task.Factory.FromAsync(MessageQueueBeginReceive,
 042                MessageQueueEndReceive, _messageQueue, _queueReceiveTimeOut, null);
 43
 044            if(message == null)
 045                return null;
 46
 047            var reader = PipeReader.Create(message.BodyStream);
 48            try
 49            {
 050                await MsmqDecodeHelper.DecodeTransportDatagram(reader);
 051            }
 52            catch (MsmqPoisonMessageException)
 53            {
 054                await _deadLetterQueueSender.SendToSystem(reader, _baseAddress);
 055                return null;
 56            }
 57
 058            return GetContext(reader, _baseAddress);
 059        }
 60
 61        private MSMQM.Message MessageQueueEndReceive(IAsyncResult result)
 62        {
 063            MSMQM.Message message = null;
 64            try
 65            {
 066                message = _messageQueue.EndReceive(result);
 067            }
 68            catch (Exception e)
 69            {
 070                Console.WriteLine(e);
 071            }
 72
 073            return message;
 74        }
 75
 76        private QueueMessageContext GetContext(PipeReader reader, Uri uri)
 77        {
 078            var context = new QueueMessageContext
 079            {
 080                QueueMessageReader = reader,
 081                LocalAddress = new EndpointAddress(uri)
 082            };
 083            var receiveContext = new MsmqReceiveContext(context, _deadLetterQueueSender);
 084            context.ReceiveContext = receiveContext;
 85
 086            return context;
 87        }
 88
 89        private static IAsyncResult MessageQueueBeginReceive(MessageQueue messageQueue, TimeSpan timeout,
 90            AsyncCallback callback, object state)
 91        {
 092            return messageQueue.BeginReceive(timeout, state, callback);
 93        }
 94
 95        public void Dispose()
 96        {
 097            _messageQueue.Close();
 098        }
 99    }
 100}