| | | 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 Confluent.Kafka; |
| | | 7 | | |
| | | 8 | | namespace CoreWCF.Channels; |
| | | 9 | | |
| | | 10 | | public sealed class KafkaMessageProperty |
| | | 11 | | { |
| | 713 | 12 | | private readonly IList<KafkaMessageHeader> _headers = new List<KafkaMessageHeader>(); |
| | | 13 | | |
| | | 14 | | public const string Name = "CoreWCF.Channels.KafkaMessageProperty"; |
| | | 15 | | |
| | 713 | 16 | | internal KafkaMessageProperty(ConsumeResult<byte[], byte[]> consumeResult) |
| | | 17 | | { |
| | 1434 | 18 | | foreach (IHeader messageHeader in consumeResult.Message.Headers) |
| | | 19 | | { |
| | 4 | 20 | | _headers.Add(new KafkaMessageHeader(messageHeader.Key, messageHeader.GetValueBytes())); |
| | | 21 | | } |
| | | 22 | | |
| | 713 | 23 | | PartitionKey = consumeResult.Message.Key; |
| | 713 | 24 | | Topic = consumeResult.Topic; |
| | 713 | 25 | | } |
| | | 26 | | |
| | 4 | 27 | | public IReadOnlyCollection<KafkaMessageHeader> Headers => _headers as IReadOnlyCollection<KafkaMessageHeader>; |
| | 4 | 28 | | public ReadOnlyMemory<byte> PartitionKey { get; } |
| | 4 | 29 | | public string Topic { get; } |
| | | 30 | | } |