| | | 1 | | namespace AsyncResponse.Transports.Kafka; |
| | | 2 | | |
| | | 3 | | /// <summary> |
| | | 4 | | /// Extracts the AsyncResponse correlation id from an inbound response message: first from the |
| | | 5 | | /// configured Kafka header, then from the JSON body via the configured paths (walked by the shared |
| | | 6 | | /// <see cref="CorrelationIdJsonPaths"/>). |
| | | 7 | | /// </summary> |
| | | 8 | | internal static class KafkaCorrelationIdExtractor |
| | | 9 | | { |
| | | 10 | | /// <summary>Extracts the correlation id from the supplied message.</summary> |
| | | 11 | | public static string? Extract( |
| | | 12 | | IReadOnlyList<KafkaTransportHeader> headers, |
| | | 13 | | string messageJson, |
| | | 14 | | KafkaAsyncResponseTransportOptions options) |
| | | 15 | | { |
| | 26 | 16 | | var headerName = KafkaTransportOptionsValidator.Required( |
| | 26 | 17 | | options.CorrelationIdHeader, |
| | 26 | 18 | | nameof(options.CorrelationIdHeader)); |
| | | 19 | | |
| | 26 | 20 | | var headerValue = TryReadHeader(headers, headerName); |
| | 26 | 21 | | if (!string.IsNullOrWhiteSpace(headerValue)) |
| | 6 | 22 | | return headerValue; |
| | | 23 | | |
| | 20 | 24 | | return CorrelationIdJsonPaths.Extract(messageJson, options.CorrelationIdJsonPaths); |
| | | 25 | | } |
| | | 26 | | |
| | | 27 | | internal static string? TryReadHeader(IReadOnlyList<KafkaTransportHeader> headers, string headerName) |
| | | 28 | | { |
| | | 29 | | // Kafka headers allow duplicate keys; the first match wins, mirroring broker tooling. |
| | 1110 | 30 | | foreach (var header in headers) |
| | | 31 | | { |
| | 118 | 32 | | if (StringComparer.Ordinal.Equals(header.Key, headerName)) |
| | 116 | 33 | | return header.ValueUtf8; |
| | | 34 | | } |
| | | 35 | | |
| | 379 | 36 | | return null; |
| | 116 | 37 | | } |
| | | 38 | | } |