< Summary - AsyncResponse (Release / net8.0+net10.0 / unit+integration)

Information
Class: AsyncResponse.Transports.DbCorrelationIdExtractor
Assembly: AsyncResponse.Transports.PostgreSQL
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/Transports/Shared/DbTransportShared.cs
Line coverage
100%
Covered lines: 51
Uncovered lines: 0
Coverable lines: 51
Total lines: 428
Line coverage: 100%
Branch coverage
100%
Covered branches: 46
Total branches: 46
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
Extract(...)100%1818100%
TryReadPath(...)100%1212100%
TryGetProperty(...)100%66100%
UnwrapJsonString(...)100%1010100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/Transports/Shared/DbTransportShared.cs

#LineLine coverage
 1using Microsoft.Extensions.Logging;
 2using System.Text.Json;
 3using System.Text.Json.Nodes;
 4using System.Threading.Channels;
 5
 6namespace AsyncResponse.Transports;
 7
 8// Shared source for the database-backed worker transports (PostgreSQL, SQL Server, MongoDB),
 9// mirroring src/Channels/Shared/DbChannelShared.cs: each transport csproj pulls this file in via
 10// <Compile Include="..\Shared\DbTransportShared.cs" />, so the base class compiles INTO each
 11// provider assembly against that provider's concrete seam types. The seam is bound per project
 12// with global using aliases (declared at the top of the provider's MessageDispatcher file):
 13//
 14//   DbTransportOptions          -> the provider's transport options (e.g. PostgreSqlAsyncResponseTransportOptions)
 15//   DbSubscriberOptions         -> the provider's subscriber options (e.g. PostgreSqlSubscriberOptions)
 16//   DbTransportDelivery         -> the provider's claimed-delivery type (e.g. PostgreSqlTransportDelivery)
 17//   DbSubscriberRole            -> the provider's subscriber-role enum
 18//   DbAckMode                   -> the provider's ack-mode enum
 19//   DbTransportOptionsValidator -> the provider's static options validator
 20//   DbBackgroundFailureContext  -> the provider's OnBackgroundFailure context type
 21//
 22// Because the aliases resolve to concrete sealed types at compile time, delivery calls stay
 23// direct — no interface dispatch on the per-message path. The only provider-specific inputs are
 24// three display strings supplied by the derived constructor: the provider name rendered into log
 25// messages, the queue-item noun ("row"/"document"), and the lowercase telemetry tag. Rendered log
 26// output and activity tags are byte-identical to the pre-extraction per-provider sources.
 27
 28/// <summary>
 29/// Applies acknowledgement, redelivery, and dead-letter policy to database transport deliveries:
 30/// ack-after-handler with fenced lease renewal, opt-in early ACK behind a bounded in-process
 31/// queue with drain-on-dispose, attempt-capped dead-lettering, and the consumer receive span.
 32/// Derived dispatchers supply only the provider display name, queue-item noun, and telemetry tag.
 33/// </summary>
 34internal abstract class DbMessageDispatcherBase : IAsyncDisposable
 35{
 36    private readonly Func<DbTransportDelivery, CancellationToken, Task> _handler;
 37    private readonly DbTransportOptions _options;
 38    private readonly DbSubscriberOptions _subscriberOptions;
 39    private readonly ILogger _logger;
 40    private readonly DbSubscriberRole _role;
 41    private readonly string _providerName;
 42    private readonly string _unitNoun;
 43    private readonly string _receiveActivityName;
 44    private readonly string _transportTag;
 45    private readonly string _roleTagName;
 46    private readonly string _ackModeTagName;
 47
 48    private readonly Channel<DbTransportDelivery>? _backgroundQueue;
 49    private readonly Task[]? _backgroundWorkers;
 50    private readonly CancellationTokenSource? _backgroundCts;
 51
 52    protected DbMessageDispatcherBase(
 53        Func<DbTransportDelivery, CancellationToken, Task> handler,
 54        DbTransportOptions options,
 55        DbSubscriberOptions subscriberOptions,
 56        ILogger logger,
 57        DbSubscriberRole role,
 58        string providerName,
 59        string unitNoun,
 60        string telemetryName)
 61    {
 62        DbTransportOptionsValidator.ValidateSubscriber(options, subscriberOptions, role.ToString());
 63
 64        _handler = handler;
 65        _options = options;
 66        _subscriberOptions = subscriberOptions;
 67        _logger = logger;
 68        _role = role;
 69        _providerName = providerName;
 70        _unitNoun = unitNoun;
 71        _receiveActivityName = $"asyncresponse.{telemetryName}.receive";
 72        _transportTag = telemetryName;
 73        _roleTagName = $"asyncresponse.{telemetryName}.role";
 74        _ackModeTagName = $"asyncresponse.{telemetryName}.ack_mode";
 75
 76        if (subscriberOptions.AckMode is DbAckMode.AckAfterEnqueue)
 77        {
 78            _backgroundQueue = Channel.CreateBounded<DbTransportDelivery>(new BoundedChannelOptions(subscriberOptions.Ba
 79            {
 80                SingleReader = false,
 81                SingleWriter = true,
 82                FullMode = BoundedChannelFullMode.Wait
 83            });
 84            _backgroundCts = new CancellationTokenSource();
 85            _backgroundWorkers = new Task[subscriberOptions.BackgroundWorkerCount];
 86            for (var i = 0; i < _backgroundWorkers.Length; i++)
 87                _backgroundWorkers[i] = Task.Run(() => BackgroundWorkerLoopAsync(_backgroundCts.Token));
 88        }
 89    }
 90
 91    /// <summary>Handles one claimed queue item.</summary>
 92    public async Task HandleAsync(DbTransportDelivery delivery, CancellationToken cancellationToken)
 93    {
 94        if (_subscriberOptions.AckMode is DbAckMode.AckAfterEnqueue)
 95        {
 96            await HandleEarlyAckAsync(delivery).ConfigureAwait(false);
 97            return;
 98        }
 99
 100        try
 101        {
 102            // While the handler runs, a fenced heartbeat keeps extending the claim's lease at
 103            // LockTimeout/2 cadence so a slow handler does not let the lock lapse and a competing
 104            // subscriber re-claim (and duplicate-process) the queue item.
 105            using var renewalCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
 106            var renewalTask = RenewLeaseLoopAsync(delivery, renewalCancellation.Token);
 107            try
 108            {
 109                await ExecuteHandlerAsync(delivery, cancellationToken).ConfigureAwait(false);
 110            }
 111            finally
 112            {
 113                renewalCancellation.Cancel();
 114                await renewalTask.ConfigureAwait(false);
 115            }
 116
 117            await delivery.AckAsync().ConfigureAwait(false);
 118        }
 119        catch (Exception ex)
 120        {
 121            await HandleFailureAsync(delivery, ex, cancellationToken).ConfigureAwait(false);
 122        }
 123    }
 124
 125    private async Task RenewLeaseLoopAsync(DbTransportDelivery delivery, CancellationToken cancellationToken)
 126    {
 127        var interval = TimeSpan.FromTicks(Math.Max(1, _options.LockTimeout.Ticks / 2));
 128        try
 129        {
 130            while (true)
 131            {
 132                await Task.Delay(interval, cancellationToken).ConfigureAwait(false);
 133
 134                bool renewed;
 135                try
 136                {
 137                    renewed = await delivery.RenewAsync().ConfigureAwait(false);
 138                }
 139                catch (Exception ex)
 140                {
 141                    _logger.LogWarning(
 142                        ex,
 143                        "Failed to renew the lease of {Provider} message {MessageId} on queue {Queue} ({Role}); retrying
 144                        _providerName,
 145                        delivery.Id,
 146                        delivery.Queue,
 147                        _role);
 148                    continue;
 149                }
 150
 151                if (!renewed)
 152                {
 153                    // The lock_id fence no longer matches: the lease expired and another subscriber
 154                    // claimed the queue item. Stop renewing; the fenced ack/NAK will no-op for this
 155                    // claim.
 156                    _logger.LogWarning(
 157                        "Lease of {Provider} message {MessageId} on queue {Queue} ({Role}) was lost; another subscriber 
 158                        _providerName,
 159                        delivery.Id,
 160                        delivery.Queue,
 161                        _role);
 162                    return;
 163                }
 164            }
 165        }
 166        catch (OperationCanceledException)
 167        {
 168            // The handler finished or the subscriber is stopping.
 169        }
 170    }
 171
 172    // Single choke point for handler execution so both ACK modes emit the consumer receive span.
 173    private async Task ExecuteHandlerAsync(DbTransportDelivery delivery, CancellationToken cancellationToken)
 174    {
 175        using var activity = AsyncResponseDiagnostics.StartActivity(
 176            _receiveActivityName,
 177            System.Diagnostics.ActivityKind.Consumer);
 178        activity?.SetTag("asyncresponse.transport", _transportTag);
 179        activity?.SetTag(_roleTagName, _role.ToString());
 180        activity?.SetTag(_ackModeTagName, _subscriberOptions.AckMode.ToString());
 181        activity?.SetTag("messaging.system", _transportTag);
 182        activity?.SetTag("messaging.destination.name", delivery.Queue);
 183        activity?.SetTag("messaging.message.id", delivery.Id.ToString());
 184        activity?.SetTag("messaging.message.delivery_attempt", delivery.Attempt);
 185
 186        if (delivery.Headers.TryGetValue(_options.CorrelationIdHeader, out var correlationId))
 187            AsyncResponseDiagnostics.SetCorrelationId(activity, correlationId);
 188
 189        try
 190        {
 191            await _handler(delivery, cancellationToken).ConfigureAwait(false);
 192        }
 193        catch (Exception ex)
 194        {
 195            AsyncResponseDiagnostics.SetError(activity, ex);
 196            throw;
 197        }
 198    }
 199
 200    private async Task HandleEarlyAckAsync(DbTransportDelivery delivery)
 201    {
 202        if (_backgroundQueue!.Writer.TryWrite(delivery))
 203        {
 204            await delivery.AckAsync().ConfigureAwait(false);
 205        }
 206        else
 207        {
 208            _logger.LogDebug("Background queue full for {Provider} {Role}; releasing {Unit} for redelivery.", _providerN
 209            await delivery.NakAsync(_subscriberOptions.RedeliveryDelay).ConfigureAwait(false);
 210        }
 211    }
 212
 213    private async Task BackgroundWorkerLoopAsync(CancellationToken cancellationToken)
 214    {
 215        // Token-less ReadAllAsync: on shutdown the queue is completed and fully drained, so every
 216        // already-ACKed queue item is attempted (with the drain token once the drain budget lapses)
 217        // instead of being silently dropped; each failure is dead-lettered and surfaced below.
 218        await foreach (var delivery in _backgroundQueue!.Reader.ReadAllAsync().ConfigureAwait(false))
 219        {
 220            try
 221            {
 222                await ExecuteHandlerAsync(delivery, cancellationToken).ConfigureAwait(false);
 223            }
 224            catch (Exception ex)
 225            {
 226                _logger.LogError(ex, "{Provider} background handler failed for {Role} on queue {Queue} after early ACK."
 227                if (!await delivery.DeadLetterAsync(ex, false, CancellationToken.None).ConfigureAwait(false))
 228                {
 229                    _logger.LogError(
 230                        "Failed to dead-letter already-ACKed {Provider} message {MessageId} on queue {Queue} ({Role}); t
 231                        _providerName,
 232                        delivery.Id,
 233                        delivery.Queue,
 234                        _role);
 235                }
 236
 237                await InvokeBackgroundFailureAsync(delivery, ex).ConfigureAwait(false);
 238            }
 239        }
 240    }
 241
 242    private async Task HandleFailureAsync(DbTransportDelivery delivery, Exception exception, CancellationToken cancellat
 243    {
 244        var maxAttempts = _subscriberOptions.MaxDeliveryAttempts;
 245        if (maxAttempts > 0 && delivery.Attempt >= maxAttempts)
 246        {
 247            _logger.LogError(
 248                exception,
 249                "{Provider} message on queue {Queue} ({Role}) failed after {Attempts} attempts; dead-lettering.",
 250                _providerName,
 251                delivery.Queue,
 252                _role,
 253                delivery.Attempt);
 254
 255            var deadLettered = await delivery.DeadLetterAsync(exception, true, cancellationToken).ConfigureAwait(false);
 256            if (!deadLettered)
 257            {
 258                _logger.LogWarning(exception, "{Provider} dead-letter publish failed for queue {Queue} ({Role}); releasi
 259                await delivery.NakAsync(_subscriberOptions.RedeliveryDelay).ConfigureAwait(false);
 260            }
 261        }
 262        else
 263        {
 264            _logger.LogWarning(
 265                exception,
 266                "{Provider} message on queue {Queue} ({Role}) failed on attempt {Attempt}; releasing for redelivery.",
 267                _providerName,
 268                delivery.Queue,
 269                _role,
 270                delivery.Attempt);
 271            await delivery.NakAsync(_subscriberOptions.RedeliveryDelay).ConfigureAwait(false);
 272        }
 273    }
 274
 275    private async Task InvokeBackgroundFailureAsync(DbTransportDelivery delivery, Exception exception)
 276    {
 277        if (_subscriberOptions.OnBackgroundFailure is null)
 278            return;
 279
 280        try
 281        {
 282            delivery.Headers.TryGetValue(_options.CorrelationIdHeader, out var correlationId);
 283            var context = new DbBackgroundFailureContext(delivery.Queue, _role.ToString(), delivery.Attempt, correlation
 284            await _subscriberOptions.OnBackgroundFailure(context).ConfigureAwait(false);
 285        }
 286        catch (Exception ex)
 287        {
 288            _logger.LogError(ex, "{Provider} OnBackgroundFailure callback threw for {Role}.", _providerName, _role);
 289        }
 290    }
 291
 292    /// <inheritdoc />
 293    public async ValueTask DisposeAsync()
 294    {
 295        if (_backgroundQueue is null)
 296            return;
 297
 298        _backgroundQueue.Writer.TryComplete();
 299        try
 300        {
 301            await Task.WhenAll(_backgroundWorkers!).WaitAsync(_subscriberOptions.BackgroundDrainTimeout).ConfigureAwait(
 302            _backgroundCts!.Dispose();
 303        }
 304        catch (TimeoutException)
 305        {
 306            _logger.LogWarning("{Provider} background handlers for {Role} did not drain within {Timeout}.", _providerNam
 307            await _backgroundCts!.CancelAsync().ConfigureAwait(false);
 308
 309            // The workers are still running and observe _backgroundCts.Token inside ReadAllAsync, so disposing
 310            // it now would throw ObjectDisposedException inside them. Dispose once they actually finish, off
 311            // the shutdown path, so the source is not leaked either.
 312            _ = Task.WhenAll(_backgroundWorkers!).ContinueWith(
 313                _ => _backgroundCts.Dispose(),
 314                CancellationToken.None,
 315                TaskContinuationOptions.ExecuteSynchronously,
 316                TaskScheduler.Default);
 317        }
 318        catch (Exception ex)
 319        {
 320            // WhenAll only completes once every worker has finished, so the source is safe to dispose here.
 321            _logger.LogDebug(ex, "{Provider} background worker drain for {Role} ended with an error.", _providerName, _r
 322            _backgroundCts!.Dispose();
 323        }
 324    }
 325}
 326
 327/// <summary>
 328/// Extracts the AsyncResponse correlation id from the queue item's metadata first, then from the
 329/// JSON response body via configured paths. Shared verbatim by the three database transports —
 330/// the header name and JSON paths both come from the aliased options type.
 331/// </summary>
 332internal static class DbCorrelationIdExtractor
 333{
 334    public static string? Extract(
 335        IReadOnlyDictionary<string, string>? headers,
 336        string messageJson,
 337        DbTransportOptions options)
 338    {
 3339        var headerName = DbTransportOptionsValidator.Required(options.CorrelationIdHeader, nameof(options.CorrelationIdH
 3340        if (headers is not null && headers.TryGetValue(headerName, out var headerValue) && !string.IsNullOrWhiteSpace(he
 3341            return headerValue;
 342
 3343        var jsonPaths = options.CorrelationIdJsonPaths;
 3344        if (jsonPaths is null || jsonPaths.Length == 0 || string.IsNullOrWhiteSpace(messageJson))
 3345            return null;
 346
 347        JsonNode? root;
 348        try
 349        {
 3350            root = JsonNode.Parse(messageJson);
 3351        }
 3352        catch (JsonException)
 353        {
 3354            return null;
 355        }
 356
 3357        if (root is null)
 3358            return null;
 359
 3360        foreach (var path in jsonPaths)
 361        {
 3362            var value = TryReadPath(root, path);
 3363            if (!string.IsNullOrWhiteSpace(value))
 3364                return value;
 365        }
 366
 3367        return null;
 3368    }
 369
 370    private static string? TryReadPath(JsonNode root, string path)
 371    {
 3372        if (string.IsNullOrWhiteSpace(path))
 3373            return null;
 374
 3375        var current = root;
 3376        foreach (var segment in path.Split('.', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries))
 377        {
 3378            current = UnwrapJsonString(current);
 3379            if (current is not JsonObject obj)
 3380                return null;
 381
 3382            current = TryGetProperty(obj, segment);
 3383            if (current is null)
 3384                return null;
 385        }
 386
 3387        current = UnwrapJsonString(current);
 3388        return current switch
 3389        {
 3390            JsonValue value when value.TryGetValue<string>(out var s) => s,
 3391            JsonValue value => value.ToString(),
 3392            _ => null
 3393        };
 394    }
 395
 396    private static JsonNode? TryGetProperty(JsonObject obj, string name)
 397    {
 3398        if (obj.TryGetPropertyValue(name, out var exact))
 3399            return exact;
 400
 3401        foreach (var property in obj)
 402        {
 2403            if (string.Equals(property.Key, name, StringComparison.OrdinalIgnoreCase))
 2404                return property.Value;
 405        }
 406
 2407        return null;
 3408    }
 409
 410    private static JsonNode? UnwrapJsonString(JsonNode? node)
 411    {
 3412        if (node is not JsonValue value || !value.TryGetValue<string>(out var text))
 3413            return node;
 414
 3415        var trimmed = text.AsSpan().TrimStart();
 3416        if (trimmed.Length == 0 || (trimmed[0] != '{' && trimmed[0] != '['))
 3417            return node;
 418
 419        try
 420        {
 2421            return JsonNode.Parse(text);
 422        }
 2423        catch (JsonException)
 424        {
 2425            return node;
 426        }
 2427    }
 428}