| | | 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.Buffers; |
| | | 6 | | using System.IO; |
| | | 7 | | using System.IO.Pipelines; |
| | | 8 | | using System.Threading.Tasks; |
| | | 9 | | using MSMQ.Messaging; |
| | | 10 | | |
| | | 11 | | namespace CoreWCF.Channels |
| | | 12 | | { |
| | | 13 | | internal class DeadLetterQueueSender |
| | | 14 | | { |
| | | 15 | | public async Task Send(PipeReader message, Uri endpoint) |
| | | 16 | | { |
| | 0 | 17 | | string nativeQueueName = MsmqQueueNameConverter.GetMsmqFormatQueueName(endpoint); |
| | 0 | 18 | | if (!MessageQueue.Exists(nativeQueueName)) |
| | | 19 | | { |
| | 0 | 20 | | MessageQueue.Create(nativeQueueName); |
| | | 21 | | } |
| | 0 | 22 | | MemoryStream memStream = await ConvertToStream(message); |
| | 0 | 23 | | var queue = new MessageQueue(nativeQueueName); |
| | 0 | 24 | | var messageForQueue = new MSMQ.Messaging.Message { BodyStream = memStream }; |
| | 0 | 25 | | queue.Send(messageForQueue); |
| | 0 | 26 | | } |
| | | 27 | | |
| | | 28 | | public async Task SendToSystem(PipeReader message, Uri endpoint) |
| | | 29 | | { |
| | 0 | 30 | | string nativeQueueName = MsmqQueueNameConverter.GetMsmqFormatQueueName(endpoint); |
| | 0 | 31 | | MemoryStream memStream = await ConvertToStream(message); |
| | 0 | 32 | | var queue = new MessageQueue(nativeQueueName); |
| | 0 | 33 | | var messageForQueue = new MSMQ.Messaging.Message |
| | 0 | 34 | | { |
| | 0 | 35 | | BodyStream = memStream, UseDeadLetterQueue = true, TimeToBeReceived = TimeSpan.FromSeconds(0), |
| | 0 | 36 | | }; |
| | 0 | 37 | | queue.Send(messageForQueue); |
| | 0 | 38 | | } |
| | | 39 | | |
| | | 40 | | private static async Task<MemoryStream> ConvertToStream(PipeReader stream) |
| | | 41 | | { |
| | 0 | 42 | | var readResult = await stream.ReadAsync(); |
| | 0 | 43 | | var memStream = new MemoryStream(readResult.Buffer.ToArray()); |
| | 0 | 44 | | return memStream; |
| | 0 | 45 | | } |
| | | 46 | | } |
| | | 47 | | } |