From 134896258aa118b592acbe6ff45caa88cda5fa3f Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 5 Oct 2021 17:25:41 +0200 Subject: [PATCH] Retry temporal startup job upon failure This prevents workflows from remaining "stuck" in case this job fails due to e.g. a DB exception. --- .../Elsa.Activities.Temporal.Common.csproj | 10 +++++-- .../HostedServices/StartJobs.cs | 29 ++++++++++++++----- 2 files changed, 29 insertions(+), 10 deletions(-) diff --git a/src/activities/Elsa.Activities.Temporal.Common/Elsa.Activities.Temporal.Common.csproj b/src/activities/Elsa.Activities.Temporal.Common/Elsa.Activities.Temporal.Common.csproj index cd72e7429..80610381b 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/Elsa.Activities.Temporal.Common.csproj +++ b/src/activities/Elsa.Activities.Temporal.Common/Elsa.Activities.Temporal.Common.csproj @@ -1,7 +1,7 @@ - - + + netstandard2.1 @@ -22,7 +22,11 @@ - + + + + + diff --git a/src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs b/src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs index ba70b770e..9bf57f0cd 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/HostedServices/StartJobs.cs @@ -1,4 +1,5 @@ -using System.Threading; +using System; +using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Temporal.Common.Bookmarks; using Elsa.Activities.Temporal.Common.Services; @@ -7,6 +8,8 @@ using Elsa.Services; using Elsa.Services.Bookmarks; using Microsoft.Extensions.Logging; using Open.Linq.AsyncExtensions; +using Polly; +using Polly.Retry; namespace Elsa.Activities.Temporal.Common.HostedServices { @@ -22,6 +25,7 @@ namespace Elsa.Activities.Temporal.Common.HostedServices private readonly IWorkflowInstanceScheduler _workflowScheduler; private readonly IDistributedLockProvider _distributedLockProvider; private readonly ILogger _logger; + private readonly AsyncRetryPolicy _retryPolicy; public StartJobs(IBookmarkFinder bookmarkFinder, IWorkflowInstanceScheduler workflowScheduler, IDistributedLockProvider distributedLockProvider, ILogger logger) { @@ -29,6 +33,12 @@ namespace Elsa.Activities.Temporal.Common.HostedServices _workflowScheduler = workflowScheduler; _distributedLockProvider = distributedLockProvider; _logger = logger; + + _retryPolicy = Policy + .Handle() + .WaitAndRetryForeverAsync(retryAttempt => + TimeSpan.FromSeconds(5) + ); } public async Task ExecuteAsync(CancellationToken stoppingToken) @@ -38,9 +48,14 @@ namespace Elsa.Activities.Temporal.Common.HostedServices if (handle == null) return; - await ScheduleTimerEventWorkflowsAsync(stoppingToken); - await ScheduleCronEventWorkflowsAsync(stoppingToken); - await ScheduleStartAtWorkflowsAsync(stoppingToken); + await _retryPolicy.ExecuteAsync(async () => await ExecuteInternalAsync(stoppingToken)); + } + + private async Task ExecuteInternalAsync(CancellationToken cancellationToken) + { + await ScheduleTimerEventWorkflowsAsync(cancellationToken); + await ScheduleCronEventWorkflowsAsync(cancellationToken); + await ScheduleStartAtWorkflowsAsync(cancellationToken); } private async Task ScheduleStartAtWorkflowsAsync(CancellationToken cancellationToken) @@ -53,7 +68,7 @@ namespace Elsa.Activities.Temporal.Common.HostedServices foreach (var result in bookmarkResults) { - var bookmark = (StartAtBookmark) result.Bookmark; + var bookmark = (StartAtBookmark)result.Bookmark; await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, bookmark.ExecuteAt, null, cancellationToken); index++; @@ -71,7 +86,7 @@ namespace Elsa.Activities.Temporal.Common.HostedServices foreach (var result in bookmarkResults) { - var bookmark = (TimerBookmark) result.Bookmark; + var bookmark = (TimerBookmark)result.Bookmark; await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, bookmark.ExecuteAt, null, cancellationToken); index++; @@ -89,7 +104,7 @@ namespace Elsa.Activities.Temporal.Common.HostedServices foreach (var result in bookmarkResults) { - var trigger = (CronBookmark) result.Bookmark; + var trigger = (CronBookmark)result.Bookmark; await _workflowScheduler.ScheduleAsync(result.WorkflowInstanceId!, result.ActivityId, trigger.ExecuteAt!.Value, null, cancellationToken); index++;