diff --git a/src/core/Elsa.Core/WorkflowProviders/DatabaseWorkflowProvider.cs b/src/core/Elsa.Core/WorkflowProviders/DatabaseWorkflowProvider.cs index 8f040786d..1dea1fe3d 100644 --- a/src/core/Elsa.Core/WorkflowProviders/DatabaseWorkflowProvider.cs +++ b/src/core/Elsa.Core/WorkflowProviders/DatabaseWorkflowProvider.cs @@ -1,5 +1,7 @@ +using System; using System.Collections.Generic; using System.Linq; +using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; using Elsa.Models; @@ -7,6 +9,7 @@ using Elsa.Persistence; using Elsa.Persistence.Specifications; using Elsa.Services; using Elsa.Services.Models; +using Microsoft.Extensions.Logging; namespace Elsa.WorkflowProviders { @@ -17,17 +20,40 @@ namespace Elsa.WorkflowProviders { private readonly IWorkflowDefinitionStore _workflowDefinitionStore; private readonly IWorkflowBlueprintMaterializer _workflowBlueprintMaterializer; + private readonly ILogger _logger; - public DatabaseWorkflowProvider(IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowBlueprintMaterializer workflowBlueprintMaterializer) + public DatabaseWorkflowProvider(IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowBlueprintMaterializer workflowBlueprintMaterializer, ILogger logger) { _workflowDefinitionStore = workflowDefinitionStore; _workflowBlueprintMaterializer = workflowBlueprintMaterializer; + _logger = logger; } - protected override async ValueTask> OnGetWorkflowsAsync(CancellationToken cancellationToken) + public override async IAsyncEnumerable GetWorkflowsAsync([EnumeratorCancellation] CancellationToken cancellationToken) { var workflowDefinitions = await _workflowDefinitionStore.FindManyAsync(Specification.Identity, cancellationToken: cancellationToken); - return await Task.WhenAll(workflowDefinitions.Select(async x => await _workflowBlueprintMaterializer.CreateWorkflowBlueprintAsync(x, cancellationToken))); + + foreach (var workflowDefinition in workflowDefinitions) + { + var workflowBlueprint = await TryMaterializeBlueprintAsync(workflowDefinition, cancellationToken); + + if (workflowBlueprint != null) + yield return workflowBlueprint; + } + } + + private async Task TryMaterializeBlueprintAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken) + { + try + { + return await _workflowBlueprintMaterializer.CreateWorkflowBlueprintAsync(workflowDefinition, cancellationToken); + } + catch (Exception e) + { + _logger.LogWarning(e, "Failed to materialize workflow definition {WorkflowDefinitionId} with version {WorkflowDefinitionVersion}", workflowDefinition.DefinitionId, workflowDefinition.Version); + } + + return null; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/WorkflowProviders/StorageWorkflowProvider.cs b/src/core/Elsa.Core/WorkflowProviders/StorageWorkflowProvider.cs index 19c7453da..5fb3b9a2a 100644 --- a/src/core/Elsa.Core/WorkflowProviders/StorageWorkflowProvider.cs +++ b/src/core/Elsa.Core/WorkflowProviders/StorageWorkflowProvider.cs @@ -1,11 +1,14 @@ -using System.Collections.Generic; +using System; +using System.Collections.Generic; using System.Runtime.CompilerServices; using System.Text; using System.Threading; +using System.Threading.Tasks; using Elsa.Models; using Elsa.Serialization; using Elsa.Services; using Elsa.Services.Models; +using Microsoft.Extensions.Logging; using Storage.Net.Blobs; namespace Elsa.WorkflowProviders @@ -15,12 +18,14 @@ namespace Elsa.WorkflowProviders private readonly IBlobStorage _storage; private readonly IWorkflowBlueprintMaterializer _workflowBlueprintMaterializer; private readonly IContentSerializer _contentSerializer; + private readonly ILogger _logger; - public StorageWorkflowProvider(IBlobStorage storage, IWorkflowBlueprintMaterializer workflowBlueprintMaterializer, IContentSerializer contentSerializer) + public StorageWorkflowProvider(IBlobStorage storage, IWorkflowBlueprintMaterializer workflowBlueprintMaterializer, IContentSerializer contentSerializer, ILogger logger) { _storage = storage; _workflowBlueprintMaterializer = workflowBlueprintMaterializer; _contentSerializer = contentSerializer; + _logger = logger; } public override async IAsyncEnumerable GetWorkflowsAsync([EnumeratorCancellation] CancellationToken cancellationToken) @@ -31,9 +36,25 @@ namespace Elsa.WorkflowProviders { var json = await _storage.ReadTextAsync(blob.FullPath, Encoding.UTF8, cancellationToken); var model = _contentSerializer.Deserialize(json); - var blueprint = await _workflowBlueprintMaterializer.CreateWorkflowBlueprintAsync(model, cancellationToken); - yield return blueprint; + var blueprint = await TryMaterializeBlueprintAsync(model, cancellationToken); + + if(blueprint != null) + yield return blueprint; } } + + private async Task TryMaterializeBlueprintAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken) + { + try + { + return await _workflowBlueprintMaterializer.CreateWorkflowBlueprintAsync(workflowDefinition, cancellationToken); + } + catch (Exception e) + { + _logger.LogWarning(e, "Failed to materialize workflow definition {WorkflowDefinitionId} with version {WorkflowDefinitionVersion}", workflowDefinition.DefinitionId, workflowDefinition.Version); + } + + return null; + } } } \ No newline at end of file