| | | 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.IO.Pipelines; |
| | | 7 | | using System.Threading; |
| | | 8 | | using System.Threading.Tasks; |
| | | 9 | | using CoreWCF.Channels; |
| | | 10 | | using CoreWCF.Runtime; |
| | | 11 | | |
| | | 12 | | namespace CoreWCF.Queue.Common |
| | | 13 | | { |
| | | 14 | | public class QueueMessageContext : RequestContext |
| | | 15 | | { |
| | 3441 | 16 | | public PipeReader QueueMessageReader { get; set; } |
| | 2416 | 17 | | public virtual IDictionary<string, object> Properties { get { return _properties.Value; } } |
| | | 18 | | private Message _requestMessage; |
| | | 19 | | private Exception _requestMessageException; |
| | | 20 | | private ReceiveContext _receiveContext; |
| | 1817 | 21 | | private readonly Lazy<Dictionary<string, object>> _properties = new(); |
| | | 22 | | |
| | | 23 | | public override Message RequestMessage |
| | | 24 | | { |
| | | 25 | | get |
| | | 26 | | { |
| | 5154 | 27 | | if (_requestMessageException != null) |
| | | 28 | | { |
| | 45 | 29 | | throw DiagnosticUtility.ExceptionUtility.ThrowHelperError(_requestMessageException); |
| | | 30 | | } |
| | | 31 | | |
| | 5109 | 32 | | return _requestMessage; |
| | | 33 | | } |
| | | 34 | | } |
| | | 35 | | |
| | | 36 | | public virtual ReceiveContext ReceiveContext |
| | | 37 | | { |
| | | 38 | | get |
| | | 39 | | { |
| | 16 | 40 | | return _receiveContext; |
| | | 41 | | } |
| | | 42 | | set |
| | | 43 | | { |
| | 1719 | 44 | | _receiveContext = value; |
| | 1719 | 45 | | if (_requestMessage != null) |
| | | 46 | | { |
| | | 47 | | // Attach _receiveContext to the message |
| | 0 | 48 | | SetRequestMessage(_requestMessage); |
| | | 49 | | } |
| | 1719 | 50 | | } |
| | | 51 | | } |
| | | 52 | | |
| | | 53 | | protected virtual void OnRequestMessageSet(Message message) |
| | | 54 | | { |
| | | 55 | | |
| | 995 | 56 | | } |
| | | 57 | | |
| | | 58 | | internal void SetRequestMessage(Message requestMessage) |
| | | 59 | | { |
| | | 60 | | Fx.Assert(_requestMessageException == null, "Cannot have both a requestMessage and a requestException."); |
| | | 61 | | |
| | 1703 | 62 | | if (_receiveContext != null) |
| | | 63 | | { |
| | 1703 | 64 | | requestMessage.Properties[ReceiveContext.Name] = _receiveContext; |
| | | 65 | | } |
| | | 66 | | |
| | 1703 | 67 | | _requestMessage = requestMessage; |
| | 1703 | 68 | | OnRequestMessageSet(_requestMessage); |
| | 1703 | 69 | | } |
| | | 70 | | |
| | | 71 | | internal void SetRequestMessage(Exception requestMessageException) |
| | | 72 | | { |
| | | 73 | | Fx.Assert(_requestMessage == null, "Cannot have both a requestMessage and a requestException."); |
| | 15 | 74 | | _requestMessageException = requestMessageException; |
| | 15 | 75 | | } |
| | | 76 | | |
| | 8689 | 77 | | public QueueTransportContext QueueTransportContext { get; set; } |
| | 2969 | 78 | | public EndpointAddress LocalAddress { get; set; } |
| | | 79 | | |
| | | 80 | | public override void Abort() |
| | | 81 | | { |
| | 54 | 82 | | _requestMessage.Close(); |
| | 54 | 83 | | } |
| | | 84 | | |
| | | 85 | | public override Task ReplyAsync(Message message) |
| | | 86 | | { |
| | 1718 | 87 | | if (message != null) |
| | | 88 | | { |
| | 15 | 89 | | return message.IsFault |
| | 15 | 90 | | ? ReceiveContext.AbandonAsync(default) |
| | 15 | 91 | | : ReceiveContext.CompleteAsync(default); |
| | | 92 | | } |
| | | 93 | | |
| | 1703 | 94 | | return Task.CompletedTask; |
| | | 95 | | } |
| | | 96 | | |
| | | 97 | | public override Task ReplyAsync(Message message, CancellationToken token) |
| | | 98 | | { |
| | 0 | 99 | | return ReplyAsync(message); |
| | | 100 | | } |
| | | 101 | | |
| | | 102 | | public override Task CloseAsync() |
| | | 103 | | { |
| | 1650 | 104 | | return Task.CompletedTask; |
| | | 105 | | } |
| | | 106 | | |
| | | 107 | | public override Task CloseAsync(CancellationToken token) |
| | | 108 | | { |
| | 0 | 109 | | return Task.CompletedTask; |
| | | 110 | | } |
| | | 111 | | } |
| | | 112 | | } |