|
4 | 4 | namespace Open.Threading.Dataflow |
5 | 5 | { |
6 | 6 | internal class TargetBlockFilter<T> : TargetBlockFilterBase<T> |
7 | | - { |
| 7 | + { |
8 | 8 | private readonly DataflowMessageStatus _filterDeclineStatus; |
9 | 9 | private readonly Func<T, bool> _filter; |
10 | 10 |
|
11 | 11 | public TargetBlockFilter( |
12 | 12 | ITargetBlock<T> target, |
13 | 13 | DataflowMessageStatus filterDeclineStatus, |
14 | | - Func<T, bool> filter) :base(target) |
15 | | - { |
| 14 | + Func<T, bool> filter) : base(target) |
| 15 | + { |
16 | 16 | _filter = filter; |
| 17 | + switch (filterDeclineStatus) |
| 18 | + { |
| 19 | + case DataflowMessageStatus.Postponed: |
| 20 | + case DataflowMessageStatus.NotAvailable: |
| 21 | + throw new ArgumentException("Block filter does not support: " + Enum.GetName(typeof(DataflowMessageStatus), filterDeclineStatus)); |
| 22 | + |
| 23 | + } |
| 24 | + |
17 | 25 | _filterDeclineStatus = filterDeclineStatus; |
| 26 | + |
18 | 27 | } |
19 | 28 |
|
| 29 | + private static readonly ITargetBlock<T> NullTarget = DataflowBlock.NullTarget<T>(); |
| 30 | + |
20 | 31 | protected virtual bool Accept(T messageValue) |
21 | 32 | => _filter(messageValue); |
22 | 33 |
|
23 | 34 | public override DataflowMessageStatus OfferMessage( |
24 | 35 | DataflowMessageHeader messageHeader, T messageValue, ISourceBlock<T> source, bool consumeToAccept) |
25 | | - => Accepting && !Accept(messageValue) |
26 | | - ? _filterDeclineStatus |
27 | | - : base.OfferMessage(messageHeader, messageValue, source, consumeToAccept); |
| 36 | + { |
| 37 | + if (!Accepting) |
| 38 | + return DataflowMessageStatus.DecliningPermanently; |
| 39 | + |
| 40 | + if (Accept(messageValue)) |
| 41 | + return base.OfferMessage(messageHeader, messageValue, source, consumeToAccept); |
| 42 | + |
| 43 | + if (_filterDeclineStatus == DataflowMessageStatus.Accepted) |
| 44 | + return NullTarget.OfferMessage(messageHeader, messageValue, source, consumeToAccept); |
| 45 | + |
| 46 | + return _filterDeclineStatus; |
| 47 | + } |
28 | 48 | } |
29 | 49 |
|
| 50 | + |
30 | 51 | public static partial class DataFlowExtensions |
31 | 52 | { |
32 | 53 | /// <summary> |
|
0 commit comments