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

Information
Class: AsyncResponse.WorkerJobExecutor
Assembly: AsyncResponse.Core
File(s): /home/runner/work/AsyncResponse/AsyncResponse/src/AsyncResponse.Core/WorkerJobExecutor.cs
Line coverage
100%
Covered lines: 33
Uncovered lines: 0
Coverable lines: 33
Total lines: 71
Line coverage: 100%
Branch coverage
100%
Covered branches: 4
Total branches: 4
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
ExecuteAsync()100%44100%

File(s)

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

#LineLine coverage
 1using Microsoft.Extensions.DependencyInjection;
 2using Microsoft.Extensions.Logging;
 3using System.Diagnostics;
 4
 5namespace AsyncResponse;
 6
 7/// <summary>
 8/// Executes <see cref="WorkerJobEnvelope"/>s: restores the correlation context and invokes the
 9/// described service method through the DI container. Shared by the broker ingress
 10/// (<see cref="IAsyncResponseIngress.HandleWorkerMessageAsync"/>) and the in-process worker
 11/// transport, so every transport executes jobs identically.
 12/// </summary>
 313internal sealed class WorkerJobExecutor(IServiceScopeFactory _scopeFactory, ILogger<WorkerJobExecutor> _logger)
 14{
 15    /// <summary>
 16    /// Executes the job. Exceptions propagate to the caller — transports decide whether to log,
 17    /// retry, or dead-letter.
 18    /// </summary>
 19    public async Task ExecuteAsync(WorkerJobEnvelope job)
 20    {
 321        ArgumentNullException.ThrowIfNull(job);
 22
 23        // Reject a job stamped with an unsupported schema rather than invoke a possibly-incompatible
 24        // method shape. Throwing routes the job through the transport's normal
 25        // failure/dead-letter handling. This is the single choke point every transport shares.
 326        if (!WorkerJobEnvelopeSchema.IsReadable(job.SchemaVersion))
 27        {
 228            _logger.LogWarning(
 229                "Worker job for correlationId {CorrelationId} has unsupported schema version {SchemaVersion} (current: {
 230                job.CorrelationId, job.SchemaVersion, WorkerJobEnvelopeSchema.Current);
 231            AsyncResponseDiagnostics.RecordWorkerOutcome("rejected");
 232            throw new InvalidOperationException(
 233                $"Worker job schema version {job.SchemaVersion} is not supported by this build " +
 234                $"(current: {WorkerJobEnvelopeSchema.Current}) and cannot be executed safely.");
 35        }
 36
 337        using var activity = AsyncResponseDiagnostics.StartActivity(
 338            "asyncresponse.worker.execute",
 339            ActivityKind.Consumer,
 340            job.CorrelationId);
 341        AsyncResponseDiagnostics.SetReplyTarget(activity, job.ReplyTarget);
 342        AsyncResponseDiagnostics.SetWorker(activity, job.Call);
 43
 344        _logger.LogDebug("Executing worker job {Target}.{Method} (correlationId: {CorrelationId}, replyTarget: {ReplyTar
 45
 46        try
 47        {
 48            // Scope the restored ambient context so one job cannot inherit or leak another job's
 49            // correlation id or reply target.
 350            using var asyncResponseScope = AsyncResponseContext.PushContext(job.CorrelationId, job.ReplyTarget);
 51
 352            var invocation = ReflectionExtensions.ResolveCallback(
 353                job.Call,
 354                payload: null,
 355                exception: null,
 356                correlationId: job.CorrelationId);
 57
 358            await using var scope = _scopeFactory.CreateAsyncScope();
 359            await scope.ServiceProvider.InvokeAsync(invocation).ConfigureAwait(false);
 60
 361            _logger.LogDebug("Executed worker job {Target}.{Method} successfully.", job.Call.ServiceInterfaceFullName, j
 362            AsyncResponseDiagnostics.RecordWorkerOutcome("executed");
 363        }
 264        catch (Exception ex)
 65        {
 266            AsyncResponseDiagnostics.SetError(activity, ex);
 267            AsyncResponseDiagnostics.RecordWorkerOutcome("failed");
 268            throw;
 69        }
 370    }
 71}