diff --git a/Directory.Build.props b/Directory.Build.props index abd6ae599..d901a2f0a 100644 --- a/Directory.Build.props +++ b/Directory.Build.props @@ -30,6 +30,6 @@ $(NoWarn);CS0162;CS1591 - 3.2.0-rc3.417 + 3.2.0-rc3.443 \ No newline at end of file diff --git a/Directory.Packages.props b/Directory.Packages.props index c4fe6f1a8..68c68bcf0 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -33,7 +33,7 @@ - + @@ -44,9 +44,9 @@ - + - + @@ -81,6 +81,8 @@ + + @@ -103,6 +105,7 @@ + @@ -134,8 +137,6 @@ - - @@ -168,8 +169,8 @@ + - \ No newline at end of file diff --git a/Elsa.sln.DotSettings b/Elsa.sln.DotSettings index 2dfe2a8f6..fcce66d42 100644 --- a/Elsa.sln.DotSettings +++ b/Elsa.sln.DotSettings @@ -1,4 +1,4 @@ - + False 1000 1 @@ -12,10 +12,12 @@ EF True True + True True True True True + True True True True diff --git a/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj b/src/bundles/Elsa.Server.Web/Elsa.Server.Web.csproj index f92ae959c..eb5be7ffe 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.Server.Web/Endpoints/Api1/Get/Endpoint.cs b/src/bundles/Elsa.Server.Web/Endpoints/Api1/Get/Endpoint.cs deleted file mode 100644 index 30b0830b9..000000000 --- a/src/bundles/Elsa.Server.Web/Endpoints/Api1/Get/Endpoint.cs +++ /dev/null @@ -1,29 +0,0 @@ -using Elsa.Abstractions; -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 73ca59231..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.Options; -using Elsa.Workflows.Runtime; -using Elsa.Workflows.Runtime.Contracts; -using Elsa.Workflows.Runtime.Parameters; - -namespace Elsa.Server.Web.Endpoints.DynamicWorkflows.Post; - -public class Post(IWorkflowInvoker workflowInvoker) : 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 workflowInvoker.InvokeAsync(workflow, cancellationToken: ct); - } -} \ No newline at end of file diff --git a/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs b/src/bundles/Elsa.Server.Web/Enums/ApplicationRole.cs index 6eac67b27..a4e8c7cf0 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, + Monitor } diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index 17c97ebdb..450d71301 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -49,6 +49,7 @@ const bool useMemoryStores = false; const bool useCaching = true; const bool useAzureServiceBusModule = false; const bool useReadOnlyMode = false; +const bool useSignalR = true; const WorkflowRuntime workflowRuntime = WorkflowRuntime.ProtoActor; const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit; const MassTransitBroker massTransitBroker = MassTransitBroker.Memory; @@ -258,7 +259,6 @@ services { api.AddFastEndpointsAssembly(); }) - .UseRealTimeWorkflows() .UseCSharp(options => { options.AppendScript("string Greet(string name) => $\"Hello {name}!\";"); @@ -330,6 +330,11 @@ services elsa.UseQuartz(quartz => { quartz.UseSqlite(sqliteConnectionString); }); } + if (useSignalR) + { + elsa.UseRealTimeWorkflows(); + } + if (useMassTransit) { elsa.UseMassTransit(massTransit => @@ -443,19 +448,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/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/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()); 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/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs b/src/clients/Elsa.Api.Client/Resources/WorkflowInstances/Contracts/IWorkflowInstancesApi.cs index 61ce1daaa..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; @@ -46,6 +47,15 @@ public interface IWorkflowInstancesApi [Post("/workflow-instances/{workflowInstanceId}/journal")] Task> GetFilteredJournalAsync(string workflowInstanceId, GetFilteredJournalRequest? filter, int? skip = default, int? take = default, CancellationToken cancellationToken = default); + /// + /// Returns the execution state of the specified workflow instance. + /// + /// The ID of the workflow instance for which to return its execution state. + /// The cancellation token. + /// 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/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/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs b/src/common/Elsa.Api.Common/Extensions/WebApplicationExtensions.cs index aaf4722ce..7d6f498a9 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(DateTimeOffset.TryParse(s.ToString(),CultureInfo.InvariantCulture,DateTimeStyles.RoundtripKind, out var result), result)); }); /// 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.Common/Features/SystemClockFeature.cs b/src/modules/Elsa.Common/Features/SystemClockFeature.cs index 76ebc55f3..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 { /// diff --git a/src/modules/Elsa.Dapper/Elsa.Dapper.csproj b/src/modules/Elsa.Dapper/Elsa.Dapper.csproj index d35cf74d5..39461f1aa 100644 --- a/src/modules/Elsa.Dapper/Elsa.Dapper.csproj +++ b/src/modules/Elsa.Dapper/Elsa.Dapper.csproj @@ -34,6 +34,7 @@ + diff --git a/src/modules/Elsa.EntityFrameworkCore.Common/Elsa.EntityFrameworkCore.Common.csproj b/src/modules/Elsa.EntityFrameworkCore.Common/Elsa.EntityFrameworkCore.Common.csproj index 2531e7402..35dc2514a 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..2a69f65ae 100644 --- a/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj +++ b/src/modules/Elsa.EntityFrameworkCore.SqlServer/Elsa.EntityFrameworkCore.SqlServer.csproj @@ -20,5 +20,7 @@ + + \ 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 ef4e0f35e..95f066cc2 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.Http/Handlers/InvalidateHttpWorkflowsCache.cs b/src/modules/Elsa.Http/Handlers/InvalidateHttpWorkflowsCache.cs index ff281af71..0fddca961 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,13 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo await InvalidateCacheAsync(notification.IndexedWorkflowTriggers.Workflow.Identity.DefinitionId); } + /// + public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken) + { + foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions) + await InvalidateCacheAsync(reloadedWorkflowDefinition.DefinitionId); + } + private async Task InvalidateCacheAsync(string workflowDefinitionId) { await httpWorkflowsCacheManager.EvictWorkflowAsync(workflowDefinitionId); @@ -103,17 +111,17 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo WorkflowDefinitionVersionId = workflowDefinitionVersionId }; var triggers = await triggerStore.FindManyAsync(filter, cancellationToken); - + await InvalidateTriggerCacheAsync(triggers, cancellationToken); } 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 = 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.JavaScript/Services/JintJavaScriptEvaluator.cs b/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs index 5902b2c9c..6b39549fb 100644 --- a/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs +++ b/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs @@ -80,6 +80,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)); diff --git a/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs b/src/modules/Elsa.MassTransit/Consumers/WorkflowDefinitionEventsConsumer.cs index 6fed5a9a5..4f7821a63 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.ReloadedWorkflowDefinitions); + 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 3783944d3..10e654479 100644 --- a/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs +++ b/src/modules/Elsa.MassTransit/Handlers/DistributedWorkflowDefinitionNotificationsHandler.cs @@ -17,7 +17,8 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IBus bus) : INotificationHandler, INotificationHandler, INotificationHandler, - INotificationHandler + INotificationHandler, + INotificationHandler { /// public Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken) @@ -72,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 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 new file mode 100644 index 000000000..138492683 --- /dev/null +++ b/src/modules/Elsa.MassTransit/Messages/WorkflowDefinitionsReloaded.cs @@ -0,0 +1,10 @@ +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 reloadedWorkflowDefinitions) +{ + /// The reloaded workflow definitions. + public ICollection ReloadedWorkflowDefinitions { get; set; } = reloadedWorkflowDefinitions; +} \ 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.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..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,5 +23,7 @@ + + \ 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.Scheduling/Features/SchedulingFeature.cs b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs index 4ee9fd2f2..496106b27 100644 --- a/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs +++ b/src/modules/Elsa.Scheduling/Features/SchedulingFeature.cs @@ -41,6 +41,9 @@ public class SchedulingFeature : FeatureBase public override void Apply() { Services + .AddSingleton() + .AddSingleton() + .AddSingleton(CronParser) .AddScoped() .AddScoped() .AddSingleton() diff --git a/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs b/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs index f005c4de0..a2c732906 100644 --- a/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs +++ b/src/modules/Elsa.Scheduling/Services/DefaultWorkflowScheduler.cs @@ -6,22 +6,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, ScheduleNewWorkflowInstanceRequest 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); } /// @@ -29,7 +19,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); } /// @@ -37,7 +27,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); } /// @@ -45,7 +35,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); } /// @@ -53,7 +43,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); } /// @@ -61,12 +51,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.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 @@ - + + + + + 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 fe636216c..f076d9d5f 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.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.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/RealTime/Handlers/BroadcastWorkflowProgress.cs b/src/modules/Elsa.Workflows.Api/RealTime/Handlers/BroadcastWorkflowProgress.cs index 79d033dce..56dc338f5 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); 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/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() diff --git a/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs b/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs index 0aef88290..01cd75881 100644 --- a/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs +++ b/src/modules/Elsa.Workflows.Management/Handlers/DeleteWorkflowInstances.cs @@ -1,12 +1,14 @@ using Elsa.Mediator.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.Management/Services/WorkflowDefinitionImporter.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs index 1cbb75299..4b7366185 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionImporter.cs @@ -62,7 +62,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(); diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs index f42fee7b6..8b0f93481 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionStorePopulator.cs @@ -11,27 +11,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 new file mode 100644 index 000000000..f17d06edf --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowDefinitionsReloader.cs @@ -0,0 +1,8 @@ +namespace Elsa.Workflows.Runtime.Contracts; + +/// Reloads all workflows by invoking the populator. +public interface IWorkflowDefinitionsReloader +{ + /// 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/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/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 5b0886b63..7359d7236 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -10,7 +10,6 @@ using Elsa.Workflows.Contracts; using Elsa.Workflows.Features; using Elsa.Workflows.Management; 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; @@ -189,6 +188,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() @@ -207,6 +207,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() @@ -250,10 +251,10 @@ public class WorkflowRuntimeFeature : FeatureBase .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() - .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() .AddNotificationHandler() // Workflow activation strategies. 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..d2d52613c --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/InvalidateWorkflowsCache.cs @@ -0,0 +1,24 @@ +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 the workflow 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 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 new file mode 100644 index 000000000..336347ff4 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionsReloaded.cs @@ -0,0 +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 ReloadedWorkflowDefinitions) : INotification; \ No newline at end of file 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 a27916890..bcb1042ca 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -47,37 +47,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 workflowDefinitions = new List(); + foreach (var provider in providers) { var results = await provider.GetWorkflowsAsync(cancellationToken).AsTask().ToList(); - 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 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) @@ -85,13 +95,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 { @@ -99,7 +109,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; @@ -129,8 +139,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. @@ -171,7 +181,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); @@ -187,7 +197,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/StoreWorkflowExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs index fda34fc4b..5aa5f87c7 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,35 +8,13 @@ namespace Elsa.Workflows.Runtime; /// /// 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(); - - await store.AddManyAsync(records, context.CancellationToken); - await notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationToken); + var records = extractor.ExtractWorkflowExecutionLogs(context).ToList(); + await store.AddManyAsync(records, context.CancellationTokens.SystemCancellationToken); + await notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationTokens.SystemCancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs index 15b5ff7c8..98b2d2736 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..f1a179e51 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloader.cs @@ -0,0 +1,19 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; +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 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/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionContext.cs new file mode 100644 index 000000000..577559bdb --- /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.Services; + +/// +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 diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index 82e981703..28c11218b 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -75,6 +75,12 @@ Always + + 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 73681de34..eb1caed0e 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -1,4 +1,4 @@ -using System.Net.Http.Headers; +using System.Net.Http.Headers; using System.Reflection; using Elsa.Alterations.Extensions; using Elsa.Common.Contracts; @@ -15,6 +15,7 @@ using Elsa.Testing.Shared; using Elsa.Testing.Shared.Handlers; using Elsa.Testing.Shared.Services; using Elsa.Workflows.ComponentTests.Consumers; +using Elsa.Workflows.ComponentTests.Helpers.Materializers; using Elsa.Workflows.ComponentTests.Helpers.Services; using Elsa.Workflows.Runtime.Distributed.Extensions; using FluentStorage; @@ -119,13 +120,16 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl builder.ConfigureTestServices(services => { - services.AddSingleton(); - services.AddSingleton(); - services.AddSingleton(); - services.AddSingleton(); - services.AddScoped(); - services.AddNotificationHandlersFrom(); - services.AddNotificationHandlersFrom(); + services + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddScoped() + .AddNotificationHandlersFrom() + .AddWorkflowDefinitionProvider() + .AddNotificationHandlersFrom() + ; }); } 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/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/ReloadWorkflowTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs new file mode 100644 index 000000000..9a0fea9ca --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/ReloadWorkflowTests.cs @@ -0,0 +1,115 @@ +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 Humanizer; +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 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. + 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(activityTypeName)!; + 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 = definitionId; + builder.Root = new WriteLine($"Version {version}"); + builder.WorkflowOptions.UsableAsActivity = true; + var workflow = await builder.BuildWorkflowAsync(); + workflow.Name = definitionId; + return new MaterializedWorkflow(workflow, _testWorkflowProvider.Name, TestWorkflowMaterializer.MaterializerName); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/http-workflow.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/http-workflow.json new file mode 100644 index 000000000..c885f29b6 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionReload/Workflows/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