From 6ba50020fee919af340ca7c48520a3656e7c8f6c Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 24 Feb 2025 22:30:56 +0100 Subject: [PATCH] Add HTTP resiliency support using Polly and pipeline builder Introduce configurable resiliency mechanisms for HTTP requests, including retries, circuit breakers, and timeouts, leveraging Microsoft.Extensions.Resilience and Polly. Refactor `SendHttpRequestBase` to include an `EnableResiliency` input and encapsulate resiliency logic in a dedicated pipeline. Update project references to include necessary dependencies. --- .../Activities/SendHttpRequestBase.cs | 83 +++++++++++++++++-- src/modules/Elsa.Http/Elsa.Http.csproj | 20 +++-- 2 files changed, 86 insertions(+), 17 deletions(-) diff --git a/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs b/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs index 04085dbd8..3da8792b4 100644 --- a/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs +++ b/src/modules/Elsa.Http/Activities/SendHttpRequestBase.cs @@ -1,3 +1,4 @@ +using System.Net; using System.Net.Http.Headers; using Elsa.Extensions; using Elsa.Http.ContentWriters; @@ -7,6 +8,7 @@ using Elsa.Workflows.Attributes; using Elsa.Workflows.UIHints; using Elsa.Workflows.Models; using Microsoft.Extensions.Logging; +using Polly; using HttpHeaders = Elsa.Http.Models.HttpHeaders; namespace Elsa.Http; @@ -25,8 +27,7 @@ public abstract class SendHttpRequestBase : Activity /// /// The URL to send the request to. /// - [Input] - public Input Url { get; set; } = default!; + [Input] public Input Url { get; set; } = default!; /// /// The HTTP method to use when sending the request. @@ -81,6 +82,11 @@ public abstract class SendHttpRequestBase : Activity )] public Input RequestHeaders { get; set; } = new(new HttpHeaders()); + /// + /// Indicates whether resiliency mechanisms should be enabled for the HTTP request. + /// + public Input EnableResiliency { get; set; } = default!; + /// /// The HTTP response status code /// @@ -122,15 +128,16 @@ public abstract class SendHttpRequestBase : Activity private async Task TrySendAsync(ActivityExecutionContext context) { - var request = PrepareRequest(context); + var logger = (ILogger)context.GetRequiredService(typeof(ILogger<>).MakeGenericType(GetType())); var httpClientFactory = context.GetRequiredService(); var httpClient = httpClientFactory.CreateClient(nameof(SendHttpRequestBase)); var cancellationToken = context.CancellationToken; + var resiliencyEnabled = EnableResiliency.GetOrDefault(context, () => false); try { - var response = await httpClient.SendAsync(request, cancellationToken); + var response = await SendRequestAsync(); var parsedContent = await ParseContentAsync(context, response); var statusCode = (int)response.StatusCode; var responseHeaders = new HttpHeaders(response.Headers); @@ -147,7 +154,7 @@ public abstract class SendHttpRequestBase : Activity logger.LogWarning(e, "An error occurred while sending an HTTP request"); context.AddExecutionLogEntry("Error", e.Message, payload: new { - StackTrace = e.StackTrace + e.StackTrace }); context.JournalData.Add("Error", e.Message); await HandleRequestExceptionAsync(context, e); @@ -157,11 +164,30 @@ public abstract class SendHttpRequestBase : Activity logger.LogWarning(e, "An error occurred while sending an HTTP request"); context.AddExecutionLogEntry("Error", e.Message, payload: new { - StackTrace = e.StackTrace + e.StackTrace }); context.JournalData.Add("Cancelled", true); await HandleTaskCanceledExceptionAsync(context, e); } + + return; + + async Task SendRequestAsync() + { + if (resiliencyEnabled) + { + var pipeline = BuildResiliencyPipeline(context); + return await pipeline.ExecuteAsync(async ct => await SendRequestAsyncCore(ct), cancellationToken); + } + + return await SendRequestAsyncCore(); + } + + async Task SendRequestAsyncCore(CancellationToken ct = default) + { + var request = PrepareRequest(context); + return await httpClient.SendAsync(request, ct); + } } private async Task ParseContentAsync(ActivityExecutionContext context, HttpResponseMessage httpResponse) @@ -195,7 +221,7 @@ public abstract class SendHttpRequestBase : Activity { var method = Method.GetOrDefault(context) ?? "GET"; var url = Url.Get(context); - var request = new HttpRequestMessage(new HttpMethod(method), url); + var request = new HttpRequestMessage(new(method), url); var headers = context.GetHeaders(RequestHeaders); var authorization = Authorization.GetOrDefault(context); var addAuthorizationWithoutValidation = DisableAuthorizationHeaderValidation.GetOrDefault(context); @@ -218,7 +244,7 @@ public abstract class SendHttpRequestBase : Activity var factory = SelectContentWriter(contentType, factories); request.Content = factory.CreateHttpContent(content, contentType); } - + return request; } @@ -230,4 +256,45 @@ public abstract class SendHttpRequestBase : Activity var parsedContentType = new System.Net.Mime.ContentType(contentType); return factories.FirstOrDefault(httpContentFactory => httpContentFactory.SupportedContentTypes.Any(c => c == parsedContentType.MediaType)) ?? new JsonContentFactory(); } + + private ResiliencePipeline BuildResiliencyPipeline(ActivityExecutionContext context) + { + var pipelineBuilder = new ResiliencePipelineBuilder() + .AddRetry(new() + { + ShouldHandle = new PredicateBuilder() + .Handle() // Specific timeout exception + .Handle(ex => IsTransientStatusCode(ex.StatusCode)) // Network errors or transient HTTP codes + .HandleResult(response => IsTransientStatusCode(response.StatusCode)), + MaxRetryAttempts = 3, + Delay = TimeSpan.FromSeconds(Math.Min(Random.Shared.NextDouble() * 2, 8)), // Jittered delay capped at 8 secs + BackoffType = DelayBackoffType.Exponential + }) + .AddCircuitBreaker(new() + { + FailureRatio = 0.5, + SamplingDuration = TimeSpan.FromSeconds(30), + MinimumThroughput = 10, + BreakDuration = TimeSpan.FromSeconds(60) + }) + .AddTimeout(TimeSpan.FromSeconds(60)); // Outer timeout + + return pipelineBuilder.Build(); + } + + // Helper method to identify transient status codes. + private static bool IsTransientStatusCode(HttpStatusCode? statusCode) + { + if (!statusCode.HasValue) return true; // No status code (e.g., network failure) is worth retrying + return statusCode switch + { + HttpStatusCode.RequestTimeout => true, // 408 + HttpStatusCode.TooManyRequests => true, // 429 + HttpStatusCode.InternalServerError => true, // 500 + HttpStatusCode.BadGateway => true, // 502 + HttpStatusCode.ServiceUnavailable => true, // 503 + HttpStatusCode.GatewayTimeout => true, // 504 + _ => false + }; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Http/Elsa.Http.csproj b/src/modules/Elsa.Http/Elsa.Http.csproj index 26f35e9b5..caada764c 100644 --- a/src/modules/Elsa.Http/Elsa.Http.csproj +++ b/src/modules/Elsa.Http/Elsa.Http.csproj @@ -8,20 +8,22 @@ - + + + - + - + - - - - - - + + + + + +