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

Information
Class: AsyncResponse.RecoverableAsyncResponseBuilder
Assembly: AsyncResponse.Core
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/AsyncResponse.Core/AsyncResponseBuilder.cs
Line coverage
100%
Covered lines: 5
Uncovered lines: 0
Coverable lines: 5
Total lines: 488
Line coverage: 100%
Branch coverage
N/A
Covered branches: 0
Total branches: 0
Branch coverage: N/A
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
For<T>(...)100%11100%
For<T>()100%11100%
AsyncResponse.IAsyncResponseBuilder.For<T>(...)100%11100%
AsyncResponse.IAsyncResponseBuilder.For<T>()100%11100%

File(s)

/home/runner/work/AsyncResponse/AsyncResponse/src/AsyncResponse.Core/AsyncResponseBuilder.cs

#LineLine coverage
 1using System.Diagnostics;
 2using System.Diagnostics.CodeAnalysis;
 3using System.Linq.Expressions;
 4
 5namespace AsyncResponse;
 6
 7internal abstract class AsyncResponseBuilderBase(
 8    IWorkerTransport? _workerTransport = null,
 9    IAsyncResponseReplyTargetProvider? _replyTargetProvider = null,
 10    AsyncResponseContextPropagation? _propagation = null)
 11{
 12    protected IAsyncResponseReplyTargetProvider? ReplyTargetProvider => _replyTargetProvider;
 13
 14    /// <summary>Validates the supplied options.</summary>
 15    protected static string ValidateCorrelationId(string correlationId)
 16        => !string.IsNullOrWhiteSpace(correlationId)
 17            ? correlationId
 18            : throw new ArgumentNullException(nameof(correlationId), "CorrelationId must not be empty or whitespace.");
 19
 20    /// <inheritdoc cref="IAsyncResponseBuilder.EnqueueWorkerAsync(ReflectionCallDto, CancellationToken)" />
 21    [RequiresUnreferencedCode("The descriptor names its target service and method as strings, resolved by reflection whe
 22                              "job executes; trimming may have removed them. Use the expression-based EnqueueWorkerAsync
 23                              "overloads, which root the service's public methods automatically.")]
 24    public Task EnqueueWorkerAsync(ReflectionCallDto work, CancellationToken cancellationToken = default)
 25        => EnqueueWorkerCoreAsync(work, cancellationToken);
 26
 27    // Shared by the annotation-free expression overloads (whose TService is rooted via
 28    // DynamicallyAccessedMembers) and the RequiresUnreferencedCode DTO overload above.
 29    private async Task EnqueueWorkerCoreAsync(ReflectionCallDto work, CancellationToken cancellationToken)
 30    {
 31        ArgumentNullException.ThrowIfNull(work);
 32        cancellationToken.ThrowIfCancellationRequested();
 33
 34        using var activity = AsyncResponseDiagnostics.StartActivity(
 35            "asyncresponse.enqueue_worker",
 36            ActivityKind.Producer,
 37            AsyncResponseContext.CorrelationId);
 38        AsyncResponseDiagnostics.SetWorker(activity, work);
 39        AsyncResponseDiagnostics.SetReplyTarget(activity, AsyncResponseContext.ReplyTarget);
 40
 41        try
 42        {
 43            var transport = _workerTransport ?? throw new InvalidOperationException(
 44                "No IWorkerTransport is registered. Call .WithInMemoryTransport() for in-process execution, " +
 45                ".WithGooglePubSubTransport(...) for Google Pub/Sub, " +
 46                "or install another full AsyncResponse transport package.");
 47
 48            await transport.PublishAsync(
 49                new WorkerJobEnvelope
 50                {
 51                    Call = work,
 52                    CorrelationId = AsyncResponseContext.CorrelationId,
 53                    ReplyTarget = AsyncResponseContext.ReplyTarget,
 54                    Context = _propagation?.Capture()
 55                },
 56                cancellationToken).ConfigureAwait(false);
 57        }
 58        catch (Exception ex)
 59        {
 60            AsyncResponseDiagnostics.SetError(activity, ex);
 61            throw;
 62        }
 63    }
 64
 65    /// <inheritdoc cref="IAsyncResponseBuilder.EnqueueWorkerAsync{TService}(Expression{Action{TService}}, CancellationT
 66    public Task EnqueueWorkerAsync<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicMethods)] TService>(
 67    {
 68        cancellationToken.ThrowIfCancellationRequested();
 69        return EnqueueWorkerCoreAsync(CallbackExpressionConverter.ToReflectionCall(work), cancellationToken);
 70    }
 71
 72    /// <inheritdoc cref="IAsyncResponseBuilder.EnqueueWorkerAsync{TService}(Expression{Func{TService, Task}}, Cancellat
 73    public Task EnqueueWorkerAsync<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicMethods)] TService>(
 74    {
 75        cancellationToken.ThrowIfCancellationRequested();
 76        return EnqueueWorkerCoreAsync(CallbackExpressionConverter.ToReflectionCall(work), cancellationToken);
 77    }
 78
 79    /// <inheritdoc cref="IAsyncResponseBuilder.EnqueueWorkerAsync{TService}(Expression{Func{TService, ValueTask}}, Canc
 80    public Task EnqueueWorkerAsync<[DynamicallyAccessedMembers(DynamicallyAccessedMemberTypes.PublicMethods)] TService>(
 81    {
 82        cancellationToken.ThrowIfCancellationRequested();
 83        return EnqueueWorkerCoreAsync(CallbackExpressionConverter.ToReflectionCall(work), cancellationToken);
 84    }
 85}
 86
 87/// <inheritdoc cref="IAsyncResponseBuilder"/>
 88internal sealed class AsyncResponseBuilder(
 89    IAsyncResponseSubscriber _subscriber,
 90    IWorkerTransport? workerTransport = null,
 91    IAsyncResponseReplyTargetProvider? replyTargetProvider = null,
 92    AsyncResponseContextPropagation? propagation = null)
 93    : AsyncResponseBuilderBase(workerTransport, replyTargetProvider, propagation),
 94        IAsyncResponseBuilder
 95{
 96    /// <inheritdoc />
 97    public IAsyncResponseAttachedBuilder<T> For<T>(string correlationId) where T : IAsyncResponsePayload
 98        => new AsyncResponseBuilder<T>(_subscriber, ReplyTargetProvider, ValidateCorrelationId(correlationId));
 99
 100    /// <inheritdoc />
 101    public IAsyncResponseTriggeredBuilder<T> For<T>() where T : IAsyncResponsePayload
 102        => new AsyncResponseBuilder<T>(_subscriber, ReplyTargetProvider, AsyncResponseContext.CreateCorrelationId());
 103}
 104
 105/// <inheritdoc cref="IRecoverableAsyncResponseBuilder"/>
 106internal sealed class RecoverableAsyncResponseBuilder(
 107    IRecoverableAsyncResponseSubscriber _subscriber,
 108    IWorkerTransport? workerTransport = null,
 109    IAsyncResponseReplyTargetProvider? replyTargetProvider = null,
 110    AsyncResponseContextPropagation? propagation = null)
 3111    : AsyncResponseBuilderBase(workerTransport, replyTargetProvider, propagation),
 112        IRecoverableAsyncResponseBuilder
 113{
 114    /// <inheritdoc />
 115    public IRecoverableAsyncResponseAttachedBuilder<T> For<T>(string correlationId) where T : IAsyncResponsePayload
 3116        => new RecoverableAsyncResponseBuilder<T>(_subscriber, ReplyTargetProvider, ValidateCorrelationId(correlationId)
 117
 118    /// <inheritdoc />
 119    public IRecoverableAsyncResponseTriggeredBuilder<T> For<T>() where T : IAsyncResponsePayload
 3120        => new RecoverableAsyncResponseBuilder<T>(_subscriber, ReplyTargetProvider, AsyncResponseContext.CreateCorrelati
 121
 122    IAsyncResponseAttachedBuilder<T> IAsyncResponseBuilder.For<T>(string correlationId)
 3123        => For<T>(correlationId);
 124
 125    IAsyncResponseTriggeredBuilder<T> IAsyncResponseBuilder.For<T>()
 3126        => For<T>();
 127}
 128
 129/// <inheritdoc cref="IAsyncResponseAttachedBuilder{T}" />
 130internal class AsyncResponseBuilder<T> : IAsyncResponseAttachedBuilder<T>, IAsyncResponseTriggeredBuilder<T>
 131    where T : IAsyncResponsePayload
 132{
 133    private readonly IAsyncResponseSubscriber _subscriber;
 134    private readonly IAsyncResponseReplyTargetProvider? _replyTargetProvider;
 135    protected readonly string _correlationId;
 136    protected Func<T, ValueTask<bool>>? _completionPredicate;
 137    protected TimeSpan? _timeout;
 138    private bool _useReplyTarget;
 139    private string? _replyTargetName;
 140    private AsyncResponseReplyTarget? _replyTarget;
 141    private int _consumed;
 142
 143    internal AsyncResponseBuilder(
 144        IAsyncResponseSubscriber subscriber,
 145        IAsyncResponseReplyTargetProvider? replyTargetProvider,
 146        string correlationId)
 147    {
 148        _subscriber = subscriber;
 149        _replyTargetProvider = replyTargetProvider;
 150        _correlationId = correlationId;
 151    }
 152
 153    /// <inheritdoc />
 154    public IAsyncResponseAttachedBuilder<T> WithTimeout(TimeSpan timeout)
 155    {
 156        if (timeout <= TimeSpan.Zero)
 157            throw new ArgumentOutOfRangeException(nameof(timeout), "Timeout must be greater than zero.");
 158
 159        _timeout = timeout;
 160        return this;
 161    }
 162
 163    /// <inheritdoc />
 164    public IAsyncResponseAttachedBuilder<T> WithReplyTarget()
 165    {
 166        _useReplyTarget = true;
 167        _replyTargetName = null;
 168        _replyTarget = null;
 169        return this;
 170    }
 171
 172    /// <inheritdoc />
 173    public IAsyncResponseAttachedBuilder<T> WithReplyTarget(string name)
 174    {
 175        _useReplyTarget = true;
 176        _replyTargetName = !string.IsNullOrWhiteSpace(name)
 177            ? name
 178            : throw new ArgumentException("Reply target name cannot be null or whitespace.", nameof(name));
 179        _replyTarget = null;
 180        return this;
 181    }
 182
 183    /// <inheritdoc />
 184    public IAsyncResponseAttachedBuilder<T> WithReplyTarget(AsyncResponseReplyTarget replyTarget)
 185    {
 186        ArgumentNullException.ThrowIfNull(replyTarget);
 187        ValidateReplyTarget(replyTarget);
 188
 189        _useReplyTarget = true;
 190        _replyTargetName = null;
 191        _replyTarget = replyTarget;
 192        return this;
 193    }
 194
 195    /// <inheritdoc />
 196    public IAsyncResponseAttachedBuilder<T> Until(Func<T, bool> predicate)
 197    {
 198        _completionPredicate = predicate != null
 199            ? payload => new ValueTask<bool>(predicate(payload))
 200            : throw new ArgumentNullException(nameof(predicate));
 201        return this;
 202    }
 203
 204    /// <inheritdoc />
 205    public IAsyncResponseAttachedBuilder<T> Until(Func<T, Task<bool>> predicate)
 206    {
 207        _completionPredicate = predicate != null
 208            ? payload => new ValueTask<bool>(predicate(payload))
 209            : throw new ArgumentNullException(nameof(predicate));
 210        return this;
 211    }
 212
 213    /// <inheritdoc cref="IAsyncResponseAttachedBuilder{T}.WaitAsync" />
 214    public Task<T> WaitAsync()
 215        => WaitCoreAsync((Func<AsyncResponseRequestContext, Task>?)null);
 216
 217    /// <inheritdoc cref="IAsyncResponseTriggeredBuilder{T}.WaitAsync(Func{AsyncResponseRequestContext, Task})" />
 218    public Task<T> WaitAsync(Func<AsyncResponseRequestContext, Task> trigger)
 219    {
 220        ArgumentNullException.ThrowIfNull(trigger);
 221        return WaitCoreAsync(trigger);
 222    }
 223
 224    /// <summary>Creates the requested resource.</summary>
 225    protected virtual Task<IAsyncResponseWaiter<T>> CreateWaiterAsync()
 226        => _subscriber.CreateResponseWaiter<T>(_correlationId, _completionPredicate, _timeout);
 227
 228    private async Task<T> WaitCoreAsync(Func<AsyncResponseRequestContext, Task>? trigger)
 229    {
 230        // Builders are single-use: reuse would re-register the SAME correlation id and re-fire
 231        // the trigger â€” exactly the double-send the attached/triggered split exists to prevent.
 232        if (Interlocked.Exchange(ref _consumed, 1) != 0)
 233        {
 234            throw new InvalidOperationException(
 235                "This async-response builder has already been awaited. Builders are single-use â€” " +
 236                "call For<T>() again so every wait gets its own correlation id and registration.");
 237        }
 238
 239        await using var waiter = await CreateWaiterAsync().ConfigureAwait(false);
 240        var replyTarget = ResolveReplyTarget();
 241        var requestContext = new AsyncResponseRequestContext(_correlationId, replyTarget);
 242
 243        // Subscribe-before-send by construction: the trigger runs only once the subscription and
 244        // the recovery state exist, so the first response can never race the registration. A
 245        // failing trigger means the operation never started â€” the waiter (and with it the
 246        // recovery state) is torn down by the await-using disposal as the exception propagates.
 247        if (trigger != null)
 248        {
 249            using var contextScope = AsyncResponseContext.PushContext(_correlationId, replyTarget);
 250            await trigger(requestContext).ConfigureAwait(false);
 251        }
 252
 253        return await waiter.ResponseTask.ConfigureAwait(false);
 254    }
 255
 256    private AsyncResponseReplyTarget? ResolveReplyTarget()
 257    {
 258        if (!_useReplyTarget)
 259            return null;
 260
 261        var replyTarget = _replyTarget
 262            ?? (_replyTargetProvider ?? throw new InvalidOperationException(
 263                "No async-response reply target provider is registered. Register a transport package " +
 264                "that provides reply targets, such as .WithGooglePubSubTransport(...), or pass an " +
 265                "explicit AsyncResponseReplyTarget to .WithReplyTarget(...)."))
 266            .GetReplyTarget(_replyTargetName);
 267
 268        ValidateReplyTarget(replyTarget);
 269        return replyTarget;
 270    }
 271
 272    private static void ValidateReplyTarget(AsyncResponseReplyTarget replyTarget)
 273    {
 274        ArgumentException.ThrowIfNullOrWhiteSpace(replyTarget.Name);
 275        ArgumentException.ThrowIfNullOrWhiteSpace(replyTarget.Transport);
 276        ArgumentException.ThrowIfNullOrWhiteSpace(replyTarget.Address);
 277    }
 278
 279    // -----------------------------------------------------------------------------------------
 280    // IAsyncResponseTriggeredBuilder<T> â€” the builder handed out by For<T>() (generated
 281    // correlation id). Same shared state and behavior; only the static return type differs, so
 282    // the trigger-required WaitAsync terminal is preserved through the fluent chain. The public
 283    // WaitAsync(Func<AsyncResponseRequestContext, Task>) overload above satisfies its terminal.
 284
 285    IAsyncResponseTriggeredBuilder<T> IAsyncResponseTriggeredBuilder<T>.WithTimeout(TimeSpan timeout)
 286    {
 287        WithTimeout(timeout);
 288        return this;
 289    }
 290
 291    IAsyncResponseTriggeredBuilder<T> IAsyncResponseTriggeredBuilder<T>.WithReplyTarget()
 292    {
 293        WithReplyTarget();
 294        return this;
 295    }
 296
 297    IAsyncResponseTriggeredBuilder<T> IAsyncResponseTriggeredBuilder<T>.WithReplyTarget(string name)
 298    {
 299        WithReplyTarget(name);
 300        return this;
 301    }
 302
 303    IAsyncResponseTriggeredBuilder<T> IAsyncResponseTriggeredBuilder<T>.WithReplyTarget(AsyncResponseReplyTarget replyTa
 304    {
 305        WithReplyTarget(replyTarget);
 306        return this;
 307    }
 308
 309    IAsyncResponseTriggeredBuilder<T> IAsyncResponseTriggeredBuilder<T>.Until(Func<T, bool> predicate)
 310    {
 311        Until(predicate);
 312        return this;
 313    }
 314
 315    IAsyncResponseTriggeredBuilder<T> IAsyncResponseTriggeredBuilder<T>.Until(Func<T, Task<bool>> predicate)
 316    {
 317        Until(predicate);
 318        return this;
 319    }
 320}
 321
 322/// <inheritdoc cref="IRecoverableAsyncResponseAttachedBuilder{T}" />
 323internal sealed class RecoverableAsyncResponseBuilder<T> :
 324    AsyncResponseBuilder<T>,
 325    IRecoverableAsyncResponseAttachedBuilder<T>,
 326    IRecoverableAsyncResponseTriggeredBuilder<T>
 327    where T : IAsyncResponsePayload
 328{
 329    private readonly IRecoverableAsyncResponseSubscriber _subscriber;
 330    private ReflectionCallDto? _resumeCallback;
 331    private ReflectionCallDto? _failureCallback;
 332
 333    internal RecoverableAsyncResponseBuilder(
 334        IRecoverableAsyncResponseSubscriber subscriber,
 335        IAsyncResponseReplyTargetProvider? replyTargetProvider,
 336        string correlationId)
 337        : base(subscriber, replyTargetProvider, correlationId)
 338    {
 339        _subscriber = subscriber;
 340    }
 341
 342    /// <summary>Creates the requested resource.</summary>
 343    protected override Task<IAsyncResponseWaiter<T>> CreateWaiterAsync()
 344        => _subscriber.CreateRecoverableResponseWaiter<T>(
 345            _correlationId,
 346            _resumeCallback,
 347            _failureCallback,
 348            _completionPredicate,
 349            _timeout);
 350
 351    private void SetResumeCallback(ReflectionCallDto callback)
 352        => _resumeCallback = callback ?? throw new ArgumentNullException(nameof(callback));
 353
 354    private void SetFailureCallback(ReflectionCallDto callback)
 355        => _failureCallback = callback ?? throw new ArgumentNullException(nameof(callback));
 356
 357    [RequiresUnreferencedCode("The callback names its target service and method as strings, resolved by reflection when 
 358                              "subscriber loss; trimming may have removed them. Use the expression-based overload, which
 359                              "public methods automatically.")]
 360    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.OnLostSubscriberResume(Refle
 361    {
 362        SetResumeCallback(callback);
 363        return this;
 364    }
 365
 366    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.OnLostSubscriberResume<[Dyna
 367    {
 368        SetResumeCallback(CallbackExpressionConverter.ToReflectionCall(callback));
 369        return this;
 370    }
 371
 372    [RequiresUnreferencedCode("The callback names its target service and method as strings, resolved by reflection when 
 373                              "subscriber loss; trimming may have removed them. Use the expression-based overload, which
 374                              "public methods automatically.")]
 375    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.OnLostSubscriberFailure(Refl
 376    {
 377        SetFailureCallback(callback);
 378        return this;
 379    }
 380
 381    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.OnLostSubscriberFailure<[Dyn
 382    {
 383        SetFailureCallback(CallbackExpressionConverter.ToReflectionCall(callback));
 384        return this;
 385    }
 386
 387    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.WithTimeout(TimeSpan timeout
 388    {
 389        WithTimeout(timeout);
 390        return this;
 391    }
 392
 393    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.WithReplyTarget()
 394    {
 395        WithReplyTarget();
 396        return this;
 397    }
 398
 399    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.WithReplyTarget(string name)
 400    {
 401        WithReplyTarget(name);
 402        return this;
 403    }
 404
 405    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.WithReplyTarget(AsyncRespons
 406    {
 407        WithReplyTarget(replyTarget);
 408        return this;
 409    }
 410
 411    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.Until(Func<T, bool> predicat
 412    {
 413        Until(predicate);
 414        return this;
 415    }
 416
 417    IRecoverableAsyncResponseAttachedBuilder<T> IRecoverableAsyncResponseAttachedBuilder<T>.Until(Func<T, Task<bool>> pr
 418    {
 419        Until(predicate);
 420        return this;
 421    }
 422
 423    [RequiresUnreferencedCode("The callback names its target service and method as strings, resolved by reflection when 
 424                              "subscriber loss; trimming may have removed them. Use the expression-based overload, which
 425                              "public methods automatically.")]
 426    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.OnLostSubscriberResume(Ref
 427    {
 428        SetResumeCallback(callback);
 429        return this;
 430    }
 431
 432    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.OnLostSubscriberResume<[Dy
 433    {
 434        SetResumeCallback(CallbackExpressionConverter.ToReflectionCall(callback));
 435        return this;
 436    }
 437
 438    [RequiresUnreferencedCode("The callback names its target service and method as strings, resolved by reflection when 
 439                              "subscriber loss; trimming may have removed them. Use the expression-based overload, which
 440                              "public methods automatically.")]
 441    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.OnLostSubscriberFailure(Re
 442    {
 443        SetFailureCallback(callback);
 444        return this;
 445    }
 446
 447    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.OnLostSubscriberFailure<[D
 448    {
 449        SetFailureCallback(CallbackExpressionConverter.ToReflectionCall(callback));
 450        return this;
 451    }
 452
 453    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.WithTimeout(TimeSpan timeo
 454    {
 455        WithTimeout(timeout);
 456        return this;
 457    }
 458
 459    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.WithReplyTarget()
 460    {
 461        WithReplyTarget();
 462        return this;
 463    }
 464
 465    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.WithReplyTarget(string nam
 466    {
 467        WithReplyTarget(name);
 468        return this;
 469    }
 470
 471    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.WithReplyTarget(AsyncRespo
 472    {
 473        WithReplyTarget(replyTarget);
 474        return this;
 475    }
 476
 477    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.Until(Func<T, bool> predic
 478    {
 479        Until(predicate);
 480        return this;
 481    }
 482
 483    IRecoverableAsyncResponseTriggeredBuilder<T> IRecoverableAsyncResponseTriggeredBuilder<T>.Until(Func<T, Task<bool>> 
 484    {
 485        Until(predicate);
 486        return this;
 487    }
 488}