From f3a8ff5e479b91cd0034a18a7a137c18f23f997d Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Mon, 8 Jul 2024 16:13:58 +0200 Subject: [PATCH 01/18] Add WorkflowExecutionLogRecordExtractor Introduced a new service, WorkflowExecutionLogRecordExtractor, to abstract the logic for extracting workflow execution logs records from the WorkflowExecutionContext. Updated StoreWorkflowExecutionLogSink to use this new service, which simplifies the execution log persistence method. This modification enhances code readability and enables potential reuse of the extraction logic. --- .../IWorkflowExecutionLogRecordExtractor.cs | 10 ++++++ .../Services/StoreWorkflowExecutionLogSink.cs | 27 ++------------ .../Services/WorkflowExecutionContext.cs | 35 +++++++++++++++++++ 3 files changed, 47 insertions(+), 25 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs new file mode 100644 index 000000000..0e704bd55 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowExecutionLogRecordExtractor.cs @@ -0,0 +1,10 @@ +using Elsa.Workflows.Runtime.Entities; + +namespace Elsa.Workflows.Runtime; + +/// Extracts workflow execution log records. +public interface IWorkflowExecutionLogRecordExtractor +{ + /// Extracts workflow execution logs from a workflow execution context. + IEnumerable ExtractWorkflowExecutionLogs(WorkflowExecutionContext context); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs index ecff347d0..d32bb89db 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs @@ -1,5 +1,4 @@ using Elsa.Mediator.Contracts; -using Elsa.Workflows.Contracts; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Notifications; @@ -9,34 +8,12 @@ namespace Elsa.Workflows.Runtime.Services; /// /// This implementation saves directly through the store. /// -public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IIdentityGenerator identityGenerator, INotificationSender notificationSender) : IWorkflowExecutionLogSink +public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IWorkflowExecutionLogRecordExtractor extractor, INotificationSender notificationSender) : IWorkflowExecutionLogSink { /// public async Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken) { - var records = context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord - { - Id = identityGenerator.GenerateId(), - ActivityInstanceId = x.ActivityInstanceId, - ParentActivityInstanceId = x.ParentActivityInstanceId, - ActivityNodeId = x.NodeId, - ActivityId = x.ActivityId, - ActivityType = x.ActivityType, - ActivityTypeVersion = x.ActivityTypeVersion, - ActivityName = x.ActivityName, - Message = x.Message, - EventName = x.EventName, - WorkflowDefinitionId = context.Workflow.Identity.DefinitionId, - WorkflowDefinitionVersionId = context.Workflow.Identity.Id, - WorkflowInstanceId = context.Id, - WorkflowVersion = context.Workflow.Version, - Source = x.Source, - ActivityState = x.ActivityState, - Payload = x.Payload, - Timestamp = x.Timestamp, - Sequence = x.Sequence - }).ToList(); - + var records = extractor.ExtractWorkflowExecutionLogs(context).ToList(); await store.AddManyAsync(records, context.CancellationTokens.SystemCancellationToken); await notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationTokens.SystemCancellationToken); } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs new file mode 100644 index 000000000..7badfd3ae --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs @@ -0,0 +1,35 @@ +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Runtime.Entities; + +namespace Elsa.Workflows.Runtime.Mappers; + +/// +public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : IWorkflowExecutionLogRecordExtractor +{ + /// + public IEnumerable ExtractWorkflowExecutionLogs(WorkflowExecutionContext context) + { + return context.ExecutionLog.Select(x => new WorkflowExecutionLogRecord + { + Id = identityGenerator.GenerateId(), + ActivityInstanceId = x.ActivityInstanceId, + ParentActivityInstanceId = x.ParentActivityInstanceId, + ActivityNodeId = x.NodeId, + ActivityId = x.ActivityId, + ActivityType = x.ActivityType, + ActivityTypeVersion = x.ActivityTypeVersion, + ActivityName = x.ActivityName, + Message = x.Message, + EventName = x.EventName, + WorkflowDefinitionId = context.Workflow.Identity.DefinitionId, + WorkflowDefinitionVersionId = context.Workflow.Identity.Id, + WorkflowInstanceId = context.Id, + WorkflowVersion = context.Workflow.Version, + Source = x.Source, + ActivityState = x.ActivityState, + Payload = x.Payload, + Timestamp = x.Timestamp, + Sequence = x.Sequence + }); + } +} \ No newline at end of file From 73b00fa6074c945c20ff1ca59f800266e6493979 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Tue, 9 Jul 2024 08:32:35 +0200 Subject: [PATCH 02/18] Added missing service registration --- .../Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs | 1 + .../Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index f27221b89..314979966 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -231,6 +231,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() // Stores. .AddScoped(BookmarkStore) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs index 7badfd3ae..577559bdb 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs @@ -1,7 +1,7 @@ using Elsa.Workflows.Contracts; using Elsa.Workflows.Runtime.Entities; -namespace Elsa.Workflows.Runtime.Mappers; +namespace Elsa.Workflows.Runtime.Services; /// public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : IWorkflowExecutionLogRecordExtractor From 99c079743831821632f88baf359b7d3e3e51f5fc Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Wed, 10 Jul 2024 11:33:41 +0200 Subject: [PATCH 03/18] Add endpoint for checking for journal updates --- .../Elsa.Server.Web/Enums/ApplicationRole.cs | 3 +- .../Contracts/IWorkflowInstancesApi.cs | 10 +++++ .../Models/HasJournalUpdateRequest.cs | 11 +++++ .../Journal/HasUpdates/Endpoint.cs | 40 +++++++++++++++++++ .../Journal/HasUpdates/Models.cs | 13 ++++++ 5 files changed, 76 insertions(+), 1 deletion(-) create mode 100644 src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs diff --git a/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs index 6eac67b27..5dd9900f1 100644 --- a/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs +++ b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs @@ -4,5 +4,6 @@ public enum ApplicationRole { Default, Api, - Worker + Worker, + Monitoring } diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs index 61ce1daaa..4a50a368c 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs @@ -46,6 +46,16 @@ public interface IWorkflowInstancesApi [Post("/workflow-instances/{workflowInstanceId}/journal")] Task> GetFilteredJournalAsync(string workflowInstanceId, GetFilteredJournalRequest? filter, int? skip = default, int? take = default, CancellationToken cancellationToken = default); + /// + /// Checks if there are updates in the journal for a specific workflow instance. + /// + /// The ID of the workflow instance for which to check for updates. + /// The request containing the ID and time from which to check for updates. + /// The cancellation token. + /// Returns whether updates are available for the journal. + [Get("/workflow-instances/{workflowInstanceId}/journal/has-updates")] + Task HasJournalUpdates(string workflowInstanceId, [Query]HasJournalUpdateRequest request, CancellationToken cancellationToken = default); + /// /// Deletes a workflow instance. /// diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs new file mode 100644 index 000000000..3e5d188a5 --- /dev/null +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs @@ -0,0 +1,11 @@ +namespace Elsa.Api.Client.Resources.WorkflowInstances.Models; + +/// A request to update a journal for a workflow instance. +public class HasJournalUpdateRequest +{ + /// The unique identifier of a workflow instance. + public string WorkflowInstanceId { get; set; } = default!; + + /// The start date for checking for updates in the workflow instance. + public DateTime UpdatesSince { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs new file mode 100644 index 000000000..4bf03ae89 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs @@ -0,0 +1,40 @@ +using Elsa.Abstractions; +using Elsa.Common.Entities; +using Elsa.Common.Models; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.OrderDefinitions; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates; + +/// Endpoint that checks if there are updates for a workflow instance. +[PublicAPI] +internal class HasUpdates : ElsaEndpoint +{ + private readonly IWorkflowExecutionLogStore _store; + + /// + public HasUpdates(IWorkflowExecutionLogStore store) + { + _store = store; + } + + /// + public override void Configure() + { + Get("/workflow-instances/{id}/journal/has-updates"); + ConfigurePermissions("read:workflow-instances"); + } + + /// + public override async Task ExecuteAsync(Request request, CancellationToken cancellationToken) + { + var pageArgs = PageArgs.From(1, 1, 0, 1); + var filter = new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = request.WorkflowInstanceId }; + var order = new WorkflowExecutionLogRecordOrder(x => x.Sequence, OrderDirection.Descending); + var pageOfRecords = await _store.FindManyAsync(filter, pageArgs, order, cancellationToken); + + return pageOfRecords.Items.Any(item => item.Timestamp >= request.UpdatesSince); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs new file mode 100644 index 000000000..cda8ed719 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs @@ -0,0 +1,13 @@ +using FastEndpoints; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates; + +/// The request to check if there are updates for a workflow instance. +public class Request +{ + /// The unique identifier of a workflow instance. + [BindFrom("id")] public string WorkflowInstanceId { get; set; } = default!; + + /// The start date for checking for updates in the workflow instance. + public DateTime UpdatesSince { get; set; } +} \ No newline at end of file From bd3797dea5c10a21bd5beee65c3ce0730e900676 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Thu, 11 Jul 2024 14:54:03 +0200 Subject: [PATCH 04/18] Fix datetime conversion when sending datetime to API endpoints --- .../Extensions/WebApplicationExtensions.cs | 4 ++++ .../WorkflowInstances/Journal/HasUpdates/Endpoint.cs | 12 ++---------- 2 files changed, 6 insertions(+), 10 deletions(-) diff --git a/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs b/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs index aaf4722ce..56bfbf28c 100644 --- a/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs +++ b/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs @@ -1,3 +1,4 @@ +using System.Globalization; using System.Text.Json; using System.Text.Json.Serialization; using Elsa.Workflows.Contracts; @@ -29,6 +30,9 @@ public static class WebApplicationExtensions config.Endpoints.RoutePrefix = routePrefix; config.Serializer.RequestDeserializer = DeserializeRequestAsync; config.Serializer.ResponseSerializer = SerializeRequestAsync; + + config.Binding.ValueParserFor(s => + new(DateTime.TryParse(s.ToString(),CultureInfo.InvariantCulture,DateTimeStyles.RoundtripKind, out var result), result)); }); /// diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs index 4bf03ae89..a535e25c9 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs @@ -10,16 +10,8 @@ namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates; /// Endpoint that checks if there are updates for a workflow instance. [PublicAPI] -internal class HasUpdates : ElsaEndpoint +internal class HasUpdates(IWorkflowExecutionLogStore store) : ElsaEndpoint { - private readonly IWorkflowExecutionLogStore _store; - - /// - public HasUpdates(IWorkflowExecutionLogStore store) - { - _store = store; - } - /// public override void Configure() { @@ -33,7 +25,7 @@ internal class HasUpdates : ElsaEndpoint var pageArgs = PageArgs.From(1, 1, 0, 1); var filter = new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = request.WorkflowInstanceId }; var order = new WorkflowExecutionLogRecordOrder(x => x.Sequence, OrderDirection.Descending); - var pageOfRecords = await _store.FindManyAsync(filter, pageArgs, order, cancellationToken); + var pageOfRecords = await store.FindManyAsync(filter, pageArgs, order, cancellationToken); return pageOfRecords.Items.Any(item => item.Timestamp >= request.UpdatesSince); } From 1c068373fda87c53f978edfdf929ca73f05cee1d Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 11 Jul 2024 22:09:13 +0200 Subject: [PATCH 05/18] Refactor API client DI methods for cleaner configuration Refactor `AddElsaApiKeyClient`, `AddElsaClient`, and related methods to improve readability and maintainability. Introduced `AddDefaultApiClientsUsingApiKey`, `AddDefaultApiClients`, and `AddApiClients` for a more structured approach, consolidating redundant code and ensuring better configurability. --- .../DependencyInjectionExtensions.cs | 125 +++++++++++------- 1 file changed, 75 insertions(+), 50 deletions(-) diff --git a/src/clients/Elsa.Api.Client/Extensions/DependencyInjectionExtensions.cs b/src/clients/Elsa.Api.Client/Extensions/DependencyInjectionExtensions.cs index 2ce4bb030..1dec341fd 100644 --- a/src/clients/Elsa.Api.Client/Extensions/DependencyInjectionExtensions.cs +++ b/src/clients/Elsa.Api.Client/Extensions/DependencyInjectionExtensions.cs @@ -20,21 +20,17 @@ using static Elsa.Api.Client.RefitSettingsHelper; namespace Elsa.Api.Client.Extensions; -/// /// Provides extension methods for dependency injection. -/// [PublicAPI] public static class DependencyInjectionExtensions { - /// - /// Adds the Elsa API client configured to use an API key to the service collection. - /// - public static IServiceCollection AddElsaApiKeyClient(this IServiceCollection services, Action configureOptions) + /// Adds default Elsa API clients configured to use an API key. + public static IServiceCollection AddDefaultApiClientsUsingApiKey(this IServiceCollection services, Action configureOptions) { var options = new ElsaClientOptions(); configureOptions(options); - return services.AddElsaClient(client => + return services.AddDefaultApiClients(client => { client.BaseAddress = options.BaseAddress; client.ApiKey = options.ApiKey; @@ -42,49 +38,70 @@ public static class DependencyInjectionExtensions }); } - /// - /// Adds the Elsa client to the service collection. - /// - public static IServiceCollection AddElsaClient(this IServiceCollection services, Action configureClient) + /// Adds default Elsa API clients. + public static IServiceCollection AddDefaultApiClients(this IServiceCollection services, Action configureClient) { - var builderOptions = new ElsaClientBuilderOptions(); - configureClient.Invoke(builderOptions); - builderOptions.ConfigureHttpClientBuilder += builder => builder.AddHttpMessageHandler(sp => (DelegatingHandler)sp.GetRequiredService(builderOptions.AuthenticationHandler)); - - services.AddScoped(builderOptions.AuthenticationHandler); - - services.Configure(options => + return services.AddApiClients(configureClient, builderOptions => { - options.BaseAddress = builderOptions.BaseAddress; - options.ConfigureHttpClient = builderOptions.ConfigureHttpClient; - options.ApiKey = builderOptions.ApiKey; + var builderOptionsWithoutRetryPolicy = new ElsaClientBuilderOptions + { + ApiKey = builderOptions.ApiKey, + AuthenticationHandler = builderOptions.AuthenticationHandler, + BaseAddress = builderOptions.BaseAddress, + ConfigureHttpClient = builderOptions.ConfigureHttpClient, + ConfigureHttpClientBuilder = builderOptions.ConfigureHttpClientBuilder, + ConfigureRetryPolicy = null + }; + + services.AddApi(builderOptions); + services.AddApi(builderOptionsWithoutRetryPolicy); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); + services.AddApi(builderOptions); }); + } - var builderOptionsWithoutRetryPolicy = new ElsaClientBuilderOptions + /// + /// Adds an API client to the service collection. Requires AddElsaClient to be called exactly once. + /// + public static IServiceCollection AddApiClient(this IServiceCollection services, Action configureClient) where T : class + { + return services.AddApiClients(configureClient, builderOptions => services.AddApi(builderOptions)); + } + + /// Adds the Elsa client to the service collection. + public static IServiceCollection AddApiClients(this IServiceCollection services, Action configureClient, Action? configureServices) + { + var builderOptionsServiceDescriptor = services.FirstOrDefault(x => x.ServiceType == typeof(ElsaClientBuilderOptions)); + + if (builderOptionsServiceDescriptor == null) { - ApiKey = builderOptions.ApiKey, - AuthenticationHandler = builderOptions.AuthenticationHandler, - BaseAddress = builderOptions.BaseAddress, - ConfigureHttpClient = builderOptions.ConfigureHttpClient, - ConfigureHttpClientBuilder = builderOptions.ConfigureHttpClientBuilder, - ConfigureRetryPolicy = null - }; + var builderOptions = new ElsaClientBuilderOptions(); + configureClient.Invoke(builderOptions); + builderOptions.ConfigureHttpClientBuilder += builder => builder.AddHttpMessageHandler(sp => (DelegatingHandler)sp.GetRequiredService(builderOptions.AuthenticationHandler)); + + services.AddScoped(builderOptions.AuthenticationHandler); + + services.Configure(options => + { + options.BaseAddress = builderOptions.BaseAddress; + options.ConfigureHttpClient = builderOptions.ConfigureHttpClient; + options.ApiKey = builderOptions.ApiKey; + }); + + configureServices?.Invoke(builderOptions); + } - services.AddApi(builderOptions); - services.AddApi(builderOptionsWithoutRetryPolicy); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); - services.AddApi(builderOptions); return services; } @@ -94,11 +111,23 @@ public static class DependencyInjectionExtensions /// The service collection. /// An options object that can be used to configure the HTTP client builder. /// The type representing the API. - public static void AddApi(this IServiceCollection services, ElsaClientBuilderOptions? httpClientBuilderOptions = default) where T : class + public static IServiceCollection AddApi(this IServiceCollection services, ElsaClientBuilderOptions? httpClientBuilderOptions = default) where T : class { - var builder = services.AddRefitClient(_ => CreateRefitSettings(), typeof(T).Name).ConfigureHttpClient(ConfigureElsaApiHttpClient); + return services.AddApi(typeof(T), httpClientBuilderOptions); + } + + /// + /// Adds a refit client for the specified API type. + /// + /// The service collection. + /// The type representing the API + /// An options object that can be used to configure the HTTP client builder. + public static IServiceCollection AddApi(this IServiceCollection services, Type apiType, ElsaClientBuilderOptions? httpClientBuilderOptions = default) + { + var builder = services.AddRefitClient(apiType, _ => CreateRefitSettings(), apiType.Name).ConfigureHttpClient(ConfigureElsaApiHttpClient); httpClientBuilderOptions?.ConfigureHttpClientBuilder(builder); httpClientBuilderOptions?.ConfigureRetryPolicy?.Invoke(builder); + return services; } /// @@ -115,9 +144,7 @@ public static class DependencyInjectionExtensions httpClientBuilderOptions?.ConfigureHttpClientBuilder(builder); } - /// /// Creates an API client for the specified API type. - /// public static T CreateApi(this IServiceProvider serviceProvider, Uri baseAddress) where T : class { var httpClientFactory = serviceProvider.GetRequiredService(); @@ -126,9 +153,7 @@ public static class DependencyInjectionExtensions return CreateApi(serviceProvider, httpClient); } - /// /// Creates an API client for the specified API type. - /// public static T CreateApi(this IServiceProvider serviceProvider, HttpClient httpClient) where T : class { return RestService.For(httpClient, CreateRefitSettings()); From e52d716eaa602ef26aaee53140d83406abcdcf7c Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Fri, 12 Jul 2024 08:08:29 +0200 Subject: [PATCH 06/18] Used DateTimeOffset instead of DateTime --- src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs | 2 +- .../WorkflowInstances/Models/HasJournalUpdateRequest.cs | 4 ++-- .../Elsa.Api.Common/Extensions/WebApplicationExtensions.cs | 4 ++-- .../WorkflowInstances/Journal/HasUpdates/Models.cs | 6 +++--- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs index 5dd9900f1..a4e8c7cf0 100644 --- a/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs +++ b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs @@ -5,5 +5,5 @@ public enum ApplicationRole Default, Api, Worker, - Monitoring + Monitor } diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs index 3e5d188a5..6efb34477 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs @@ -6,6 +6,6 @@ public class HasJournalUpdateRequest /// The unique identifier of a workflow instance. public string WorkflowInstanceId { get; set; } = default!; - /// The start date for checking for updates in the workflow instance. - public DateTime UpdatesSince { get; set; } + /// The start date for checking for updates in the workflow instance journal. + public DateTimeOffset UpdatesSince { get; set; } } \ No newline at end of file diff --git a/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs b/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs index 56bfbf28c..7d6f498a9 100644 --- a/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs +++ b/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs @@ -31,8 +31,8 @@ public static class WebApplicationExtensions config.Serializer.RequestDeserializer = DeserializeRequestAsync; config.Serializer.ResponseSerializer = SerializeRequestAsync; - config.Binding.ValueParserFor(s => - new(DateTime.TryParse(s.ToString(),CultureInfo.InvariantCulture,DateTimeStyles.RoundtripKind, out var result), result)); + config.Binding.ValueParserFor(s => + new(DateTimeOffset.TryParse(s.ToString(),CultureInfo.InvariantCulture,DateTimeStyles.RoundtripKind, out var result), result)); }); /// diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs index cda8ed719..eefafc6d0 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs @@ -2,12 +2,12 @@ using FastEndpoints; namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates; -/// The request to check if there are updates for a workflow instance. +/// The request to check if there are updates for a workflow instance journal. public class Request { /// The unique identifier of a workflow instance. [BindFrom("id")] public string WorkflowInstanceId { get; set; } = default!; - /// The start date for checking for updates in the workflow instance. - public DateTime UpdatesSince { get; set; } + /// The start date for checking for updates in the workflow instance journal. + public DateTimeOffset UpdatesSince { get; set; } } \ No newline at end of file From feb7da4d28f222f8db25df02974faff1d73a2ce0 Mon Sep 17 00:00:00 2001 From: MariusVuscanNx <96233009+MariusVuscanNx@users.noreply.github.com> Date: Fri, 12 Jul 2024 09:35:46 +0300 Subject: [PATCH 07/18] Implemented workflows reload endpoint (#5732) * Implemented workflows reload endpoint * Added componenent tests for reload * Adjusted cache handling * Made the refresh endpoint return refreshed and not found definitions * Update test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs Fix typo Co-authored-by: raymonddenhaan <155616759+raymonddenhaan@users.noreply.github.com> --------- Co-authored-by: Sipke Schoorstra Co-authored-by: raymonddenhaan <155616759+raymonddenhaan@users.noreply.github.com> --- .../Handlers/InvalidateHttpWorkflowsCache.cs | 24 ++-- .../WorkflowDefinitionEventsConsumer.cs | 17 ++- ...dWorkflowDefinitionNotificationsHandler.cs | 18 ++- .../Messages/WorkflowDefinitionsReloaded.cs | 8 ++ .../Services/AmbientConsumerScope.cs | 2 +- .../WorkflowDefinitions/Refresh/Endpoint.cs | 9 +- .../WorkflowDefinitions/Refresh/Models.cs | 10 +- .../WorkflowDefinitions/Reload/Endpoint.cs | 28 +++++ .../IWorkflowDefinitionStorePopulator.cs | 6 +- .../Contracts/IWorkflowDefinitionsReloader.cs | 8 ++ .../Features/WorkflowRuntimeFeature.cs | 1 + .../Handlers/InvalidateTriggersCache.cs | 10 +- .../Handlers/InvalidateWorkflowsCache.cs | 26 +++++ .../WorkflowDefinitionsReloaded.cs | 6 + .../RefreshWorkflowDefinitionsResponse.cs | 2 +- ...DefaultWorkflowDefinitionStorePopulator.cs | 13 ++- .../Services/WorkflowDefinitionRefresher.cs | 7 +- .../Services/WorkflowDefinitionsReloader.cs | 17 +++ .../Elsa.Workflows.ComponentTests.csproj | 6 + .../DynamicEndpointTests.cs | 20 ++-- .../RemoveReloadWorkflowTests.cs | 35 ++++++ .../http-workflow.json | 104 ++++++++++++++++++ 22 files changed, 336 insertions(+), 41 deletions(-) create mode 100644 src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Reload/Endpoint.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json diff --git a/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs index 714ee1a36..bb87afe99 100644 --- a/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs +++ b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs @@ -15,17 +15,18 @@ namespace Elsa.Http.Handlers; /// A handler that invalidates the HTTP workflows cache when a workflow definition is published, retracted, or deleted or when triggers are indexed. /// [UsedImplicitly] -public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflowsCacheManager, - ITriggerStore triggerStore, - IHttpWorkflowsCacheManager cacheManager) : +public class InvalidateHttpWorkflowsCache( + IHttpWorkflowsCacheManager httpWorkflowsCacheManager, + ITriggerStore triggerStore) : INotificationHandler, INotificationHandler, INotificationHandler, - INotificationHandler, + INotificationHandler, INotificationHandler, INotificationHandler, INotificationHandler, - INotificationHandler + INotificationHandler, + INotificationHandler { /// public Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) @@ -91,6 +92,15 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo await InvalidateCacheAsync(notification.IndexedWorkflowTriggers.Workflow.Identity.DefinitionId); } + /// + public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + foreach (var workflowDefinitionId in notification.WorkflowDefinitionIds) + { + await InvalidateCacheAsync(workflowDefinitionId); + } + } + private async Task InvalidateCacheAsync(string workflowDefinitionId) { await httpWorkflowsCacheManager.EvictWorkflowAsync(workflowDefinitionId); @@ -103,7 +113,7 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo WorkflowDefinitionVersionId = workflowDefinitionVersionId }; var triggers = await triggerStore.FindManyAsync(filter, cancellationToken); - + await InvalidateTriggerCacheAsync(triggers, cancellationToken); } @@ -113,7 +123,7 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo { if (trigger?.Payload is HttpEndpointBookmarkPayload httpPayload) { - var hash = cacheManager.ComputeBookmarkHash(httpPayload.Path, httpPayload.Method); + var hash = httpWorkflowsCacheManager.ComputeBookmarkHash(httpPayload.Path, httpPayload.Method); await httpWorkflowsCacheManager.EvictTriggerAsync(hash, cancellationToken); } } diff --git a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs index 6fed5a9a5..d43d11172 100644 --- a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs +++ b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs @@ -18,7 +18,8 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr global::MassTransit.IConsumer, global::MassTransit.IConsumer, global::MassTransit.IConsumer, - global::MassTransit.IConsumer + global::MassTransit.IConsumer, + global::MassTransit.IConsumer { /// public Task Consume(ConsumeContext context) @@ -82,9 +83,19 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr { var message = context.Message; var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsRefreshed(message.WorkflowDefinitionIds); - AmbientConsumerScope.IsConsumerExecutionContext = true; + AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = true; await notificationSender.SendAsync(notification, context.CancellationToken); - AmbientConsumerScope.IsConsumerExecutionContext = false; + AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = false; + } + + /// + public async Task Consume(ConsumeContext context) + { + var message = context.Message; + var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsReloaded(message.WorkflowDefinitionIds); + AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = true; + await notificationSender.SendAsync(notification, context.CancellationToken); + AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = false; } private Task UpdateDefinition(string id, bool usableAsActivity) diff --git a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs index 12fdef392..991eb0914 100644 --- a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs +++ b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs @@ -1,4 +1,3 @@ -using Elsa.MassTransit.Contracts; using Elsa.MassTransit.Services; using Elsa.Mediator.Contracts; using Elsa.Workflows.Management.Notifications; @@ -18,7 +17,8 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IBus bus) : INotificationHandler, INotificationHandler, INotificationHandler, - INotificationHandler + INotificationHandler, + INotificationHandler { /// public Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) @@ -73,11 +73,23 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IBus bus) : public Task HandleAsync(WorkflowDefinitionsRefreshed notification, CancellationToken cancellationToken) { // Prevent re-entrance. - if (AmbientConsumerScope.IsConsumerExecutionContext) + if (AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer) return Task.CompletedTask; var definitionIds = notification.WorkflowDefinitionIds; var message = new Distributed.WorkflowDefinitionsRefreshed(definitionIds); return bus.Publish(message, cancellationToken); } + + /// + public Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + // Prevent re-entrance. + if (AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer) + return Task.CompletedTask; + + var definitionIds = notification.WorkflowDefinitionIds; + var message = new Distributed.WorkflowDefinitionsReloaded(definitionIds); + return bus.Publish(message, cancellationToken); + } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs new file mode 100644 index 000000000..bba459390 --- /dev/null +++ b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs @@ -0,0 +1,8 @@ +namespace Elsa.MassTransit.Messages; + +/// Represents a message that indicates that the specified workflow definitions have been reloaded. +public class WorkflowDefinitionsReloaded(ICollection workflowDefinitionIds) +{ + /// The workflow definition IDs that have been reloaded. + public ICollection WorkflowDefinitionIds { get; set; } = workflowDefinitionIds; +} \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Services/AmbientConsumerScope.cs b/src/modules/Elsa.MassTransit/Services/AmbientConsumerScope.cs index 2236cd969..55c585edc 100644 --- a/src/modules/Elsa.MassTransit/Services/AmbientConsumerScope.cs +++ b/src/modules/Elsa.MassTransit/Services/AmbientConsumerScope.cs @@ -6,7 +6,7 @@ public static class AmbientConsumerScope private static readonly AsyncLocal IsRaisedFromConsumerState = new(); /// Gets or sets a value that indicates if the current code is initiated from a consumer. - public static bool IsConsumerExecutionContext + public static bool IsWorkflowDefinitionEventsConsumer { get => IsRaisedFromConsumerState.Value; set => IsRaisedFromConsumerState.Value = value; diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs index 752102397..8acc4c5a9 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs @@ -1,6 +1,7 @@ using Elsa.Abstractions; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Responses; using JetBrains.Annotations; namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Refresh; @@ -18,13 +19,13 @@ internal class Refresh(IWorkflowDefinitionsRefresher workflowDefinitionsRefreshe public override async Task HandleAsync(Request request, CancellationToken cancellationToken) { - await RefreshWorkflowDefinitionsAsync(request.DefinitionIds, cancellationToken); - await SendOkAsync(cancellationToken); + var result = await RefreshWorkflowDefinitionsAsync(request.DefinitionIds, cancellationToken); + await SendOkAsync(new Response(result.Refreshed, result.NotFound), cancellationToken); } - private async Task RefreshWorkflowDefinitionsAsync(ICollection? definitionIds, CancellationToken cancellationToken) + private async Task RefreshWorkflowDefinitionsAsync(ICollection? definitionIds, CancellationToken cancellationToken) { var request = new RefreshWorkflowDefinitionsRequest(definitionIds, BatchSize); - await workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(request, cancellationToken); + return await workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(request, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Models.cs index c5970467b..f52d3e9e6 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Models.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Models.cs @@ -1,6 +1,14 @@ +using System.Text.Json.Serialization; + namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Refresh; -public class Request +internal class Request { public ICollection? DefinitionIds { get; set; } +} + +internal class Response(ICollection refreshed, ICollection notFound) +{ + [JsonPropertyName("refreshed")] public ICollection Refreshed { get; } = refreshed; + [JsonPropertyName("notFound")] public ICollection NotFound { get; } = notFound; } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Reload/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Reload/Endpoint.cs new file mode 100644 index 000000000..71d14f949 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Reload/Endpoint.cs @@ -0,0 +1,28 @@ +using Elsa.Abstractions; +using Elsa.Workflows.Runtime.Contracts; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Reload; + +[PublicAPI] +internal class Reload(IWorkflowDefinitionsReloader workflowDefinitionsReloader) : ElsaEndpointWithoutRequest +{ + private const int BatchSize = 10; + + public override void Configure() + { + Post("/actions/workflow-definitions/reload"); + ConfigurePermissions("actions:workflow-definitions:reload"); + } + + public override async Task HandleAsync(CancellationToken cancellationToken) + { + await ReloadWorkflowDefinitionsAsync(cancellationToken); + await SendOkAsync(cancellationToken); + } + + private async Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken) + { + await workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs index d9326097f..aeb5c198b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs @@ -12,14 +12,14 @@ public interface IWorkflowDefinitionStorePopulator /// Populates the with workflow definitions provided from implementations. /// /// The cancellation token. - Task PopulateStoreAsync(CancellationToken cancellationToken = default); - + Task> PopulateStoreAsync(CancellationToken cancellationToken = default); + /// /// Populates the with workflow definitions provided from implementations. /// /// Whether to index triggers. /// The cancellation token. - Task PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default); + Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default); /// /// Adds a workflow definition to the store. diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs new file mode 100644 index 000000000..e2cc2987e --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs @@ -0,0 +1,8 @@ +namespace Elsa.Workflows.Runtime.Contracts; + +/// Reloads all workflows by re-invoking the populator. +public interface IWorkflowDefinitionsReloader +{ + /// Reloads all workflows by re-invoking the populator. + Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 314979966..75a41922e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -221,6 +221,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateTriggersCache.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateTriggersCache.cs index e0347168a..d1581b49f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateTriggersCache.cs +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateTriggersCache.cs @@ -15,11 +15,19 @@ namespace Elsa.Workflows.Runtime.Handlers; /// The cache is invalidated by calling the TriggerTokenAsync method of the ICacheManager passed to the class constructor. /// [UsedImplicitly] -public class InvalidateTriggersCache(ICacheManager cacheManager) : INotificationHandler +public class InvalidateTriggersCache(ICacheManager cacheManager) : + INotificationHandler, + INotificationHandler { /// public Task HandleAsync(WorkflowDefinitionsRefreshed notification, CancellationToken cancellationToken) { return cacheManager.TriggerTokenAsync(CachingTriggerStore.CacheInvalidationTokenKey, cancellationToken).AsTask(); } + + /// + public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + await cacheManager.TriggerTokenAsync(CachingTriggerStore.CacheInvalidationTokenKey, cancellationToken); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs new file mode 100644 index 000000000..fb6dc6d45 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs @@ -0,0 +1,26 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Runtime.Notifications; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// +/// A notification handler that invalidates workflows cache when workflow definitions are reloaded. +/// +/// +/// The class implements the INotificationHandler interface and is responsible for handling WorkflowDefinitionsReloaded notifications. +/// When a WorkflowDefinitionsReloaded notification is received, the HandleAsync method is called to invalidate the http definition cache. +/// +[UsedImplicitly] +public class InvalidateWorkflowsCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) : INotificationHandler +{ + /// + public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + foreach (var workflowDefinitionId in notification.WorkflowDefinitionIds) + { + await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(workflowDefinitionId, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs new file mode 100644 index 000000000..882fa7a51 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs @@ -0,0 +1,6 @@ +using Elsa.Mediator.Contracts; + +namespace Elsa.Workflows.Runtime.Notifications; + +/// Published when workflow definitions have been reloaded. +public record WorkflowDefinitionsReloaded(ICollection WorkflowDefinitionIds) : INotification; diff --git a/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs b/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs index 36cfc5360..ce380a1ec 100644 --- a/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs +++ b/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs @@ -3,4 +3,4 @@ using Elsa.Workflows.Management.Entities; namespace Elsa.Workflows.Runtime.Responses; /// Represents a response to a request to refresh workflow definitions. -public record RefreshWorkflowDefinitionsResponse(ICollection WorkflowDefinitions); \ No newline at end of file +public record RefreshWorkflowDefinitionsResponse(ICollection Refreshed, ICollection NotFound); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs index aa0208a7b..cba487b5a 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -49,21 +49,26 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } /// - public Task PopulateStoreAsync(CancellationToken cancellationToken = default) + public Task> PopulateStoreAsync(CancellationToken cancellationToken = default) { return PopulateStoreAsync(true, cancellationToken); } /// - public async Task PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) + public async Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) { var providers = _workflowDefinitionProviders(); + var workflowDefinitionIds = new List(); foreach (var provider in providers) { var results = await provider.GetWorkflowsAsync(cancellationToken).AsTask().ToList(); + workflowDefinitionIds.AddRange(results.Select(w => w.Workflow.Id)); + foreach (var result in results) await AddAsync(result, indexTriggers, cancellationToken); } + + return workflowDefinitionIds; } /// @@ -131,8 +136,8 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP if (existingDefinitionVersion != null) { workflowDefinitionsToSave.Add(existingDefinitionVersion); - - if(existingDefinitionVersion.Id != workflow.Identity.Id) + + if (existingDefinitionVersion.Id != workflow.Identity.Id) { // It's possible that the imported workflow definition has a different ID than the existing one in the store. // In a future update, we might store this discrepancy in a "troubleshooting" table and provide tooling for managing these, and other, discrepancies. diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs index 39eb938e1..8c3f64fa1 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs @@ -39,12 +39,15 @@ public class WorkflowDefinitionsRefresher(IWorkflowDefinitionStore store, ITrigg await IndexWorkflowTriggersAsync(definitions.Items, cancellationToken); processedWorkflowDefinitions.AddRange(definitions.Items); currentPage++; + + if (definitions.Items.Count < batchSize) + break; } - var processedWorkflowDefinitionIds = processedWorkflowDefinitions.Select(x => x.Id).ToList(); + var processedWorkflowDefinitionIds = processedWorkflowDefinitions.Select(x => x.DefinitionId).ToList(); var notification = new WorkflowDefinitionsRefreshed(processedWorkflowDefinitionIds); await notificationSender.SendAsync(notification, cancellationToken); - return new RefreshWorkflowDefinitionsResponse(processedWorkflowDefinitions); + return new RefreshWorkflowDefinitionsResponse(processedWorkflowDefinitionIds, request.DefinitionIds?.Except(processedWorkflowDefinitionIds)?.ToList() ?? []); } private async Task IndexWorkflowTriggersAsync(IEnumerable definitions, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs new file mode 100644 index 000000000..3afadafd2 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs @@ -0,0 +1,17 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Notifications; + +namespace Elsa.Workflows.Runtime.Services; + +/// +public class WorkflowDefinitionsReloader(IWorkflowDefinitionStorePopulator workflowDefinitionStorePopulator, INotificationSender notificationSender) : IWorkflowDefinitionsReloader +{ + /// + public async Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken) + { + var definitionIds = await workflowDefinitionStorePopulator.PopulateStoreAsync(true, cancellationToken); + var notification = new WorkflowDefinitionsReloaded(definitionIds); + await notificationSender.SendAsync(notification, cancellationToken); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index 7d745b355..c6588808f 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -21,6 +21,9 @@ + + Always + Always @@ -75,6 +78,9 @@ Always + + Always + diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs index f3b8e3df2..2be751964 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs @@ -1,4 +1,5 @@ -using Elsa.Workflows.ComponentTests.Helpers.Services; +using System.Net; +using Elsa.Workflows.ComponentTests.Helpers.Services; using Elsa.Workflows.Runtime.Contracts; using Microsoft.Extensions.DependencyInjection; @@ -15,22 +16,19 @@ public class DynamicEndpointTests : AppComponentTest } [Fact] - public async Task HelloWorldWorkflow_ShouldRespondWithHelloWorld() + public async Task ChangingEndpointValueThenRefresh_WorkflowShouldRespondToTheNewValue() { var client = WorkflowServer.CreateHttpWorkflowClient(); - var firstResponse = await client.GetStringAsync("first-value"); + var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "first-value")); StaticValueHolder.Value = "second-value"; - await _workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(new Runtime.Requests.RefreshWorkflowDefinitionsRequest - { - BatchSize = 10, - DefinitionIds = ["f69f061159adc3ae"] - }, CancellationToken.None); + var _ = await _workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync( + new Runtime.Requests.RefreshWorkflowDefinitionsRequest() { DefinitionIds = ["f69f061159adc3ae"] }, CancellationToken.None); - var secondResponse = await client.GetStringAsync("second-value"); + var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "second-value")); - Assert.Equal("", firstResponse); - Assert.Equal("", secondResponse); + Assert.Equal(HttpStatusCode.OK, firstResponse.StatusCode); + Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); } } \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs new file mode 100644 index 000000000..390c58d14 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs @@ -0,0 +1,35 @@ +using System.Net; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowDefinitionReload; + +public class RemoveReloadWorkflowTests : AppComponentTest +{ + private readonly IWorkflowDefinitionManager _workflowDefinitionManager; + private readonly IWorkflowDefinitionsReloader _workflowDefinitionsReloader; + + public RemoveReloadWorkflowTests(App app) : base(app) + { + _workflowDefinitionManager = Scope.ServiceProvider.GetRequiredService(); + _workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService(); + } + + [Fact] + public async Task RemovingTheWorkflowThenReload_WorkflowShouldBeReachableAgain() + { + var client = WorkflowServer.CreateHttpWorkflowClient(); + + var result = await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); + + var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); + + var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + + Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); + Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json new file mode 100644 index 000000000..c885f29b6 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json @@ -0,0 +1,104 @@ +{ + "id": "c2e19917-d03e-41ed-9d19-ef5c1ddc9871", + "definitionId": "f68b09bc-2013-4617-b82f-d76b6819a624", + "name": "Workflow 1", + "createdAt": "2024-07-04T13:41:33.1905134+00:00", + "version": 1, + "toolVersion": "3.3.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": { + "Elsa:WorkflowContextProviderTypes": [] + }, + "isReadonly": false, + "isSystem": false, + "isLatest": true, + "isPublished": true, + "options": { + "autoUpdateConsumingWorkflows": false + }, + "root": { + "type": "Elsa.Flowchart", + "version": 1, + "id": "f8019bae-c506-4467-812b-e403e7aad8ff", + "nodeId": "Workflow1:f8019bae-c506-4467-812b-e403e7aad8ff", + "metadata": {}, + "customProperties": { + "source": "FlowchartJsonConverter.cs:45", + "notFoundConnections": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "activities": [ + { + "path": { + "typeName": "String", + "expression": { + "type": "Literal", + "value": "reload-test" + } + }, + "supportedMethods": { + "typeName": "String[]", + "expression": { + "type": "Literal", + "value": "[\u0022GET\u0022]" + } + }, + "authorize": { + "typeName": "Boolean", + "expression": { + "type": "Literal", + "value": false + } + }, + "policy": { + "typeName": "String", + "expression": { + "type": "Literal" + } + }, + "requestTimeout": null, + "requestSizeLimit": null, + "fileSizeLimit": null, + "allowedFileExtensions": null, + "blockedFileExtensions": null, + "allowedMimeTypes": null, + "exposeRequestTooLargeOutcome": false, + "exposeFileTooLargeOutcome": false, + "exposeInvalidFileExtensionOutcome": false, + "exposeInvalidFileMimeTypeOutcome": false, + "parsedContent": null, + "files": null, + "routeData": null, + "queryStringData": null, + "headers": null, + "result": null, + "id": "cc196a5e-04f7-4794-b25c-1a0987632d8a", + "nodeId": "Workflow1:f8019bae-c506-4467-812b-e403e7aad8ff:cc196a5e-04f7-4794-b25c-1a0987632d8a", + "name": "HttpEndpoint1", + "type": "Elsa.HttpEndpoint", + "version": 1, + "customProperties": { + "canStartWorkflow": true, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -260, + "y": 460 + }, + "size": { + "width": 176.390625, + "height": 50 + } + } + } + } + ], + "connections": [] + } +} \ No newline at end of file From 1a0d4714bd43b51083c0880966f9cb5096628dc9 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 16 Jul 2024 07:58:27 +0200 Subject: [PATCH 08/18] Update activity ID reference in BroadcastWorkflowProgress (#5767) Replaced the use of ActivityId with ActivityNodeId when retrieving distinct activity IDs. This change ensures that we are using the correct node identifier for logging workflow progress. --- .../RealTime/Handlers/BroadcastWorkflowProgress.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/modules/Elsa.Workflows.Api/RealTime/Handlers/BroadcastWorkflowProgress.cs b/src/modules/Elsa.Workflows.Api/RealTime/Handlers/BroadcastWorkflowProgress.cs index f8a5a2cd5..7609710f7 100644 --- a/src/modules/Elsa.Workflows.Api/RealTime/Handlers/BroadcastWorkflowProgress.cs +++ b/src/modules/Elsa.Workflows.Api/RealTime/Handlers/BroadcastWorkflowProgress.cs @@ -35,7 +35,7 @@ public class BroadcastWorkflowProgress : public async Task HandleAsync(ActivityExecutionLogUpdated notification, CancellationToken cancellationToken) { var workflowInstanceId = notification.WorkflowExecutionContext.Id; - var activityIds = notification.Records.Select(x => x.ActivityId).Distinct().ToList(); + var activityIds = notification.Records.Select(x => x.ActivityNodeId).Distinct().ToList(); var stats = (await _activityExecutionStatsService.GetStatsAsync(workflowInstanceId, activityIds, cancellationToken)).ToList(); var clients = _hubContext.Clients.Group(workflowInstanceId); var message = new ActivityExecutionLogUpdatedMessage(stats); From f63f823e7b939e8d264f6d997abc828bee007d16 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 16 Jul 2024 09:02:24 +0200 Subject: [PATCH 09/18] Update ElsaStudioVersion to 3.2.0-rc3.443. This change updates the ElsaStudioVersion property to the latest release candidate. It ensures compatibility with recent updates and bug fixes provided in version 3.2.0-rc3.443. --- Directory.Build.props | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Directory.Build.props b/Directory.Build.props index 4c56d1528..7828ed57c 100644 --- a/Directory.Build.props +++ b/Directory.Build.props @@ -27,6 +27,6 @@ true - 3.2.0-rc3.417 + 3.2.0-rc3.443 \ No newline at end of file From 9fb51c6cd0f87e63d830a3238084c40c9656f1f3 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 16 Jul 2024 13:55:23 +0200 Subject: [PATCH 10/18] Add functions to get workflow definition details Introduced three new functions: `getWorkflowDefinitionId`, `getWorkflowDefinitionVersionId`, and `getWorkflowDefinitionVersion`. These functions provide easy access to specific workflow definition details within the JavaScript evaluator. Removed redundant constructor in CommonFunctionsDefinitionProvider to streamline the code. --- .../Services/JintJavaScriptEvaluator.cs | 3 +++ .../CommonFunctionsDefinitionProvider.cs | 25 +++++++++++-------- 2 files changed, 17 insertions(+), 11 deletions(-) diff --git a/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs b/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs index ac9cd3908..e718f7c4b 100644 --- a/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs +++ b/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs @@ -76,6 +76,9 @@ public class JintJavaScriptEvaluator(IConfiguration configuration, INotification configureEngine?.Invoke(engine); // Add common functions. + engine.SetValue("getWorkflowDefinitionId", (Func)(() => context.GetWorkflowExecutionContext().Workflow.Identity.DefinitionId)); + engine.SetValue("getWorkflowDefinitionVersionId", (Func)(() => context.GetWorkflowExecutionContext().Workflow.Identity.Id)); + engine.SetValue("getWorkflowDefinitionVersion", (Func)(() => context.GetWorkflowExecutionContext().Workflow.Identity.Version)); engine.SetValue("getWorkflowInstanceId", (Func)(() => context.GetActivityExecutionContext().WorkflowExecutionContext.Id)); engine.SetValue("setCorrelationId", (Action)(value => context.GetActivityExecutionContext().WorkflowExecutionContext.CorrelationId = value)); engine.SetValue("getCorrelationId", (Func)(() => context.GetActivityExecutionContext().WorkflowExecutionContext.CorrelationId)); diff --git a/src/modules/Elsa.JavaScript/TypeDefinitions/Providers/CommonFunctionsDefinitionProvider.cs b/src/modules/Elsa.JavaScript/TypeDefinitions/Providers/CommonFunctionsDefinitionProvider.cs index 28943dd8e..4a37429f0 100644 --- a/src/modules/Elsa.JavaScript/TypeDefinitions/Providers/CommonFunctionsDefinitionProvider.cs +++ b/src/modules/Elsa.JavaScript/TypeDefinitions/Providers/CommonFunctionsDefinitionProvider.cs @@ -6,20 +6,23 @@ using Humanizer; namespace Elsa.JavaScript.TypeDefinitions.Providers; -/// /// Produces s for common functions. -/// -internal class CommonFunctionsDefinitionProvider : FunctionDefinitionProvider +internal class CommonFunctionsDefinitionProvider(ITypeAliasRegistry typeAliasRegistry) : FunctionDefinitionProvider { - private readonly ITypeAliasRegistry _typeAliasRegistry; - - public CommonFunctionsDefinitionProvider(ITypeAliasRegistry typeAliasRegistry) - { - _typeAliasRegistry = typeAliasRegistry; - } - protected override IEnumerable GetFunctionDefinitions(TypeDefinitionContext context) { + yield return CreateFunctionDefinition(builder => builder + .Name("getWorkflowDefinitionId") + .ReturnType("string")); + + yield return CreateFunctionDefinition(builder => builder + .Name("getWorkflowDefinitionVersionId") + .ReturnType("string")); + + yield return CreateFunctionDefinition(builder => builder + .Name("getWorkflowDefinitionVersion") + .ReturnType("number")); + yield return CreateFunctionDefinition(builder => builder .Name("getWorkflowInstanceId") .ReturnType("string")); @@ -94,7 +97,7 @@ internal class CommonFunctionsDefinitionProvider : FunctionDefinitionProvider { var pascalName = variable.Name.Pascalize(); var variableType = variable.GetVariableType(); - var typeAlias = _typeAliasRegistry.TryGetAlias(variableType, out var alias) ? alias : "any"; + var typeAlias = typeAliasRegistry.TryGetAlias(variableType, out var alias) ? alias : "any"; // get{Variable}. yield return CreateFunctionDefinition(builder => builder.Name($"get{pascalName}").ReturnType(typeAlias)); From a69492d521efa2867d4f2ea7751683b9aaae3671 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Tue, 16 Jul 2024 14:31:20 +0200 Subject: [PATCH 11/18] Overwritten System.Text.Json version due to vulnerability --- Directory.Packages.props | 104 +++++++++--------- .../Elsa.Server.Web/Elsa.Server.Web.csproj | 4 +- .../Elsa.ServerAndStudio.Web.csproj | 5 + .../Elsa.Studio.Web/Elsa.Studio.Web.csproj | 5 + .../Elsa.Api.Client/Elsa.Api.Client.csproj | 5 + .../Elsa.EntityFrameworkCore.Common.csproj | 5 + .../Elsa.EntityFrameworkCore.MySql.csproj | 5 + ...Elsa.EntityFrameworkCore.PostgreSql.csproj | 10 +- .../Elsa.EntityFrameworkCore.SqlServer.csproj | 1 + .../Elsa.EntityFrameworkCore.Sqlite.csproj | 5 + .../Elsa.EntityFrameworkCore.csproj | 5 + .../Elsa.FileStorage/Elsa.FileStorage.csproj | 5 + src/modules/Elsa.Http/Elsa.Http.csproj | 5 + src/modules/Elsa.MongoDb/Elsa.MongoDb.csproj | 5 + ...sa.Quartz.EntityFrameworkCore.MySql.csproj | 5 + ...artz.EntityFrameworkCore.PostgreSql.csproj | 4 + ...uartz.EntityFrameworkCore.SqlServer.csproj | 1 + ...a.Quartz.EntityFrameworkCore.Sqlite.csproj | 5 + src/modules/Elsa.Telnyx/Elsa.Telnyx.csproj | 5 + .../Elsa.WorkflowProviders.BlobStorage.csproj | 4 + .../Elsa.Workflows.Api.csproj | 6 +- 21 files changed, 138 insertions(+), 61 deletions(-) diff --git a/Directory.Packages.props b/Directory.Packages.props index 0663c9b3e..9dc0dd505 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -6,11 +6,11 @@ - + - - + + @@ -18,51 +18,51 @@ - + - + - + - - - + + + - + - + - + - - + + - - - + + + - - - + + + - + - - + + @@ -76,38 +76,39 @@ - - - - - + + + + + - + - - - - - + + + + + - + - - + + + - + @@ -134,21 +135,20 @@ - - - - - - - - - - - + + + + + + + + + + @@ -156,15 +156,15 @@ - + - + - + \ No newline at end of file diff --git a/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj b/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj index e1da5c473..15636b9b6 100644 --- a/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj +++ b/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj @@ -53,9 +53,9 @@ - + - + diff --git a/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj b/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj index 3cb1c9444..23d188639 100644 --- a/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj +++ b/src/bundles/Elsa.ServerAndStudio.Web/Elsa.ServerAndStudio.Web.csproj @@ -40,4 +40,9 @@ + + + + + \ No newline at end of file diff --git a/src/bundles/Elsa.Studio.Web/Elsa.Studio.Web.csproj b/src/bundles/Elsa.Studio.Web/Elsa.Studio.Web.csproj index 4190ff8f0..432ff7c6c 100644 --- a/src/bundles/Elsa.Studio.Web/Elsa.Studio.Web.csproj +++ b/src/bundles/Elsa.Studio.Web/Elsa.Studio.Web.csproj @@ -15,6 +15,11 @@ + + + + + diff --git a/src/clients/Elsa.Api.Client/Elsa.Api.Client.csproj b/src/clients/Elsa.Api.Client/Elsa.Api.Client.csproj index f8246f150..7ce828980 100644 --- a/src/clients/Elsa.Api.Client/Elsa.Api.Client.csproj +++ b/src/clients/Elsa.Api.Client/Elsa.Api.Client.csproj @@ -16,5 +16,10 @@ + + + + + diff --git a/src/modules/Elsa.EntityFrameworkCore.Common/Elsa.EntityFrameworkCore.Common.csproj b/src/modules/Elsa.EntityFrameworkCore.Common/Elsa.EntityFrameworkCore.Common.csproj index f7beacf88..eab7ce6b2 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Common/Elsa.EntityFrameworkCore.Common.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.Common/Elsa.EntityFrameworkCore.Common.csproj @@ -17,6 +17,11 @@ + + + + + diff --git a/src/modules/Elsa.EntityFrameworkCore.MySql/Elsa.EntityFrameworkCore.MySql.csproj b/src/modules/Elsa.EntityFrameworkCore.MySql/Elsa.EntityFrameworkCore.MySql.csproj index c00d3b8d7..5df8320c1 100644 --- a/src/modules/Elsa.EntityFrameworkCore.MySql/Elsa.EntityFrameworkCore.MySql.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.MySql/Elsa.EntityFrameworkCore.MySql.csproj @@ -12,6 +12,11 @@ + + + + + diff --git a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Elsa.EntityFrameworkCore.PostgreSql.csproj b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Elsa.EntityFrameworkCore.PostgreSql.csproj index c34aad583..ce9f15c01 100644 --- a/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Elsa.EntityFrameworkCore.PostgreSql.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.PostgreSql/Elsa.EntityFrameworkCore.PostgreSql.csproj @@ -12,13 +12,11 @@ - - + + + - - - - + diff --git a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj index b4c979bb2..a3c210bda 100644 --- a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj @@ -20,5 +20,6 @@ + \ No newline at end of file diff --git a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Elsa.EntityFrameworkCore.Sqlite.csproj b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Elsa.EntityFrameworkCore.Sqlite.csproj index 7b7edfa04..6580755a3 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Sqlite/Elsa.EntityFrameworkCore.Sqlite.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.Sqlite/Elsa.EntityFrameworkCore.Sqlite.csproj @@ -13,6 +13,11 @@ + + + + + diff --git a/src/modules/Elsa.EntityFrameworkCore/Elsa.EntityFrameworkCore.csproj b/src/modules/Elsa.EntityFrameworkCore/Elsa.EntityFrameworkCore.csproj index 8a9b1901e..3dd2e4a6a 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Elsa.EntityFrameworkCore.csproj +++ b/src/modules/Elsa.EntityFrameworkCore/Elsa.EntityFrameworkCore.csproj @@ -11,6 +11,11 @@ + + + + + diff --git a/src/modules/Elsa.FileStorage/Elsa.FileStorage.csproj b/src/modules/Elsa.FileStorage/Elsa.FileStorage.csproj index f043d374f..0393bc645 100644 --- a/src/modules/Elsa.FileStorage/Elsa.FileStorage.csproj +++ b/src/modules/Elsa.FileStorage/Elsa.FileStorage.csproj @@ -15,6 +15,11 @@ + + + + + diff --git a/src/modules/Elsa.Http/Elsa.Http.csproj b/src/modules/Elsa.Http/Elsa.Http.csproj index bda010f88..1024756e2 100644 --- a/src/modules/Elsa.Http/Elsa.Http.csproj +++ b/src/modules/Elsa.Http/Elsa.Http.csproj @@ -10,6 +10,11 @@ + + + + + diff --git a/src/modules/Elsa.MongoDb/Elsa.MongoDb.csproj b/src/modules/Elsa.MongoDb/Elsa.MongoDb.csproj index 91fc3bdff..39d38e218 100644 --- a/src/modules/Elsa.MongoDb/Elsa.MongoDb.csproj +++ b/src/modules/Elsa.MongoDb/Elsa.MongoDb.csproj @@ -13,6 +13,11 @@ + + + + + diff --git a/src/modules/Elsa.Quartz.EntityFrameworkCore.MySql/Elsa.Quartz.EntityFrameworkCore.MySql.csproj b/src/modules/Elsa.Quartz.EntityFrameworkCore.MySql/Elsa.Quartz.EntityFrameworkCore.MySql.csproj index 20bb854d8..a2e95de32 100644 --- a/src/modules/Elsa.Quartz.EntityFrameworkCore.MySql/Elsa.Quartz.EntityFrameworkCore.MySql.csproj +++ b/src/modules/Elsa.Quartz.EntityFrameworkCore.MySql/Elsa.Quartz.EntityFrameworkCore.MySql.csproj @@ -13,6 +13,11 @@ + + + + + diff --git a/src/modules/Elsa.Quartz.EntityFrameworkCore.PostgreSql/Elsa.Quartz.EntityFrameworkCore.PostgreSql.csproj b/src/modules/Elsa.Quartz.EntityFrameworkCore.PostgreSql/Elsa.Quartz.EntityFrameworkCore.PostgreSql.csproj index 0290e7236..8cca95901 100644 --- a/src/modules/Elsa.Quartz.EntityFrameworkCore.PostgreSql/Elsa.Quartz.EntityFrameworkCore.PostgreSql.csproj +++ b/src/modules/Elsa.Quartz.EntityFrameworkCore.PostgreSql/Elsa.Quartz.EntityFrameworkCore.PostgreSql.csproj @@ -14,12 +14,16 @@ + + + + diff --git a/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj b/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj index f9efcc8a3..25f942e21 100644 --- a/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj +++ b/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj @@ -23,5 +23,6 @@ + \ No newline at end of file diff --git a/src/modules/Elsa.Quartz.EntityFrameworkCore.Sqlite/Elsa.Quartz.EntityFrameworkCore.Sqlite.csproj b/src/modules/Elsa.Quartz.EntityFrameworkCore.Sqlite/Elsa.Quartz.EntityFrameworkCore.Sqlite.csproj index 63abad7b2..283ff1ffb 100644 --- a/src/modules/Elsa.Quartz.EntityFrameworkCore.Sqlite/Elsa.Quartz.EntityFrameworkCore.Sqlite.csproj +++ b/src/modules/Elsa.Quartz.EntityFrameworkCore.Sqlite/Elsa.Quartz.EntityFrameworkCore.Sqlite.csproj @@ -13,6 +13,11 @@ + + + + + diff --git a/src/modules/Elsa.Telnyx/Elsa.Telnyx.csproj b/src/modules/Elsa.Telnyx/Elsa.Telnyx.csproj index 9633deb6c..ae23a5915 100644 --- a/src/modules/Elsa.Telnyx/Elsa.Telnyx.csproj +++ b/src/modules/Elsa.Telnyx/Elsa.Telnyx.csproj @@ -22,4 +22,9 @@ + + + + + diff --git a/src/modules/Elsa.WorkflowProviders.BlobStorage/Elsa.WorkflowProviders.BlobStorage.csproj b/src/modules/Elsa.WorkflowProviders.BlobStorage/Elsa.WorkflowProviders.BlobStorage.csproj index 3ac581466..c6c83beb6 100644 --- a/src/modules/Elsa.WorkflowProviders.BlobStorage/Elsa.WorkflowProviders.BlobStorage.csproj +++ b/src/modules/Elsa.WorkflowProviders.BlobStorage/Elsa.WorkflowProviders.BlobStorage.csproj @@ -15,4 +15,8 @@ + + + + diff --git a/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj b/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj index 4fff95dc9..95079450f 100644 --- a/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj +++ b/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj @@ -14,5 +14,9 @@ - + + + + + From a0846effaec4656d7eb32d629cbd279de894eb94 Mon Sep 17 00:00:00 2001 From: Mohamed Ali Date: Tue, 16 Jul 2024 15:54:20 +0300 Subject: [PATCH 12/18] bugfixing issues related to new PropertyBag in workflow definition (#5778) --- .../Elsa.Api.Client/Extensions/PropertyBagExtensions.cs | 2 +- .../Endpoints/WorkflowDefinitions/Post/Endpoint.cs | 3 ++- .../Services/WorkflowDefinitionImporter.cs | 3 ++- 3 files changed, 5 insertions(+), 3 deletions(-) diff --git a/src/clients/Elsa.Api.Client/Extensions/PropertyBagExtensions.cs b/src/clients/Elsa.Api.Client/Extensions/PropertyBagExtensions.cs index c27200c88..a72b7b556 100644 --- a/src/clients/Elsa.Api.Client/Extensions/PropertyBagExtensions.cs +++ b/src/clients/Elsa.Api.Client/Extensions/PropertyBagExtensions.cs @@ -22,7 +22,7 @@ public static class PropertyBagExtensions if (!propertyBag.TryGetValue(key, out var value)) return defaultValue(); - var json = (string)value; + var json = value.ToString(); return JsonSerializer.Deserialize(json); } diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs index df0e87a15..9c899e3a1 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Post/Endpoint.cs @@ -79,7 +79,8 @@ internal class Post( draft.Name = model.Name?.Trim(); draft.ToolVersion = model.ToolVersion; draft.Description = model.Description?.Trim(); - draft.CustomProperties = model.CustomProperties ?? new Dictionary(); + draft.CustomProperties = model.CustomProperties ?? new Dictionary(); + draft.PropertyBag = model.PropertyBag ?? new PropertyBag(); draft.Variables = variables; draft.Inputs = inputs; draft.Outputs = outputs; diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs index f09a4ef7b..47df45637 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs @@ -63,7 +63,8 @@ namespace Elsa.Workflows.Management.Services draft.MaterializerName = JsonWorkflowMaterializer.MaterializerName; draft.Name = model.Name?.Trim(); draft.Description = model.Description?.Trim(); - draft.CustomProperties = model.CustomProperties ?? new Dictionary(); + draft.CustomProperties = model.CustomProperties ?? new Dictionary(); + draft.PropertyBag = model.PropertyBag ?? new PropertyBag(); draft.Variables = variables; draft.Inputs = model.Inputs ?? new List(); draft.Outputs = model.Outputs ?? new List(); From 85c701e5d07ec52cde98432986197387f071886d Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Tue, 16 Jul 2024 15:04:46 +0200 Subject: [PATCH 13/18] Overwriting System.Formats.Asn1 due to vulnerability --- src/common/Elsa.DropIns/Elsa.DropIns.csproj | 5 +++++ src/modules/Elsa.Dapper/Elsa.Dapper.csproj | 1 + .../Elsa.EntityFrameworkCore.SqlServer.csproj | 1 + .../Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj | 1 + 4 files changed, 8 insertions(+) diff --git a/src/common/Elsa.DropIns/Elsa.DropIns.csproj b/src/common/Elsa.DropIns/Elsa.DropIns.csproj index b1cf52e98..d77a55517 100644 --- a/src/common/Elsa.DropIns/Elsa.DropIns.csproj +++ b/src/common/Elsa.DropIns/Elsa.DropIns.csproj @@ -15,6 +15,11 @@ + + + + + diff --git a/src/modules/Elsa.Dapper/Elsa.Dapper.csproj b/src/modules/Elsa.Dapper/Elsa.Dapper.csproj index c0b8c8e4b..dc9a4279f 100644 --- a/src/modules/Elsa.Dapper/Elsa.Dapper.csproj +++ b/src/modules/Elsa.Dapper/Elsa.Dapper.csproj @@ -33,6 +33,7 @@ + diff --git a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj index a3c210bda..2a69f65ae 100644 --- a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj @@ -20,6 +20,7 @@ + \ No newline at end of file diff --git a/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj b/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj index 25f942e21..f33bc9af8 100644 --- a/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj +++ b/src/modules/Elsa.Quartz.EntityFrameworkCore.SqlServer/Elsa.Quartz.EntityFrameworkCore.SqlServer.csproj @@ -23,6 +23,7 @@ + \ No newline at end of file From f448e9520a025613092ae1d73ea7e3af50fd28cd Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 17 Jul 2024 09:51:45 +0200 Subject: [PATCH 14/18] Fix reloading logic for workflow definitions (#5781) * Add reloading logic for workflow definitions Introduced a new mechanism to handle reloaded workflow definitions. This includes creating a "ReloadedWorkflowDefinition" model, updating the caching logic, and modifying notification handlers to work with the enhanced workflow reloading logic. This ensures workflow definitions are updated and managed correctly when published, retracted, or deleted. * Improve RefreshActivityRegistry documentation Updated the XML documentation to clarify that `RefreshActivityRegistry` refreshes the `IActivityRegistry` for `WorkflowDefinitionActivityProvider` whenever workflow definitions are reloaded, instead of when they are published, retracted, or deleted. * Add 'materializer' to user dictionary The term 'materializer' has been added to the user dictionary to improve code spelling and naming consistency. This change ensures that 'materializer' is recognized as a correct term in the codebase. * Organize test files by adding a 'Workflows' directory Renamed 'http-workflow.json' to indicate it belongs under 'Workflows'. This improves file organization and clarity within the 'WorkflowDefinitionReload' scenario. * Add TestWorkflowMaterializer and TestWorkflowProvider Introduce `TestWorkflowMaterializer` for deserializing workflows from `TestWorkflowProvider`. Added integration of these new components in the `ReloadWorkflowTests` and `WorkflowServer`. Also renamed `RemoveReloadWorkflowTests` to `ReloadWorkflowTests`. * Rename and expand workflow reload tests Renamed `RemoveReloadWorkflowTests` to `ReloadWorkflowTests` to better reflect its purpose and expanded with additional test cases. Added tests to verify workflow and activity registry updates after source provider changes and workflow reloads. * Update copy settings for workflow test scenarios Reorganized and added 'CopyToOutputDirectory' settings for JSON files in workflow test scenarios. Ensured all necessary files are correctly included and copied during output directory builds to maintain test consistency. --- Elsa.sln.DotSettings | 2 + .../Handlers/InvalidateHttpWorkflowsCache.cs | 10 +- .../WorkflowDefinitionEventsConsumer.cs | 2 +- ...dWorkflowDefinitionNotificationsHandler.cs | 4 +- .../Messages/WorkflowDefinitionsReloaded.cs | 10 +- ...rkflowDefinitionActivityRegistryUpdater.cs | 5 +- .../Handlers/DeleteWorkflowInstances.cs | 2 + .../IWorkflowDefinitionStorePopulator.cs | 13 +- .../Contracts/IWorkflowDefinitionsReloader.cs | 4 +- .../Contracts/IWorkflowRuntime.cs | 2 +- .../Features/CachingWorkflowRuntimeFeature.cs | 3 +- .../Features/WorkflowRuntimeFeature.cs | 3 +- .../Handlers/InvalidateWorkflowsCache.cs | 8 +- .../Handlers/RefreshActivityRegistry.cs | 31 +++++ .../Models/ReloadedWorkflowDefinition.cs | 30 +++++ .../WorkflowDefinitionsReloaded.cs | 3 +- ...DefaultWorkflowDefinitionStorePopulator.cs | 35 +++--- .../Services/WorkflowDefinitionsReloader.cs | 6 +- .../Elsa.Workflows.ComponentTests.csproj | 6 +- .../Helpers/Fixtures/WorkflowServer.cs | 17 ++- .../Materializers/TestWorkflowMaterializer.cs | 26 ++++ .../WorkflowProviders/TestWorkflowProvider.cs | 14 +++ .../ReloadWorkflowTests.cs | 113 ++++++++++++++++++ .../RemoveReloadWorkflowTests.cs | 35 ------ .../{ => Workflows}/http-workflow.json | 0 25 files changed, 290 insertions(+), 94 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs delete mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs rename test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/{ => Workflows}/http-workflow.json (100%) diff --git a/Elsa.sln.DotSettings b/Elsa.sln.DotSettings index 19414abf6..b291186ee 100644 --- a/Elsa.sln.DotSettings +++ b/Elsa.sln.DotSettings @@ -12,10 +12,12 @@ EF True True + True True True True True + True True True True diff --git a/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs index bb87afe99..f07a32072 100644 --- a/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs +++ b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs @@ -95,10 +95,8 @@ public class InvalidateHttpWorkflowsCache( /// public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) { - foreach (var workflowDefinitionId in notification.WorkflowDefinitionIds) - { - await InvalidateCacheAsync(workflowDefinitionId); - } + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await InvalidateCacheAsync(reloadedWorkflowDefinition.DefinitionId); } private async Task InvalidateCacheAsync(string workflowDefinitionId) @@ -119,9 +117,9 @@ public class InvalidateHttpWorkflowsCache( private async Task InvalidateTriggerCacheAsync(IEnumerable triggers, CancellationToken cancellationToken) { - foreach (StoredTrigger trigger in triggers) + foreach (var trigger in triggers) { - if (trigger?.Payload is HttpEndpointBookmarkPayload httpPayload) + if (trigger.Payload is HttpEndpointBookmarkPayload httpPayload) { var hash = httpWorkflowsCacheManager.ComputeBookmarkHash(httpPayload.Path, httpPayload.Method); await httpWorkflowsCacheManager.EvictTriggerAsync(hash, cancellationToken); diff --git a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs index d43d11172..4f7821a63 100644 --- a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs +++ b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs @@ -92,7 +92,7 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr public async Task Consume(ConsumeContext context) { var message = context.Message; - var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsReloaded(message.WorkflowDefinitionIds); + var notification = new Elsa.Workflows.Runtime.Notifications.WorkflowDefinitionsReloaded(message.ReloadedWorkflowDefinitions); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = true; await notificationSender.SendAsync(notification, context.CancellationToken); AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer = false; diff --git a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs index 991eb0914..10e654479 100644 --- a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs +++ b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs @@ -88,8 +88,8 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IBus bus) : if (AmbientConsumerScope.IsWorkflowDefinitionEventsConsumer) return Task.CompletedTask; - var definitionIds = notification.WorkflowDefinitionIds; - var message = new Distributed.WorkflowDefinitionsReloaded(definitionIds); + var reloadedWorkflowDefinitions = notification.ReloadedWorkflowDefinitions; + var message = new Distributed.WorkflowDefinitionsReloaded(reloadedWorkflowDefinitions); return bus.Publish(message, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs index bba459390..138492683 100644 --- a/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs +++ b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs @@ -1,8 +1,10 @@ -namespace Elsa.MassTransit.Messages; +using Elsa.Workflows.Runtime.Models; + +namespace Elsa.MassTransit.Messages; /// Represents a message that indicates that the specified workflow definitions have been reloaded. -public class WorkflowDefinitionsReloaded(ICollection workflowDefinitionIds) +public class WorkflowDefinitionsReloaded(ICollection reloadedWorkflowDefinitions) { - /// The workflow definition IDs that have been reloaded. - public ICollection WorkflowDefinitionIds { get; set; } = workflowDefinitionIds; + /// The reloaded workflow definitions. + public ICollection ReloadedWorkflowDefinitions { get; set; } = reloadedWorkflowDefinitions; } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs index 454ac3dbe..826b37c8c 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionActivityRegistryUpdater.cs @@ -8,9 +8,9 @@ public interface IWorkflowDefinitionActivityRegistryUpdater /// /// Tries to add a workflow as an activity to the registry. /// - /// The ID of the workflow definition. + /// The version ID of the workflow definition. /// The cancellation token. - Task AddToRegistry(string workflowDefinitionId, CancellationToken cancellationToken = default); + Task AddToRegistry(string workflowDefinitionVersionId, CancellationToken cancellationToken = default); /// /// Removes workflow definition activities from the . @@ -18,7 +18,6 @@ public interface IWorkflowDefinitionActivityRegistryUpdater /// The ID of the workflow definition to remove. void RemoveDefinitionFromRegistry(string workflowDefinitionId); - /// /// Removes a workflow definition version activity from the . /// diff --git a/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs b/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs index 1a0178476..66141f945 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs @@ -2,12 +2,14 @@ using Elsa.Mediator.Contracts; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Notifications; +using JetBrains.Annotations; namespace Elsa.Workflows.Management.Handlers; /// /// Deletes workflow instances when a workflow definition or version is deleted. /// +[UsedImplicitly] public class DeleteWorkflowInstances : INotificationHandler, INotificationHandler, diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs index aeb5c198b..5cdf956dc 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs @@ -1,4 +1,5 @@ using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Runtime.Models; namespace Elsa.Workflows.Runtime.Contracts; @@ -12,27 +13,27 @@ public interface IWorkflowDefinitionStorePopulator /// Populates the with workflow definitions provided from implementations. /// /// The cancellation token. - Task> PopulateStoreAsync(CancellationToken cancellationToken = default); + Task> PopulateStoreAsync(CancellationToken cancellationToken = default); /// /// Populates the with workflow definitions provided from implementations. /// /// Whether to index triggers. /// The cancellation token. - Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default); + Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default); /// /// Adds a workflow definition to the store. /// /// A materialized workflow. /// An optional cancellation token. - Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default); - + Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default); + /// /// Adds a workflow definition to the store. /// /// A materialized workflow. - /// /// Whether to index triggers. + /// Whether to index triggers. /// An optional cancellation token. - Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default); + Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs index e2cc2987e..f17d06edf 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs @@ -1,8 +1,8 @@ namespace Elsa.Workflows.Runtime.Contracts; -/// Reloads all workflows by re-invoking the populator. +/// Reloads all workflows by invoking the populator. public interface IWorkflowDefinitionsReloader { - /// Reloads all workflows by re-invoking the populator. + /// Reloads all workflows by invoking the populator. Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs index a60d31f88..be218983b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowRuntime.cs @@ -15,7 +15,7 @@ namespace Elsa.Workflows.Runtime.Contracts; public interface IWorkflowRuntime { /// - /// Returns a value whether or not the specified workflow definition can create a new instance. + /// Returns a value whether the specified workflow definition can create a new instance. /// Task CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = default); diff --git a/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs index 9c0e26e73..bb31de5d4 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/CachingWorkflowRuntimeFeature.cs @@ -26,6 +26,7 @@ public class CachingWorkflowRuntimeFeature : FeatureBase .Decorate() // Handlers. - .AddNotificationHandler(); + .AddNotificationHandler() + .AddNotificationHandler(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index 75a41922e..6bdbb7d99 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -7,7 +7,6 @@ using Elsa.Features.Attributes; using Elsa.Features.Services; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Management.Handlers; using Elsa.Workflows.Management.Services; using Elsa.Workflows.Runtime.ActivationValidators; using Elsa.Workflows.Runtime.Contracts; @@ -272,12 +271,12 @@ public class WorkflowRuntimeFeature : FeatureBase .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() - .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() // Workflow activation strategies. .AddScoped() diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs index fb6dc6d45..d2d52613c 100644 --- a/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs @@ -6,7 +6,7 @@ using JetBrains.Annotations; namespace Elsa.Workflows.Runtime.Handlers; /// -/// A notification handler that invalidates workflows cache when workflow definitions are reloaded. +/// A notification handler that invalidates the workflow cache when workflow definitions are reloaded. /// /// /// The class implements the INotificationHandler interface and is responsible for handling WorkflowDefinitionsReloaded notifications. @@ -18,9 +18,7 @@ public class InvalidateWorkflowsCache(IWorkflowDefinitionCacheManager workflowDe /// public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) { - foreach (var workflowDefinitionId in notification.WorkflowDefinitionIds) - { - await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(workflowDefinitionId, cancellationToken); - } + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(reloadedWorkflowDefinition.DefinitionId, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs new file mode 100644 index 000000000..607751270 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/RefreshActivityRegistry.cs @@ -0,0 +1,31 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Runtime.Notifications; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// Refreshes the for the provider whenever workflow definitions are reloaded. +[PublicAPI] +public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater) : INotificationHandler +{ + /// + public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await UpdateDefinition(reloadedWorkflowDefinition.DefinitionVersionId, reloadedWorkflowDefinition.UsableAsActivity); + } + + private Task UpdateDefinition(string definitionVersionId, bool? usableAsActivity) + { + // A workflow should remain in the activity registry unless no longer being marked as an activity. + if (usableAsActivity.GetValueOrDefault()) + return workflowDefinitionActivityRegistryUpdater.AddToRegistry(definitionVersionId); + + workflowDefinitionActivityRegistryUpdater.RemoveDefinitionVersionFromRegistry(definitionVersionId); + return Task.CompletedTask; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs b/src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs new file mode 100644 index 000000000..7a30f1f3b --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Models/ReloadedWorkflowDefinition.cs @@ -0,0 +1,30 @@ +using Elsa.Workflows.Management.Entities; + +namespace Elsa.Workflows.Runtime.Models; + +/// +/// Represents a reloaded workflow definition with necessary properties for +/// identification, versioning, and usability status as an activity. +/// +/// The unique identifier for the workflow definition. +/// The unique identifier for the specific version of the workflow definition. +/// The version number of the workflow definition. +/// Indicates whether the workflow definition can be used as an activity. +public record ReloadedWorkflowDefinition(string DefinitionId, string DefinitionVersionId, int Version, bool UsableAsActivity) +{ + /// + /// Creates an instance of from a given . + /// + /// The workflow definition used to create the reloaded workflow definition. + /// A new instance of . + public static ReloadedWorkflowDefinition FromDefinition(WorkflowDefinition workflowDefinition) + { + return new ReloadedWorkflowDefinition + ( + workflowDefinition.DefinitionId, + workflowDefinition.Id, + workflowDefinition.Version, + workflowDefinition.Options.UsableAsActivity ?? false + ); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs index 882fa7a51..336347ff4 100644 --- a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs @@ -1,6 +1,7 @@ using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Models; namespace Elsa.Workflows.Runtime.Notifications; /// Published when workflow definitions have been reloaded. -public record WorkflowDefinitionsReloaded(ICollection WorkflowDefinitionIds) : INotification; +public record WorkflowDefinitionsReloaded(ICollection ReloadedWorkflowDefinitions) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs index cba487b5a..fad0f8759 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -49,42 +49,47 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } /// - public Task> PopulateStoreAsync(CancellationToken cancellationToken = default) + public Task> PopulateStoreAsync(CancellationToken cancellationToken = default) { return PopulateStoreAsync(true, cancellationToken); } /// - public async Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) + public async Task> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default) { var providers = _workflowDefinitionProviders(); - var workflowDefinitionIds = new List(); + var workflowDefinitions = new List(); + foreach (var provider in providers) { var results = await provider.GetWorkflowsAsync(cancellationToken).AsTask().ToList(); - workflowDefinitionIds.AddRange(results.Select(w => w.Workflow.Id)); - - foreach (var result in results) await AddAsync(result, indexTriggers, cancellationToken); + foreach (var result in results) + { + var workflowDefinition = await AddAsync(result, indexTriggers, cancellationToken); + workflowDefinitions.Add(workflowDefinition); + } } - return workflowDefinitionIds; + return workflowDefinitions; } /// - public Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) + public Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) { return AddAsync(materializedWorkflow, true, cancellationToken); } /// - public async Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default) + public async Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default) { await AssignIdentities(materializedWorkflow.Workflow, cancellationToken); - await AddOrUpdateAsync(materializedWorkflow, cancellationToken); + var workflowDefinition = await AddOrUpdateAsync(materializedWorkflow, cancellationToken); if (indexTriggers) await IndexTriggersAsync(materializedWorkflow, cancellationToken); + + return workflowDefinition; } private async Task AssignIdentities(Workflow workflow, CancellationToken cancellationToken) @@ -92,13 +97,13 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP await _identityGraphService.AssignIdentitiesAsync(workflow, cancellationToken); } - private async Task AddOrUpdateAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) + private async Task AddOrUpdateAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) { await _semaphore.WaitAsync(cancellationToken); try { - await AddOrUpdateCoreAsync(materializedWorkflow, cancellationToken); + return await AddOrUpdateCoreAsync(materializedWorkflow, cancellationToken); } finally { @@ -106,7 +111,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } } - private async Task AddOrUpdateCoreAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) + private async Task AddOrUpdateCoreAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default) { var workflow = materializedWorkflow.Workflow; var definitionId = workflow.Identity.DefinitionId; @@ -178,7 +183,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP if (existingDefinitionVersion is null && workflowDefinitionsToSave.Any(w => w.Id == workflowDefinition.Id)) { _logger.LogInformation("Workflow with ID {WorkflowId} already exists", workflowDefinition.Id); - return; + return workflowDefinition; } workflowDefinitionsToSave.Add(workflowDefinition); @@ -194,7 +199,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } await _workflowDefinitionStore.SaveManyAsync(workflowDefinitionsToSave, cancellationToken); - return; + return workflowDefinition; async Task UpdateIsLatest() { diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs index 3afadafd2..f1a179e51 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs @@ -1,5 +1,6 @@ using Elsa.Mediator.Contracts; using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; using Elsa.Workflows.Runtime.Notifications; namespace Elsa.Workflows.Runtime.Services; @@ -10,8 +11,9 @@ public class WorkflowDefinitionsReloader(IWorkflowDefinitionStorePopulator workf /// public async Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken) { - var definitionIds = await workflowDefinitionStorePopulator.PopulateStoreAsync(true, cancellationToken); - var notification = new WorkflowDefinitionsReloaded(definitionIds); + var workflowDefinitions = await workflowDefinitionStorePopulator.PopulateStoreAsync(true, cancellationToken); + var reloadedWorkflowDefinitions = workflowDefinitions.Select(ReloadedWorkflowDefinition.FromDefinition).ToList(); + var notification = new WorkflowDefinitionsReloaded(reloadedWorkflowDefinitions); await notificationSender.SendAsync(notification, cancellationToken); } } \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index c6588808f..deb69e04d 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -21,9 +21,6 @@ - - Always - Always @@ -81,6 +78,9 @@ Always + + Always + diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs index aaaa665b4..a026a35a4 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -6,8 +6,11 @@ using Elsa.Extensions; using Elsa.Identity.Providers; using Elsa.MassTransit.Extensions; using Elsa.Workflows.ComponentTests.Consumers; +using Elsa.Workflows.ComponentTests.Helpers.Materializers; using Elsa.Workflows.ComponentTests.Helpers.Services; +using Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; using Elsa.Workflows.ComponentTests.Services; +using Elsa.Workflows.Management.Contracts; using FluentStorage; using Hangfire.Annotations; using Microsoft.AspNetCore.Hosting; @@ -94,11 +97,15 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl builder.ConfigureTestServices(services => { - services.AddSingleton(); - services.AddSingleton(); - services.AddSingleton(); - services.AddSingleton(); - services.AddNotificationHandlersFrom(); + services + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddScoped() + .AddNotificationHandlersFrom() + .AddWorkflowDefinitionProvider() + ; }); } diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs new file mode 100644 index 000000000..e011422c4 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Materializers/TestWorkflowMaterializer.cs @@ -0,0 +1,26 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Runtime.Contracts; + +namespace Elsa.Workflows.ComponentTests.Helpers.Materializers; + +/// A workflow materializer that deserializes workflows created from . +public class TestWorkflowMaterializer(IEnumerable workflowProviders) : IWorkflowMaterializer +{ + /// The name of the materializer. + public const string MaterializerName = "Test"; + + /// + public string Name => MaterializerName; + + /// + public ValueTask MaterializeAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + { + var testProvider = (TestWorkflowProvider)workflowProviders.Single(x => x is TestWorkflowProvider); + var materializedWorkflow = testProvider.MaterializedWorkflows.First(x => x.Workflow.Identity.Id == definition.Id); + + return ValueTask.FromResult(materializedWorkflow.Workflow); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs new file mode 100644 index 000000000..5473cf8eb --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/WorkflowProviders/TestWorkflowProvider.cs @@ -0,0 +1,14 @@ +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; + +namespace Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; + +public class TestWorkflowProvider : IWorkflowProvider +{ + public string Name => "Test"; + public ICollection MaterializedWorkflows { get; set; } = new List(); + public ValueTask> GetWorkflowsAsync(CancellationToken cancellationToken = default) + { + return new(MaterializedWorkflows); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs new file mode 100644 index 000000000..120f0fe42 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -0,0 +1,113 @@ +using System.Net; +using Elsa.Common.Models; +using Elsa.Workflows.Activities; +using Elsa.Workflows.ComponentTests.Helpers.Materializers; +using Elsa.Workflows.ComponentTests.Helpers.WorkflowProviders; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Materializers; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowDefinitionReload; + +public class ReloadWorkflowTests : AppComponentTest +{ + private readonly IWorkflowDefinitionManager _workflowDefinitionManager; + private readonly IWorkflowDefinitionsReloader _workflowDefinitionsReloader; + private readonly IWorkflowBuilderFactory _workflowBuilderFactory; + private readonly TestWorkflowProvider _testWorkflowProvider; + private readonly IWorkflowDefinitionService _workflowDefinitionService; + private readonly IActivityRegistry _activityRegistry; + + public ReloadWorkflowTests(App app) : base(app) + { + _workflowDefinitionManager = Scope.ServiceProvider.GetRequiredService(); + _workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService(); + _workflowBuilderFactory = Scope.ServiceProvider.GetRequiredService(); + _workflowDefinitionService = Scope.ServiceProvider.GetRequiredService(); + _activityRegistry = Scope.ServiceProvider.GetRequiredService(); + var workflowProviders = Scope.ServiceProvider.GetRequiredService>(); + _testWorkflowProvider = (TestWorkflowProvider)workflowProviders.First(x => x is TestWorkflowProvider); + } + + [Fact] + public async Task Reloading_AfterRemovingTheWorkflow_ShouldMakeWorkflowReachableAgain() + { + var client = WorkflowServer.CreateHttpWorkflowClient(); + await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); + var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); + var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); + Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); + Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); + } + + [Fact] + public async Task Reloading_AfterUpdatingSourceProvider_ShouldRefreshCaches() + { + var definitionId = Guid.NewGuid().ToString(); + var definitionVersionId1 = Guid.NewGuid().ToString(); + var workflowV1 = await BuildWorkflowAsync(definitionId, definitionVersionId1, 1); + + // Set up the initial workflow version. + _testWorkflowProvider.MaterializedWorkflows = [workflowV1]; + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + var definitionV1 = await _workflowDefinitionService.FindWorkflowGraphAsync(definitionId, VersionOptions.Latest); + Assert.Equal(definitionVersionId1, definitionV1!.Workflow.Identity.Id); + + // Simulate the workflow provider to have a new version available. + var definitionVersionId2 = Guid.NewGuid().ToString(); + var workflowV2 = await BuildWorkflowAsync(definitionId, definitionVersionId2, 2); + _testWorkflowProvider.MaterializedWorkflows = [workflowV1, workflowV2]; + + // Reload the workflow definitions. + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + + // Assert that the workflow definition service finds the updated workflow version. + var definitionV2 = await _workflowDefinitionService.FindWorkflowGraphAsync(definitionId, VersionOptions.Latest); + Assert.Equal(definitionVersionId2, definitionV2!.Workflow.Identity.Id); + } + + [Fact] + public async Task Reloading_AfterUpdatingSourceProvider_ShouldRefreshActivityRegistry() + { + var definitionId = Guid.NewGuid().ToString(); + var definitionVersionId1 = Guid.NewGuid().ToString(); + var workflowV1 = await BuildWorkflowAsync(definitionId, definitionVersionId1, 1); + + // Set up the initial workflow version. + _testWorkflowProvider.MaterializedWorkflows = [workflowV1]; + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + var activityV1 = _activityRegistry.Find(workflowV1.Workflow.Name!); + Assert.Equal(1, activityV1!.Version); + + // Simulate the workflow provider to have a new version available. + var definitionVersionId2 = Guid.NewGuid().ToString(); + var workflowV2 = await BuildWorkflowAsync(definitionId, definitionVersionId2, 2); + _testWorkflowProvider.MaterializedWorkflows = [workflowV1, workflowV2]; + + // Reload the workflow definitions. + await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); + + // Assert that the activity registry contains a new activity descriptor representing the new workflow version. + var activityV2 = _activityRegistry.Find(workflowV2.Workflow.Name!)!; + Assert.Equal(2, activityV2.Version); + } + + private async Task BuildWorkflowAsync(string definitionId, string definitionVersionId, int version) + { + var builder = _workflowBuilderFactory.CreateBuilder(); + builder.DefinitionId = definitionId; + builder.Id = definitionVersionId; + builder.Version = version; + builder.Name = Guid.NewGuid().ToString(); + builder.Root = new WriteLine($"Version {version}"); + builder.WorkflowOptions.UsableAsActivity = true; + var workflow = await builder.BuildWorkflowAsync(); + workflow.Name = builder.Name; + return new MaterializedWorkflow(workflow, _testWorkflowProvider.Name, TestWorkflowMaterializer.MaterializerName); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs deleted file mode 100644 index 390c58d14..000000000 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/RemoveReloadWorkflowTests.cs +++ /dev/null @@ -1,35 +0,0 @@ -using System.Net; -using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Runtime.Contracts; -using Microsoft.Extensions.DependencyInjection; - -namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowDefinitionReload; - -public class RemoveReloadWorkflowTests : AppComponentTest -{ - private readonly IWorkflowDefinitionManager _workflowDefinitionManager; - private readonly IWorkflowDefinitionsReloader _workflowDefinitionsReloader; - - public RemoveReloadWorkflowTests(App app) : base(app) - { - _workflowDefinitionManager = Scope.ServiceProvider.GetRequiredService(); - _workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService(); - } - - [Fact] - public async Task RemovingTheWorkflowThenReload_WorkflowShouldBeReachableAgain() - { - var client = WorkflowServer.CreateHttpWorkflowClient(); - - var result = await _workflowDefinitionManager.DeleteByDefinitionIdAsync("f68b09bc-2013-4617-b82f-d76b6819a624", CancellationToken.None); - - var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); - - await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(CancellationToken.None); - - var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "reload-test")); - - Assert.Equal(HttpStatusCode.NotFound, firstResponse.StatusCode); - Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); - } -} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/http-workflow.json similarity index 100% rename from test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/http-workflow.json rename to test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/http-workflow.json From 149f91e53815d991ae1fd1a6ab446006dee4be85 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 18 Jul 2024 07:51:46 +0200 Subject: [PATCH 15/18] Update API endpoint for Polling Observer (#5787) * Rename and refactor journal update endpoint Replaced `/workflow-instances/{id}/journal/has-updates` endpoint with `/workflow-instances/{id}/updated-at` to simplify API responses. Deleted `HasUpdates` related classes and introduced `GetUpdatedAtResponse` for consistency and clarity. Updated client contracts accordingly. * Remove HasUpdates endpoint and refactor workflow observer Deleted the HasUpdates endpoint and refactored related code to use an updated timestamp approach instead. Improved nullable handling in WorkflowInstanceDesigner and ensured proper observer disposal to avoid memory leaks. Updated workflow observer factory and observer implementations to support observer names and enhanced logging. * Rename updated workflow instance endpoint and handle execution state Renamed the endpoint from "/updated-at" to "/execution-state" to better reflect its purpose. Updated related response models and documentation to capture workflow execution state details such as status, sub-status, and last updated timestamp. * Enable SignalR for real-time workflows Add a flag to use SignalR and activate real-time workflows when enabled. Refactor code to wrap SignalR setup in conditional checks based on the new flag. This enhances the application's interactivity through real-time capabilities. * Remove obsolete endpoints and rename execution state paths Deleted the outdated Api1 and DynamicWorkflows endpoints under Elsa.Server.Web. Also, renamed paths related to execution state models and endpoint to remove "Journal" from the namespace for better clarity and organization. --- .../Endpoints/Api1/Get/Endpoint.cs | 33 ----------------- .../DynamicWorkflows/Post/Endpoint.cs | 37 ------------------- src/bundles/Elsa.Server.Web/Program.cs | 18 +++++---- .../Contracts/IWorkflowInstancesApi.cs | 12 +++--- .../Models/HasJournalUpdateRequest.cs | 11 ------ .../WorkflowInstanceExecutionStateResponse.cs | 6 +++ .../ExecutionState/Endpoint.cs | 37 +++++++++++++++++++ .../ExecutionState/Models.cs | 13 +++++++ .../Journal/HasUpdates/Endpoint.cs | 32 ---------------- .../Journal/HasUpdates/Models.cs | 13 ------- 10 files changed, 73 insertions(+), 139 deletions(-) delete mode 100644 src/bundles/Elsa.Server.Web/Endpoints/Api1/Get/Endpoint.cs delete mode 100644 src/bundles/Elsa.Server.Web/Endpoints/DynamicWorkflows/Post/Endpoint.cs delete mode 100644 src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs create mode 100644 src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Responses/WorkflowInstanceExecutionStateResponse.cs create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Endpoint.cs create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Models.cs delete mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs delete mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs diff --git a/src/bundles/Elsa.Server.Web/Endpoints/Api1/Get/Endpoint.cs b/src/bundles/Elsa.Server.Web/Endpoints/Api1/Get/Endpoint.cs deleted file mode 100644 index 5d727bb33..000000000 --- a/src/bundles/Elsa.Server.Web/Endpoints/Api1/Get/Endpoint.cs +++ /dev/null @@ -1,33 +0,0 @@ -using Elsa.Abstractions; -using Elsa.Workflows.Activities; -using Elsa.Workflows.Models; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Parameters; -using JetBrains.Annotations; - -namespace Elsa.Server.Web.Endpoints.Api1.Get; - -/// -/// Returns a message. -/// -[UsedImplicitly] -public class Get : ElsaEndpointWithoutRequest -{ - /// - public override void Configure() - { - Get("/api-1"); - AllowAnonymous(); - } - - /// - public override async Task HandleAsync(CancellationToken ct) - { - await Task.Delay(1000, ct); - var response = new - { - Message = "OK" - }; - await SendOkAsync(response, ct); - } -} \ No newline at end of file diff --git a/src/bundles/Elsa.Server.Web/Endpoints/DynamicWorkflows/Post/Endpoint.cs b/src/bundles/Elsa.Server.Web/Endpoints/DynamicWorkflows/Post/Endpoint.cs deleted file mode 100644 index f72fb576e..000000000 --- a/src/bundles/Elsa.Server.Web/Endpoints/DynamicWorkflows/Post/Endpoint.cs +++ /dev/null @@ -1,37 +0,0 @@ -using Elsa.Abstractions; -using Elsa.Workflows.Activities; -using Elsa.Workflows.Models; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Options; -using Elsa.Workflows.Runtime.Parameters; - -namespace Elsa.Server.Web.Endpoints.DynamicWorkflows.Post; - -public class Post(IWorkflowRegistry workflowRegistry, IWorkflowRuntime workflowRuntime) : ElsaEndpointWithoutRequest -{ - public override void Configure() - { - Post("/dynamic-workflows"); - AllowAnonymous(); - } - - public override async Task HandleAsync(CancellationToken ct) - { - var workflow = new Workflow - { - Identity = new WorkflowIdentity("DynamicWorkflow1", 1, "DynamicWorkflow1:v1"), - Root = new Sequence - { - Activities = - { - new WriteLine("Step 1"), - new WriteLine("Step 2"), - new WriteLine("Step 3") - } - } - }; - - await workflowRegistry.RegisterAsync(workflow, ct); - await workflowRuntime.StartWorkflowAsync("DynamicWorkflow1", new StartWorkflowRuntimeParams()); - } -} \ No newline at end of file diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index 9c41549da..c4f259015 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -48,6 +48,7 @@ const bool runEFCoreMigrations = true; const bool useMemoryStores = false; const bool useCaching = true; const bool useReadOnlyMode = false; +const bool useSignalR = true; const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit; const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory; @@ -250,7 +251,6 @@ services { api.AddFastEndpointsAssembly(); }) - .UseRealTimeWorkflows() .UseCSharp(options => { options.AppendScript("string Greet(string name) => $\"Hello {name}!\";"); @@ -322,6 +322,11 @@ services elsa.UseQuartz(quartz => { quartz.UseSqlite(sqliteConnectionString); }); } + if (useSignalR) + { + elsa.UseRealTimeWorkflows(); + } + if (useMassTransit) { elsa.UseMassTransit(massTransit => @@ -418,19 +423,18 @@ if (app.Environment.IsDevelopment()) } // SignalR. -app.UseWorkflowsSignalRHubs(); +if (useSignalR) +{ + app.UseWorkflowsSignalRHubs(); +} // Run. app.Run(); -/// /// The main entry point for the application made public for end to end testing. -/// [UsedImplicitly] public partial class Program { - /// /// Set by the test runner to configure the module for testing. - /// public static Action? ConfigureForTest { get; set; } -} +} \ No newline at end of file diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs index 4a50a368c..a225e57ad 100644 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs @@ -1,6 +1,7 @@ using Elsa.Api.Client.Resources.WorkflowDefinitions.Models; using Elsa.Api.Client.Resources.WorkflowInstances.Models; using Elsa.Api.Client.Resources.WorkflowInstances.Requests; +using Elsa.Api.Client.Resources.WorkflowInstances.Responses; using Elsa.Api.Client.Shared.Models; using Refit; @@ -47,14 +48,13 @@ public interface IWorkflowInstancesApi Task> GetFilteredJournalAsync(string workflowInstanceId, GetFilteredJournalRequest? filter, int? skip = default, int? take = default, CancellationToken cancellationToken = default); /// - /// Checks if there are updates in the journal for a specific workflow instance. + /// Returns the execution state of the specified workflow instance. /// - /// The ID of the workflow instance for which to check for updates. - /// The request containing the ID and time from which to check for updates. + /// The ID of the workflow instance for which to return its execution state. /// The cancellation token. - /// Returns whether updates are available for the journal. - [Get("/workflow-instances/{workflowInstanceId}/journal/has-updates")] - Task HasJournalUpdates(string workflowInstanceId, [Query]HasJournalUpdateRequest request, CancellationToken cancellationToken = default); + /// Returns a response containing the execution state. + [Get("/workflow-instances/{workflowInstanceId}/execution-state")] + Task GetExecutionStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default); /// /// Deletes a workflow instance. diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs deleted file mode 100644 index 6efb34477..000000000 --- a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Models/HasJournalUpdateRequest.cs +++ /dev/null @@ -1,11 +0,0 @@ -namespace Elsa.Api.Client.Resources.WorkflowInstances.Models; - -/// A request to update a journal for a workflow instance. -public class HasJournalUpdateRequest -{ - /// The unique identifier of a workflow instance. - public string WorkflowInstanceId { get; set; } = default!; - - /// The start date for checking for updates in the workflow instance journal. - public DateTimeOffset UpdatesSince { get; set; } -} \ No newline at end of file diff --git a/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Responses/WorkflowInstanceExecutionStateResponse.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Responses/WorkflowInstanceExecutionStateResponse.cs new file mode 100644 index 000000000..e3262e1f8 --- /dev/null +++ b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Responses/WorkflowInstanceExecutionStateResponse.cs @@ -0,0 +1,6 @@ +using Elsa.Api.Client.Resources.WorkflowInstances.Enums; + +namespace Elsa.Api.Client.Resources.WorkflowInstances.Responses; + +/// Represents the response containing the last updated timestamp of a workflow instance. +public record WorkflowInstanceExecutionStateResponse(WorkflowStatus Status, WorkflowSubStatus WorkflowSubStatus, DateTimeOffset UpdatedAt); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Endpoint.cs new file mode 100644 index 000000000..96ccc853d --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Endpoint.cs @@ -0,0 +1,37 @@ +using Elsa.Abstractions; +using Elsa.Extensions; +using Elsa.Workflows.Management.Contracts; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.ExecutionState; + +/// Returns the execution state of the specified workflow instance. +[PublicAPI] +internal class ExecutionState(IWorkflowInstanceStore store) : ElsaEndpoint +{ + /// + public override void Configure() + { + Get("/workflow-instances/{id}/execution-state"); + ConfigurePermissions("read:workflow-instances"); + } + + /// + public override async Task HandleAsync(Request request, CancellationToken cancellationToken) + { + var workflowInstance = await store.FindAsync(request.WorkflowInstanceId, cancellationToken); + + if (workflowInstance == null) + { + await SendNotFoundAsync(cancellationToken); + return; + } + + var response = new Response( + workflowInstance.Status, + workflowInstance.SubStatus, + workflowInstance.UpdatedAt); + + await SendOkAsync(response, cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Models.cs new file mode 100644 index 000000000..4a0605239 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/ExecutionState/Models.cs @@ -0,0 +1,13 @@ +using FastEndpoints; + +namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.ExecutionState; + +/// The request to check the execution state of the workflow instance. +public class Request +{ + /// The unique identifier of a workflow instance. + [BindFrom("id")] public string WorkflowInstanceId { get; set; } = default!; +} + +/// Represents the response containing the last updated timestamp of a workflow instance. +public record Response(WorkflowStatus Status, WorkflowSubStatus SubStatus, DateTimeOffset UpdatedAt); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs deleted file mode 100644 index a535e25c9..000000000 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Endpoint.cs +++ /dev/null @@ -1,32 +0,0 @@ -using Elsa.Abstractions; -using Elsa.Common.Entities; -using Elsa.Common.Models; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Filters; -using Elsa.Workflows.Runtime.OrderDefinitions; -using JetBrains.Annotations; - -namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates; - -/// Endpoint that checks if there are updates for a workflow instance. -[PublicAPI] -internal class HasUpdates(IWorkflowExecutionLogStore store) : ElsaEndpoint -{ - /// - public override void Configure() - { - Get("/workflow-instances/{id}/journal/has-updates"); - ConfigurePermissions("read:workflow-instances"); - } - - /// - public override async Task ExecuteAsync(Request request, CancellationToken cancellationToken) - { - var pageArgs = PageArgs.From(1, 1, 0, 1); - var filter = new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = request.WorkflowInstanceId }; - var order = new WorkflowExecutionLogRecordOrder(x => x.Sequence, OrderDirection.Descending); - var pageOfRecords = await store.FindManyAsync(filter, pageArgs, order, cancellationToken); - - return pageOfRecords.Items.Any(item => item.Timestamp >= request.UpdatesSince); - } -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs deleted file mode 100644 index eefafc6d0..000000000 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/HasUpdates/Models.cs +++ /dev/null @@ -1,13 +0,0 @@ -using FastEndpoints; - -namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates; - -/// The request to check if there are updates for a workflow instance journal. -public class Request -{ - /// The unique identifier of a workflow instance. - [BindFrom("id")] public string WorkflowInstanceId { get; set; } = default!; - - /// The start date for checking for updates in the workflow instance journal. - public DateTimeOffset UpdatesSince { get; set; } -} \ No newline at end of file From 481545813d0bcc3e73079d9165c6698a6a02ca75 Mon Sep 17 00:00:00 2001 From: Raymond den Haan Date: Thu, 18 Jul 2024 10:34:48 +0200 Subject: [PATCH 16/18] Skip failing test --- .../Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs index 120f0fe42..b8bd5f819 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -71,7 +71,7 @@ public class ReloadWorkflowTests : AppComponentTest Assert.Equal(definitionVersionId2, definitionV2!.Workflow.Identity.Id); } - [Fact] + [Fact(Skip = "Not working as intended")] public async Task Reloading_AfterUpdatingSourceProvider_ShouldRefreshActivityRegistry() { var definitionId = Guid.NewGuid().ToString(); From 68c6623555412d709c3e4683e83e60988d60a772 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 18 Jul 2024 14:23:39 +0200 Subject: [PATCH 17/18] Fix Workflow Reload Test and Update Activity Registry Logic Unskip the test for workflow reload after updating the source provider and fix the issue with activity name formatting using Humanizer library. Ensure activity names are consistent for both workflow versions in the test. --- .../WorkflowDefinitionReload/ReloadWorkflowTests.cs | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs index b8bd5f819..9a0fea9ca 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -9,6 +9,7 @@ using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Materializers; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Models; +using Humanizer; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Workflows.ComponentTests.Scenarios.WorkflowDefinitionReload; @@ -71,7 +72,7 @@ public class ReloadWorkflowTests : AppComponentTest Assert.Equal(definitionVersionId2, definitionV2!.Workflow.Identity.Id); } - [Fact(Skip = "Not working as intended")] + [Fact] public async Task Reloading_AfterUpdatingSourceProvider_ShouldRefreshActivityRegistry() { var definitionId = Guid.NewGuid().ToString(); @@ -81,7 +82,8 @@ public class ReloadWorkflowTests : AppComponentTest // Set up the initial workflow version. _testWorkflowProvider.MaterializedWorkflows = [workflowV1]; await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); - var activityV1 = _activityRegistry.Find(workflowV1.Workflow.Name!); + var activityTypeName = workflowV1.Workflow.Name.Pascalize(); + var activityV1 = _activityRegistry.Find(activityTypeName); Assert.Equal(1, activityV1!.Version); // Simulate the workflow provider to have a new version available. @@ -93,7 +95,7 @@ public class ReloadWorkflowTests : AppComponentTest await _workflowDefinitionsReloader.ReloadWorkflowDefinitionsAsync(); // Assert that the activity registry contains a new activity descriptor representing the new workflow version. - var activityV2 = _activityRegistry.Find(workflowV2.Workflow.Name!)!; + var activityV2 = _activityRegistry.Find(activityTypeName)!; Assert.Equal(2, activityV2.Version); } @@ -103,11 +105,11 @@ public class ReloadWorkflowTests : AppComponentTest builder.DefinitionId = definitionId; builder.Id = definitionVersionId; builder.Version = version; - builder.Name = Guid.NewGuid().ToString(); + builder.Name = definitionId; builder.Root = new WriteLine($"Version {version}"); builder.WorkflowOptions.UsableAsActivity = true; var workflow = await builder.BuildWorkflowAsync(); - workflow.Name = builder.Name; + workflow.Name = definitionId; return new MaterializedWorkflow(workflow, _testWorkflowProvider.Name, TestWorkflowMaterializer.MaterializerName); } } \ No newline at end of file From e5ba9a8a9d31541ad24fa0ea0853cce5c4cc9601 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 18 Jul 2024 16:00:49 +0200 Subject: [PATCH 18/18] Refactor scheduling and system clock configurations Moved specific services from Scoped to Singleton in SchedulingFeature and SystemClockFeature for better performance and consistency. Refactored DefaultWorkflowScheduler to use a constructor with an IScheduler parameter and removed redundant private field. Removed attribute RequiresUnreferencedCode in WorkflowManagementFeature. --- .../Features/SystemClockFeature.cs | 4 +-- .../Features/SchedulingFeature.cs | 6 ++-- .../Services/DefaultWorkflowScheduler.cs | 28 ++++++------------- .../Features/WorkflowManagementFeature.cs | 1 - 4 files changed, 13 insertions(+), 26 deletions(-) diff --git a/src/modules/Elsa.Common/Features/SystemClockFeature.cs b/src/modules/Elsa.Common/Features/SystemClockFeature.cs index 027319e11..04aea968f 100644 --- a/src/modules/Elsa.Common/Features/SystemClockFeature.cs +++ b/src/modules/Elsa.Common/Features/SystemClockFeature.cs @@ -6,9 +6,7 @@ using Microsoft.Extensions.DependencyInjection; namespace Elsa.Common.Features; -/// /// Configures the system clock. -/// public class SystemClockFeature : FeatureBase { /// @@ -19,6 +17,6 @@ public class SystemClockFeature : FeatureBase /// public override void Apply() { - Services.AddScoped(); + Services.AddSingleton(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs index 368ab9ba2..19c634c6a 100644 --- a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs +++ b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs @@ -42,12 +42,12 @@ public class SchedulingFeature : FeatureBase public override void Apply() { Services + .AddSingleton() + .AddSingleton() + .AddSingleton(CronParser) .AddScoped() .AddScoped() - .AddScoped() .AddScoped() - .AddScoped() - .AddScoped(CronParser) .AddScoped(WorkflowScheduler) .AddHandlersFrom(); diff --git a/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs b/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs index d1a8d22f3..6bde761c2 100644 --- a/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs +++ b/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs @@ -8,22 +8,12 @@ namespace Elsa.Scheduling.Services; /// /// A default implementation of that uses the . /// -public class DefaultWorkflowScheduler : IWorkflowScheduler +public class DefaultWorkflowScheduler(IScheduler scheduler) : IWorkflowScheduler { - private readonly IScheduler _scheduler; - - /// - /// Initializes a new instance of the class. - /// - public DefaultWorkflowScheduler(IScheduler scheduler) - { - _scheduler = scheduler; - } - /// public async ValueTask ScheduleAtAsync(string taskName, DispatchWorkflowDefinitionRequest request, DateTimeOffset at, CancellationToken cancellationToken = default) { - await _scheduler.ScheduleAsync(taskName, new RunWorkflowTask(request), new SpecificInstantSchedule(at), cancellationToken); + await scheduler.ScheduleAsync(taskName, new RunWorkflowTask(request), new SpecificInstantSchedule(at), cancellationToken); } /// @@ -31,7 +21,7 @@ public class DefaultWorkflowScheduler : IWorkflowScheduler { var task = new ResumeWorkflowTask(request); var schedule = new SpecificInstantSchedule(at); - await _scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); + await scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); } /// @@ -39,7 +29,7 @@ public class DefaultWorkflowScheduler : IWorkflowScheduler { var task = new RunWorkflowTask(request); var schedule = new RecurringSchedule(startAt, interval); - await _scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); + await scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); } /// @@ -47,7 +37,7 @@ public class DefaultWorkflowScheduler : IWorkflowScheduler { var task = new ResumeWorkflowTask(request); var schedule = new RecurringSchedule(startAt, interval); - await _scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); + await scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); } /// @@ -55,7 +45,7 @@ public class DefaultWorkflowScheduler : IWorkflowScheduler { var task = new RunWorkflowTask(request); var schedule = new CronSchedule(cronExpression); - await _scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); + await scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); } /// @@ -63,12 +53,12 @@ public class DefaultWorkflowScheduler : IWorkflowScheduler { var task = new ResumeWorkflowTask(request); var schedule = new CronSchedule(cronExpression); - await _scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); + await scheduler.ScheduleAsync(taskName, task, schedule, cancellationToken); } /// - public async ValueTask UnscheduleAsync(string workflowInstanceId, CancellationToken cancellationToken = default) + public async ValueTask UnscheduleAsync(string taskName, CancellationToken cancellationToken = default) { - await _scheduler.ClearScheduleAsync(workflowInstanceId, cancellationToken); + await scheduler.ClearScheduleAsync(taskName, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index cfb45929c..3fdc5adf8 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -101,7 +101,6 @@ public class WorkflowManagementFeature : FeatureBase /// /// Adds all types implementing to the system. /// - [RequiresUnreferencedCode("The assembly containing the specified marker type will be scanned for activity types.")] public WorkflowManagementFeature AddActivitiesFrom() { var activityTypes = typeof(TMarker).Assembly.GetExportedTypes()