| | | 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.Threading; |
| | | 6 | | using System.Threading.Tasks; |
| | | 7 | | using Confluent.Kafka; |
| | | 8 | | |
| | | 9 | | namespace CoreWCF.Channels; |
| | | 10 | | |
| | | 11 | | internal class KafkaReceiveContext : ReceiveContext |
| | | 12 | | { |
| | | 13 | | private readonly ConsumeResult<byte[], byte[]> _consumeResult; |
| | | 14 | | private readonly KafkaTransportPump _kafkaTransportPump; |
| | | 15 | | |
| | 713 | 16 | | public KafkaReceiveContext(ConsumeResult<byte[], byte[]> consumeResult, KafkaTransportPump kafkaTransportPump) |
| | | 17 | | { |
| | 713 | 18 | | _consumeResult = consumeResult; |
| | 713 | 19 | | _kafkaTransportPump = kafkaTransportPump; |
| | 713 | 20 | | _kafkaTransportPump.IncrementReceiveContextCount(); |
| | 713 | 21 | | } |
| | | 22 | | |
| | | 23 | | protected override async Task OnAbandonAsync(CancellationToken token) |
| | | 24 | | { |
| | | 25 | | try |
| | | 26 | | { |
| | 6 | 27 | | if (_kafkaTransportPump.TransportBindingElement.ErrorHandlingStrategy == KafkaErrorHandlingStrategy.DeadLett |
| | | 28 | | { |
| | 2 | 29 | | await _kafkaTransportPump.Producer.ProduceAsync(_kafkaTransportPump.TransportBindingElement.DeadLetterQu |
| | | 30 | | } |
| | | 31 | | |
| | 6 | 32 | | if (_kafkaTransportPump.TransportBindingElement.DeliverySemantics == KafkaDeliverySemantics.AtLeastOnce) |
| | | 33 | | { |
| | 2 | 34 | | _kafkaTransportPump.OffsetTracker.MarkAsProcessed(_consumeResult); |
| | | 35 | | } |
| | 6 | 36 | | } |
| | | 37 | | finally |
| | | 38 | | { |
| | 6 | 39 | | _kafkaTransportPump.DecrementReceiveContextCount(); |
| | | 40 | | } |
| | 6 | 41 | | } |
| | | 42 | | |
| | | 43 | | protected override Task OnCompleteAsync(CancellationToken token) |
| | | 44 | | { |
| | | 45 | | try |
| | | 46 | | { |
| | 707 | 47 | | if (_kafkaTransportPump.TransportBindingElement.DeliverySemantics == KafkaDeliverySemantics.AtLeastOnce) |
| | | 48 | | { |
| | 441 | 49 | | _kafkaTransportPump.OffsetTracker.MarkAsProcessed(_consumeResult); |
| | | 50 | | } |
| | 707 | 51 | | } |
| | | 52 | | finally |
| | | 53 | | { |
| | 707 | 54 | | _kafkaTransportPump.DecrementReceiveContextCount(); |
| | 707 | 55 | | } |
| | | 56 | | |
| | 707 | 57 | | return Task.CompletedTask; |
| | | 58 | | } |
| | | 59 | | } |