Merge remote-tracking branch 'origin/patch/3.2.x'

This commit is contained in:
Sipke Schoorstra 2024-07-18 18:25:49 +02:00
commit 827680b3c4
78 changed files with 904 additions and 273 deletions

View file

@ -30,6 +30,6 @@
<NoWarn>$(NoWarn);CS0162;CS1591</NoWarn>
</PropertyGroup>
<PropertyGroup>
<ElsaStudioVersion>3.2.0-rc3.417</ElsaStudioVersion>
<ElsaStudioVersion>3.2.0-rc3.443</ElsaStudioVersion>
</PropertyGroup>
</Project>

View file

@ -33,7 +33,7 @@
<PackageVersion Include="FluentMigrator.Runner" Version="5.2.0" />
<PackageVersion Include="FluentStorage" Version="5.4.3" />
<PackageVersion Include="FluentStorage.Azure.Blobs" Version="5.2.3" />
<PackageVersion Include="Fluid.Core" Version="2.10.0" />
<PackageVersion Include="Fluid.Core" Version="2.11.0" />
<PackageVersion Include="Fody" Version="6.8.1" />
<PackageVersion Include="GitHubActionsTestLogger" Version="2.4.1" />
<PackageVersion Include="Google.Protobuf" Version="3.27.2" />
@ -44,9 +44,9 @@
<PackageVersion Include="Humanizer.Core" Version="2.14.1" />
<PackageVersion Include="IronCompress" Version="1.5.2" />
<PackageVersion Include="JetBrains.Annotations" Version="2024.2.0" />
<PackageVersion Include="Jint" Version="3.1.4" />
<PackageVersion Include="Jint" Version="3.1.5" />
<PackageVersion Include="LinqKit.Core" Version="1.2.5" />
<PackageVersion Include="MailKit" Version="4.7.0" />
<PackageVersion Include="MailKit" Version="4.7.1.1" />
<PackageVersion Include="MassTransit" Version="8.2.3" />
<PackageVersion Include="MassTransit.Azure.ServiceBus.Core" Version="8.2.3" />
<PackageVersion Include="MassTransit.Extensions.DependencyInjection" Version="7.3.1" />
@ -81,6 +81,8 @@
<PackageVersion Include="Quartz.Extensions.DependencyInjection" Version="3.11.0" />
<PackageVersion Include="Quartz.Extensions.Hosting" Version="3.11.0" />
<PackageVersion Include="Quartz.Serialization.Json" Version="3.11.0" />
<PackageVersion Include="Refit" Version="7.1.2" />
<PackageVersion Include="Refit.HttpClientFactory" Version="7.0.0" />
<PackageVersion Include="Scrutor" Version="4.2.2" />
<PackageVersion Include="ShortGuid" Version="2.0.1" />
<PackageVersion Include="StackExchange.Redis" Version="2.8.0" />
@ -103,6 +105,7 @@
<PackageVersion Include="AppAny.Quartz.EntityFrameworkCore.Migrations.PostgreSQL" Version="0.5.1" />
<PackageVersion Include="AppAny.Quartz.EntityFrameworkCore.Migrations.SQLite" Version="0.5.1" />
<PackageVersion Include="AppAny.Quartz.EntityFrameworkCore.Migrations.SqlServer" Version="0.5.1" />
<PackageVersion Include="System.Text.Json" Version="8.0.4" />
</ItemGroup>
<ItemGroup Condition="'$(TargetFramework)' == 'net6.0' or '$(TargetFramework)' == 'net7.0'">
<PackageVersion Include="AspNetCore.Authentication.ApiKey" Version="7.0.0" />
@ -134,8 +137,6 @@
<PackageVersion Include="Oracle.EntityFrameworkCore" Version="7.21.13" />
<PackageVersion Include="Polly" Version="7.2.4" />
<PackageVersion Include="Pomelo.EntityFrameworkCore.MySql" Version="7.0.0" />
<PackageVersion Include="Refit" Version="7.0.0" />
<PackageVersion Include="Refit.HttpClientFactory" Version="7.0.0" />
<PackageVersion Include="System.Text.Json" Version="7.0.4" />
</ItemGroup>
<ItemGroup Condition="'$(TargetFramework)' == 'net8.0'">
@ -168,8 +169,8 @@
<PackageVersion Include="Oracle.EntityFrameworkCore" Version="8.23.40" />
<PackageVersion Include="Polly" Version="8.4.1" />
<PackageVersion Include="Pomelo.EntityFrameworkCore.MySql" Version="8.0.2" />
<PackageVersion Include="System.Text.Json" Version="8.0.4" />
<PackageVersion Include="Refit" Version="7.1.2" />
<PackageVersion Include="Refit.HttpClientFactory" Version="7.1.2" />
<PackageVersion Include="System.Text.Json" Version="8.0.4" />
</ItemGroup>
</Project>

View file

@ -1,4 +1,4 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:Boolean x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/KEEP_EXISTING_INITIALIZER_ARRANGEMENT/@EntryValue">False</s:Boolean>
<s:Int64 x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/MAX_ARRAY_INITIALIZER_ELEMENTS_ON_LINE/@EntryValue">1000</s:Int64>
<s:Int64 x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/MAX_INITIALIZER_ELEMENTS_ON_LINE/@EntryValue">1</s:Int64>
@ -12,10 +12,12 @@
<s:String x:Key="/Default/CodeStyle/Naming/CSharpNaming/Abbreviations/=EF/@EntryIndexedValue">EF</s:String>
<s:Boolean x:Key="/Default/UserDictionary/Words/=downloadables/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=initializable/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=materializer/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=materializers/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Persister/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Populator/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Postgre/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=reloader/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=resumer/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=startable/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Telnyx/@EntryIndexedValue">True</s:Boolean>

View file

@ -53,9 +53,9 @@
<PackageReference Include="Proto.Persistence.SqlServer"/>
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions of Npgsql.-->
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="Npgsql" VersionOverride="8.0.3"/>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -1,29 +0,0 @@
using Elsa.Abstractions;
using JetBrains.Annotations;
namespace Elsa.Server.Web.Endpoints.Api1.Get;
/// <summary>
/// Returns a message.
/// </summary>
[UsedImplicitly]
public class Get : ElsaEndpointWithoutRequest
{
/// <inheritdoc />
public override void Configure()
{
Get("/api-1");
AllowAnonymous();
}
/// <inheritdoc />
public override async Task HandleAsync(CancellationToken ct)
{
await Task.Delay(1000, ct);
var response = new
{
Message = "OK"
};
await SendOkAsync(response, ct);
}
}

View file

@ -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);
}
}

View file

@ -4,5 +4,6 @@ public enum ApplicationRole
{
Default,
Api,
Worker
Worker,
Monitor
}

View file

@ -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<Program>();
})
.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();
/// <summary>
/// The main entry point for the application made public for end to end testing.
/// </summary>
[UsedImplicitly]
public partial class Program
{
/// <summary>
/// Set by the test runner to configure the module for testing.
/// </summary>
public static Action<IModule>? ConfigureForTest { get; set; }
}

View file

@ -40,4 +40,9 @@
<PackageReference Include="Elsa.Studio.Login.BlazorWasm"/>
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -15,6 +15,11 @@
<PackageReference Include="Elsa.Studio.Login.BlazorWasm" />
<PackageReference Include="Microsoft.AspNetCore.Components.WebAssembly.Server" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<Folder Include="wwwroot\" />

View file

@ -16,5 +16,10 @@
<PackageReference Include="Refit" />
<PackageReference Include="Refit.HttpClientFactory" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -20,21 +20,17 @@ using static Elsa.Api.Client.RefitSettingsHelper;
namespace Elsa.Api.Client.Extensions;
/// <summary>
/// Provides extension methods for dependency injection.
/// </summary>
[PublicAPI]
public static class DependencyInjectionExtensions
{
/// <summary>
/// Adds the Elsa API client configured to use an API key to the service collection.
/// </summary>
public static IServiceCollection AddElsaApiKeyClient(this IServiceCollection services, Action<ElsaClientOptions> configureOptions)
/// Adds default Elsa API clients configured to use an API key.
public static IServiceCollection AddDefaultApiClientsUsingApiKey(this IServiceCollection services, Action<ElsaClientOptions> 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
});
}
/// <summary>
/// Adds the Elsa client to the service collection.
/// </summary>
public static IServiceCollection AddElsaClient(this IServiceCollection services, Action<ElsaClientBuilderOptions> configureClient)
/// Adds default Elsa API clients.
public static IServiceCollection AddDefaultApiClients(this IServiceCollection services, Action<ElsaClientBuilderOptions> 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<ElsaClientOptions>(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<IWorkflowDefinitionsApi>(builderOptions);
services.AddApi<IExecuteWorkflowApi>(builderOptionsWithoutRetryPolicy);
services.AddApi<IWorkflowInstancesApi>(builderOptions);
services.AddApi<IActivityDescriptorsApi>(builderOptions);
services.AddApi<IActivityDescriptorOptionsApi>(builderOptions);
services.AddApi<IActivityExecutionsApi>(builderOptions);
services.AddApi<IStorageDriversApi>(builderOptions);
services.AddApi<IVariableTypesApi>(builderOptions);
services.AddApi<IWorkflowActivationStrategiesApi>(builderOptions);
services.AddApi<IIncidentStrategiesApi>(builderOptions);
services.AddApi<ILoginApi>(builderOptions);
services.AddApi<IFeaturesApi>(builderOptions);
services.AddApi<IJavaScriptApi>(builderOptions);
services.AddApi<IExpressionDescriptorsApi>(builderOptions);
services.AddApi<IWorkflowContextProviderDescriptorsApi>(builderOptions);
});
}
var builderOptionsWithoutRetryPolicy = new ElsaClientBuilderOptions
/// <summary>
/// Adds an API client to the service collection. Requires AddElsaClient to be called exactly once.
/// </summary>
public static IServiceCollection AddApiClient<T>(this IServiceCollection services, Action<ElsaClientBuilderOptions> configureClient) where T : class
{
return services.AddApiClients(configureClient, builderOptions => services.AddApi<T>(builderOptions));
}
/// Adds the Elsa client to the service collection.
public static IServiceCollection AddApiClients(this IServiceCollection services, Action<ElsaClientBuilderOptions> configureClient, Action<ElsaClientBuilderOptions>? 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<ElsaClientOptions>(options =>
{
options.BaseAddress = builderOptions.BaseAddress;
options.ConfigureHttpClient = builderOptions.ConfigureHttpClient;
options.ApiKey = builderOptions.ApiKey;
});
configureServices?.Invoke(builderOptions);
}
services.AddApi<IWorkflowDefinitionsApi>(builderOptions);
services.AddApi<IExecuteWorkflowApi>(builderOptionsWithoutRetryPolicy);
services.AddApi<IWorkflowInstancesApi>(builderOptions);
services.AddApi<IActivityDescriptorsApi>(builderOptions);
services.AddApi<IActivityDescriptorOptionsApi>(builderOptions);
services.AddApi<IActivityExecutionsApi>(builderOptions);
services.AddApi<IStorageDriversApi>(builderOptions);
services.AddApi<IVariableTypesApi>(builderOptions);
services.AddApi<IWorkflowActivationStrategiesApi>(builderOptions);
services.AddApi<IIncidentStrategiesApi>(builderOptions);
services.AddApi<ILoginApi>(builderOptions);
services.AddApi<IFeaturesApi>(builderOptions);
services.AddApi<IJavaScriptApi>(builderOptions);
services.AddApi<IExpressionDescriptorsApi>(builderOptions);
services.AddApi<IWorkflowContextProviderDescriptorsApi>(builderOptions);
return services;
}
@ -94,11 +111,23 @@ public static class DependencyInjectionExtensions
/// <param name="services">The service collection.</param>
/// <param name="httpClientBuilderOptions">An options object that can be used to configure the HTTP client builder.</param>
/// <typeparam name="T">The type representing the API.</typeparam>
public static void AddApi<T>(this IServiceCollection services, ElsaClientBuilderOptions? httpClientBuilderOptions = default) where T : class
public static IServiceCollection AddApi<T>(this IServiceCollection services, ElsaClientBuilderOptions? httpClientBuilderOptions = default) where T : class
{
var builder = services.AddRefitClient<T>(_ => CreateRefitSettings(), typeof(T).Name).ConfigureHttpClient(ConfigureElsaApiHttpClient);
return services.AddApi(typeof(T), httpClientBuilderOptions);
}
/// <summary>
/// Adds a refit client for the specified API type.
/// </summary>
/// <param name="services">The service collection.</param>
/// <param name="apiType">The type representing the API</param>
/// <param name="httpClientBuilderOptions">An options object that can be used to configure the HTTP client builder.</param>
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;
}
/// <summary>
@ -115,9 +144,7 @@ public static class DependencyInjectionExtensions
httpClientBuilderOptions?.ConfigureHttpClientBuilder(builder);
}
/// <summary>
/// Creates an API client for the specified API type.
/// </summary>
public static T CreateApi<T>(this IServiceProvider serviceProvider, Uri baseAddress) where T : class
{
var httpClientFactory = serviceProvider.GetRequiredService<IHttpClientFactory>();
@ -126,9 +153,7 @@ public static class DependencyInjectionExtensions
return CreateApi<T>(serviceProvider, httpClient);
}
/// <summary>
/// Creates an API client for the specified API type.
/// </summary>
public static T CreateApi<T>(this IServiceProvider serviceProvider, HttpClient httpClient) where T : class
{
return RestService.For<T>(httpClient, CreateRefitSettings());

View file

@ -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<T>(json);
}

View file

@ -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<PagedListResponse<WorkflowExecutionLogRecord>> GetFilteredJournalAsync(string workflowInstanceId, GetFilteredJournalRequest? filter, int? skip = default, int? take = default, CancellationToken cancellationToken = default);
/// <summary>
/// Returns the execution state of the specified workflow instance.
/// </summary>
/// <param name="workflowInstanceId">The ID of the workflow instance for which to return its execution state.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>Returns a response containing the execution state.</returns>
[Get("/workflow-instances/{workflowInstanceId}/execution-state")]
Task<WorkflowInstanceExecutionStateResponse> GetExecutionStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
/// <summary>
/// Deletes a workflow instance.
/// </summary>

View file

@ -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);

View file

@ -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<DateTimeOffset>(s =>
new(DateTimeOffset.TryParse(s.ToString(),CultureInfo.InvariantCulture,DateTimeStyles.RoundtripKind, out var result), result));
});
/// <summary>

View file

@ -15,6 +15,11 @@
<PackageReference Include="ThrottleDebounce" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Formats.Asn1" VersionOverride="8.0.1" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.DropIns.Core\Elsa.DropIns.Core.csproj" />
</ItemGroup>

View file

@ -6,9 +6,7 @@ using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Common.Features;
/// <summary>
/// Configures the system clock.
/// </summary>
public class SystemClockFeature : FeatureBase
{
/// <inheritdoc />

View file

@ -34,6 +34,7 @@
<ItemGroup>
<PackageReference Include="Azure.Identity" VersionOverride="1.11.4" />
<PackageReference Include="Microsoft.Identity.Client" VersionOverride="4.61.3" />
<PackageReference Include="System.Formats.Asn1" VersionOverride="8.0.1" />
</ItemGroup>
</Project>

View file

@ -17,6 +17,11 @@
<PackageReference Include="System.CommandLine" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Common\Elsa.Common.csproj" />
<ProjectReference Include="..\Elsa.Tenants\Elsa.Tenants.csproj" />

View file

@ -12,6 +12,11 @@
<PackageReference Include="Pomelo.EntityFrameworkCore.MySql" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.EntityFrameworkCore\Elsa.EntityFrameworkCore.csproj" />
</ItemGroup>

View file

@ -12,13 +12,11 @@
<PackageReference Include="Microsoft.EntityFrameworkCore.Design" PrivateAssets="all" />
</ItemGroup>
<ItemGroup Condition="'$(TargetFramework)' == 'net6.0' or '$(TargetFramework)' == 'net7.0'">
<PackageReference Include="Npgsql" VersionOverride="7.0.7" />
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup Condition="'$(TargetFramework)' == 'net8.0'">
<PackageReference Include="Npgsql" VersionOverride="8.0.3" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.EntityFrameworkCore\Elsa.EntityFrameworkCore.csproj" />
</ItemGroup>

View file

@ -20,5 +20,7 @@
<ItemGroup>
<PackageReference Include="Azure.Identity" VersionOverride="1.11.4" />
<PackageReference Include="Microsoft.Identity.Client" VersionOverride="4.61.3" />
<PackageReference Include="System.Formats.Asn1" VersionOverride="8.0.1" />
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -13,6 +13,11 @@
<PackageReference Include="Microsoft.EntityFrameworkCore.Design" PrivateAssets="all" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.EntityFrameworkCore\Elsa.EntityFrameworkCore.csproj" />
</ItemGroup>

View file

@ -11,6 +11,11 @@
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Alterations.Core\Elsa.Alterations.Core.csproj" />
<ProjectReference Include="..\Elsa.Alterations\Elsa.Alterations.csproj" />

View file

@ -15,6 +15,11 @@
<PackageReference Include="FluentStorage" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Workflows.Core\Elsa.Workflows.Core.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Management\Elsa.Workflows.Management.csproj" />

View file

@ -10,6 +10,11 @@
<ItemGroup>
<PackageReference Include="FluentStorage" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Liquid\Elsa.Liquid.csproj" />

View file

@ -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.
/// </summary>
[UsedImplicitly]
public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflowsCacheManager,
ITriggerStore triggerStore,
IHttpWorkflowsCacheManager cacheManager) :
public class InvalidateHttpWorkflowsCache(
IHttpWorkflowsCacheManager httpWorkflowsCacheManager,
ITriggerStore triggerStore) :
INotificationHandler<WorkflowDefinitionPublished>,
INotificationHandler<WorkflowDefinitionRetracted>,
INotificationHandler<WorkflowDefinitionVersionsUpdated>,
INotificationHandler<WorkflowDefinitionDeleted>,
INotificationHandler<WorkflowDefinitionDeleted>,
INotificationHandler<WorkflowDefinitionsDeleted>,
INotificationHandler<WorkflowDefinitionVersionDeleted>,
INotificationHandler<WorkflowDefinitionVersionsDeleted>,
INotificationHandler<WorkflowTriggersIndexed>
INotificationHandler<WorkflowTriggersIndexed>,
INotificationHandler<WorkflowDefinitionsReloaded>
{
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionPublished notification, CancellationToken cancellationToken)
@ -91,6 +92,13 @@ public class InvalidateHttpWorkflowsCache(IHttpWorkflowsCacheManager httpWorkflo
await InvalidateCacheAsync(notification.IndexedWorkflowTriggers.Workflow.Identity.DefinitionId);
}
/// <inheritdoc />
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<StoredTrigger> 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);
}
}

View file

@ -80,6 +80,9 @@ public class JintJavaScriptEvaluator(IConfiguration configuration, INotification
configureEngine?.Invoke(engine);
// Add common functions.
engine.SetValue("getWorkflowDefinitionId", (Func<string>)(() => context.GetWorkflowExecutionContext().Workflow.Identity.DefinitionId));
engine.SetValue("getWorkflowDefinitionVersionId", (Func<string>)(() => context.GetWorkflowExecutionContext().Workflow.Identity.Id));
engine.SetValue("getWorkflowDefinitionVersion", (Func<int>)(() => context.GetWorkflowExecutionContext().Workflow.Identity.Version));
engine.SetValue("getWorkflowInstanceId", (Func<string>)(() => context.GetActivityExecutionContext().WorkflowExecutionContext.Id));
engine.SetValue("setCorrelationId", (Action<string?>)(value => context.GetActivityExecutionContext().WorkflowExecutionContext.CorrelationId = value));
engine.SetValue("getCorrelationId", (Func<string?>)(() => context.GetActivityExecutionContext().WorkflowExecutionContext.CorrelationId));

View file

@ -6,20 +6,23 @@ using Humanizer;
namespace Elsa.JavaScript.TypeDefinitions.Providers;
/// <summary>
/// Produces <see cref="FunctionDefinition"/>s for common functions.
/// </summary>
internal class CommonFunctionsDefinitionProvider : FunctionDefinitionProvider
internal class CommonFunctionsDefinitionProvider(ITypeAliasRegistry typeAliasRegistry) : FunctionDefinitionProvider
{
private readonly ITypeAliasRegistry _typeAliasRegistry;
public CommonFunctionsDefinitionProvider(ITypeAliasRegistry typeAliasRegistry)
{
_typeAliasRegistry = typeAliasRegistry;
}
protected override IEnumerable<FunctionDefinition> 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));

View file

@ -18,7 +18,8 @@ public class WorkflowDefinitionEventsConsumer(IWorkflowDefinitionActivityRegistr
global::MassTransit.IConsumer<WorkflowDefinitionVersionDeleted>,
global::MassTransit.IConsumer<WorkflowDefinitionVersionsDeleted>,
global::MassTransit.IConsumer<WorkflowDefinitionVersionsUpdated>,
global::MassTransit.IConsumer<WorkflowDefinitionsRefreshed>
global::MassTransit.IConsumer<WorkflowDefinitionsRefreshed>,
global::MassTransit.IConsumer<WorkflowDefinitionsReloaded>
{
/// <inheritdoc />
public Task Consume(ConsumeContext<WorkflowDefinitionDeleted> 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;
}
/// <inheritdoc />
public async Task Consume(ConsumeContext<WorkflowDefinitionsReloaded> 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)

View file

@ -17,7 +17,8 @@ public class DistributedWorkflowDefinitionNotificationsHandler(IBus bus) :
INotificationHandler<WorkflowDefinitionVersionDeleted>,
INotificationHandler<WorkflowDefinitionVersionsDeleted>,
INotificationHandler<WorkflowDefinitionVersionsUpdated>,
INotificationHandler<WorkflowDefinitionsRefreshed>
INotificationHandler<WorkflowDefinitionsRefreshed>,
INotificationHandler<WorkflowDefinitionsReloaded>
{
/// <inheritdoc />
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);
}
/// <inheritdoc />
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);
}
}

View file

@ -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<ReloadedWorkflowDefinition> reloadedWorkflowDefinitions)
{
/// The reloaded workflow definitions.
public ICollection<ReloadedWorkflowDefinition> ReloadedWorkflowDefinitions { get; set; } = reloadedWorkflowDefinitions;
}

View file

@ -6,7 +6,7 @@ public static class AmbientConsumerScope
private static readonly AsyncLocal<bool> 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;

View file

@ -13,6 +13,11 @@
<PackageReference Include="MongoDB.Driver.Core.Extensions.DiagnosticSources" />
<PackageReference Include="MongoDB.Driver.Extensions" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\common\Elsa.Features\Elsa.Features.csproj" />

View file

@ -13,6 +13,11 @@
<PackageReference Include="Microsoft.EntityFrameworkCore.Design" PrivateAssets="all" />
<PackageReference Include="Quartz.Serialization.Json" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.EntityFrameworkCore.Common\Elsa.EntityFrameworkCore.Common.csproj" />

View file

@ -14,12 +14,16 @@
<PackageReference Include="Quartz.Serialization.Json" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup Condition="'$(TargetFramework)' == 'net6.0' or '$(TargetFramework)' == 'net7.0'">
<PackageReference Include="Npgsql.EntityFrameworkCore.PostgreSQL" VersionOverride="7.0.18" />
</ItemGroup>
<ItemGroup Condition="'$(TargetFramework)' == 'net8.0'">
<PackageReference Include="Npgsql.EntityFrameworkCore.PostgreSQL" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.EntityFrameworkCore.Common\Elsa.EntityFrameworkCore.Common.csproj" />

View file

@ -23,5 +23,7 @@
<ItemGroup>
<PackageReference Include="Azure.Identity" VersionOverride="1.11.4" />
<PackageReference Include="Microsoft.Identity.Client" VersionOverride="4.61.3" />
<PackageReference Include="System.Formats.Asn1" VersionOverride="8.0.1" />
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -13,6 +13,11 @@
<PackageReference Include="Microsoft.EntityFrameworkCore.Design" PrivateAssets="all" />
<PackageReference Include="Quartz.Serialization.Json" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.EntityFrameworkCore.Common\Elsa.EntityFrameworkCore.Common.csproj" />

View file

@ -41,6 +41,9 @@ public class SchedulingFeature : FeatureBase
public override void Apply()
{
Services
.AddSingleton<IScheduler, LocalScheduler>()
.AddSingleton<CronosCronParser>()
.AddSingleton(CronParser)
.AddScoped<ITriggerScheduler, DefaultTriggerScheduler>()
.AddScoped<IBookmarkScheduler, DefaultBookmarkScheduler>()
.AddSingleton<IScheduler, LocalScheduler>()

View file

@ -6,22 +6,12 @@ namespace Elsa.Scheduling.Services;
/// <summary>
/// A default implementation of <see cref="IWorkflowScheduler"/> that uses the <see cref="LocalScheduler"/>.
/// </summary>
public class DefaultWorkflowScheduler : IWorkflowScheduler
public class DefaultWorkflowScheduler(IScheduler scheduler) : IWorkflowScheduler
{
private readonly IScheduler _scheduler;
/// <summary>
/// Initializes a new instance of the <see cref="DefaultWorkflowScheduler"/> class.
/// </summary>
public DefaultWorkflowScheduler(IScheduler scheduler)
{
_scheduler = scheduler;
}
/// <inheritdoc />
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);
}
/// <inheritdoc />
@ -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);
}
/// <inheritdoc />
@ -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);
}
/// <inheritdoc />
@ -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);
}
/// <inheritdoc />
@ -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);
}
/// <inheritdoc />
@ -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);
}
/// <inheritdoc />
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);
}
}

View file

@ -22,4 +22,9 @@
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -15,4 +15,8 @@
<PackageReference Include="FluentStorage" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -14,5 +14,9 @@
<ProjectReference Include="..\Elsa.Workflows.Management\Elsa.Workflows.Management.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>
<!--Overridden for vulnaribility reasons with dependencies referencing older versions.-->
<ItemGroup>
<PackageReference Include="System.Text.Json" VersionOverride="8.0.4" />
</ItemGroup>
</Project>

View file

@ -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<string, object>();
draft.CustomProperties = model.CustomProperties ?? new Dictionary<string, object>();
draft.PropertyBag = model.PropertyBag ?? new PropertyBag();
draft.Variables = variables;
draft.Inputs = inputs;
draft.Outputs = outputs;

View file

@ -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<string>? definitionIds, CancellationToken cancellationToken)
private async Task<RefreshWorkflowDefinitionsResponse> RefreshWorkflowDefinitionsAsync(ICollection<string>? definitionIds, CancellationToken cancellationToken)
{
var request = new RefreshWorkflowDefinitionsRequest(definitionIds, BatchSize);
await workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(request, cancellationToken);
return await workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(request, cancellationToken);
}
}

View file

@ -1,6 +1,14 @@
using System.Text.Json.Serialization;
namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Refresh;
public class Request
internal class Request
{
public ICollection<string>? DefinitionIds { get; set; }
}
internal class Response(ICollection<string> refreshed, ICollection<string> notFound)
{
[JsonPropertyName("refreshed")] public ICollection<string> Refreshed { get; } = refreshed;
[JsonPropertyName("notFound")] public ICollection<string> NotFound { get; } = notFound;
}

View file

@ -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);
}
}

View file

@ -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<Request, Response>
{
/// <inheritdoc />
public override void Configure()
{
Get("/workflow-instances/{id}/execution-state");
ConfigurePermissions("read:workflow-instances");
}
/// <inheritdoc />
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);
}
}

View file

@ -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);

View file

@ -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);

View file

@ -8,9 +8,9 @@ public interface IWorkflowDefinitionActivityRegistryUpdater
/// <summary>
/// Tries to add a workflow as an activity to the registry.
/// </summary>
/// <param name="workflowDefinitionId">The ID of the workflow definition.</param>
/// <param name="workflowDefinitionVersionId">The version ID of the workflow definition.</param>
/// <param name="cancellationToken">The cancellation token.</param>
Task AddToRegistry(string workflowDefinitionId, CancellationToken cancellationToken = default);
Task AddToRegistry(string workflowDefinitionVersionId, CancellationToken cancellationToken = default);
/// <summary>
/// Removes workflow definition activities from the <see cref="Elsa.Workflows.Contracts.IActivityRegistry"/>.
@ -18,7 +18,6 @@ public interface IWorkflowDefinitionActivityRegistryUpdater
/// <param name="workflowDefinitionId">The ID of the workflow definition to remove.</param>
void RemoveDefinitionFromRegistry(string workflowDefinitionId);
/// <summary>
/// Removes a workflow definition version activity from the <see cref="Elsa.Workflows.Contracts.IActivityRegistry"/>.
/// </summary>

View file

@ -101,7 +101,6 @@ public class WorkflowManagementFeature : FeatureBase
/// <summary>
/// Adds all types implementing <see cref="IActivity"/> to the system.
/// </summary>
[RequiresUnreferencedCode("The assembly containing the specified marker type will be scanned for activity types.")]
public WorkflowManagementFeature AddActivitiesFrom<TMarker>()
{
var activityTypes = typeof(TMarker).Assembly.GetExportedTypes()

View file

@ -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;
/// <summary>
/// Deletes workflow instances when a workflow definition or version is deleted.
/// </summary>
[UsedImplicitly]
public class DeleteWorkflowInstances :
INotificationHandler<WorkflowDefinitionDeleting>,
INotificationHandler<WorkflowDefinitionVersionDeleting>,

View file

@ -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<string, object>();
draft.CustomProperties = model.CustomProperties ?? new Dictionary<string, object>();
draft.PropertyBag = model.PropertyBag ?? new PropertyBag();
draft.Variables = variables;
draft.Inputs = model.Inputs ?? new List<InputDefinition>();
draft.Outputs = model.Outputs ?? new List<OutputDefinition>();

View file

@ -11,27 +11,27 @@ public interface IWorkflowDefinitionStorePopulator
/// Populates the <see cref="IWorkflowDefinitionStore"/> with workflow definitions provided from <see cref="IWorkflowsProvider"/> implementations.
/// </summary>
/// <param name="cancellationToken">The cancellation token.</param>
Task PopulateStoreAsync(CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(CancellationToken cancellationToken = default);
/// <summary>
/// Populates the <see cref="IWorkflowDefinitionStore"/> with workflow definitions provided from <see cref="IWorkflowProvider"/> implementations.
/// </summary>
/// <param name="indexTriggers">Whether to index triggers.</param>
/// <param name="cancellationToken">The cancellation token.</param>
Task PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default);
/// <summary>
/// Adds a workflow definition to the store.
/// </summary>
/// <param name="materializedWorkflow">A materialized workflow.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default);
Task<WorkflowDefinition> AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default);
/// <summary>
/// Adds a workflow definition to the store.
/// </summary>
/// <param name="materializedWorkflow">A materialized workflow.</param>
/// /// <param name="indexTriggers">Whether to index triggers.</param>
/// <param name="indexTriggers">Whether to index triggers.</param>
/// <param name="cancellationToken">An optional cancellation token.</param>
Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default);
Task<WorkflowDefinition> AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default);
}

View file

@ -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);
}

View file

@ -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<WorkflowExecutionLogRecord> ExtractWorkflowExecutionLogs(WorkflowExecutionContext context);
}

View file

@ -26,6 +26,7 @@ public class CachingWorkflowRuntimeFeature : FeatureBase
.Decorate<ITriggerStore, CachingTriggerStore>()
// Handlers.
.AddNotificationHandler<InvalidateTriggersCache>();
.AddNotificationHandler<InvalidateTriggersCache>()
.AddNotificationHandler<InvalidateWorkflowsCache>();
}
}

View file

@ -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<IWorkflowDefinitionStorePopulator, DefaultWorkflowDefinitionStorePopulator>()
.AddScoped<IRegistriesPopulator, DefaultRegistriesPopulator>()
.AddScoped<IWorkflowDefinitionsRefresher, WorkflowDefinitionsRefresher>()
.AddScoped<IWorkflowDefinitionsReloader, WorkflowDefinitionsReloader>()
.AddScoped<IWorkflowRegistry, DefaultWorkflowRegistry>()
.AddScoped<IWorkflowMatcher, WorkflowMatcher>()
.AddScoped<IWorkflowInvoker, WorkflowInvoker>()
@ -207,6 +207,7 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddScoped<IBookmarkQueue, StoreBookmarkQueue>()
.AddScoped<IWorkflowCanceler, WorkflowCanceler>()
.AddScoped<IWorkflowCancellationService, WorkflowCancellationService>()
.AddScoped<IWorkflowExecutionLogRecordExtractor, WorkflowExecutionLogRecordExtractor>()
.AddScoped<IBookmarkQueueProcessor, BookmarkQueueProcessor>()
.AddScoped<StoreCommitStateHandler>()
@ -250,10 +251,10 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddNotificationHandler<CancelBackgroundActivities>()
.AddNotificationHandler<DeleteBookmarks>()
.AddNotificationHandler<DeleteTriggers>()
.AddNotificationHandler<DeleteWorkflowInstances>()
.AddNotificationHandler<DeleteActivityExecutionLogRecords>()
.AddNotificationHandler<DeleteWorkflowExecutionLogRecords>()
.AddNotificationHandler<WorkflowExecutionContextNotificationsHandler>()
.AddNotificationHandler<RefreshActivityRegistry>()
.AddNotificationHandler<SignalBookmarkQueueWorker>()
// Workflow activation strategies.

View file

@ -15,11 +15,19 @@ namespace Elsa.Workflows.Runtime.Handlers;
/// The cache is invalidated by calling the <c>TriggerTokenAsync</c> method of the <c>ICacheManager</c> passed to the class constructor.
/// </remarks>
[UsedImplicitly]
public class InvalidateTriggersCache(ICacheManager cacheManager) : INotificationHandler<WorkflowDefinitionsRefreshed>
public class InvalidateTriggersCache(ICacheManager cacheManager) :
INotificationHandler<WorkflowDefinitionsRefreshed>,
INotificationHandler<WorkflowDefinitionsReloaded>
{
/// <inheritdoc />
public Task HandleAsync(WorkflowDefinitionsRefreshed notification, CancellationToken cancellationToken)
{
return cacheManager.TriggerTokenAsync(CachingTriggerStore.CacheInvalidationTokenKey, cancellationToken).AsTask();
}
/// <inheritdoc />
public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken)
{
await cacheManager.TriggerTokenAsync(CachingTriggerStore.CacheInvalidationTokenKey, cancellationToken);
}
}

View file

@ -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;
/// <summary>
/// A notification handler that invalidates the workflow cache when workflow definitions are reloaded.
/// </summary>
/// <remarks>
/// The class implements the <c>INotificationHandler</c> interface and is responsible for handling <c>WorkflowDefinitionsReloaded</c> notifications.
/// When a <c>WorkflowDefinitionsReloaded</c> notification is received, the <c>HandleAsync</c> method is called to invalidate the http definition cache.
/// </remarks>
[UsedImplicitly]
public class InvalidateWorkflowsCache(IWorkflowDefinitionCacheManager workflowDefinitionCacheManager) : INotificationHandler<WorkflowDefinitionsReloaded>
{
/// <inheritdoc />
public async Task HandleAsync(WorkflowDefinitionsReloaded notification, CancellationToken cancellationToken)
{
foreach (var reloadedWorkflowDefinition in notification.ReloadedWorkflowDefinitions)
await workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(reloadedWorkflowDefinition.DefinitionId, cancellationToken);
}
}

View file

@ -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 <see cref="IActivityRegistry"/> for the <see cref="WorkflowDefinitionActivityProvider"/> provider whenever workflow definitions are reloaded.
[PublicAPI]
public class RefreshActivityRegistry(IWorkflowDefinitionActivityRegistryUpdater workflowDefinitionActivityRegistryUpdater) : INotificationHandler<WorkflowDefinitionsReloaded>
{
/// <inheritdoc />
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;
}
}

View file

@ -0,0 +1,30 @@
using Elsa.Workflows.Management.Entities;
namespace Elsa.Workflows.Runtime.Models;
/// <summary>
/// Represents a reloaded workflow definition with necessary properties for
/// identification, versioning, and usability status as an activity.
/// </summary>
/// <param name="DefinitionId">The unique identifier for the workflow definition.</param>
/// <param name="DefinitionVersionId">The unique identifier for the specific version of the workflow definition.</param>
/// <param name="Version">The version number of the workflow definition.</param>
/// <param name="UsableAsActivity">Indicates whether the workflow definition can be used as an activity.</param>
public record ReloadedWorkflowDefinition(string DefinitionId, string DefinitionVersionId, int Version, bool UsableAsActivity)
{
/// <summary>
/// Creates an instance of <see cref="ReloadedWorkflowDefinition"/> from a given <see cref="WorkflowDefinition"/>.
/// </summary>
/// <param name="workflowDefinition">The workflow definition used to create the reloaded workflow definition.</param>
/// <returns>A new instance of <see cref="ReloadedWorkflowDefinition"/>.</returns>
public static ReloadedWorkflowDefinition FromDefinition(WorkflowDefinition workflowDefinition)
{
return new ReloadedWorkflowDefinition
(
workflowDefinition.DefinitionId,
workflowDefinition.Id,
workflowDefinition.Version,
workflowDefinition.Options.UsableAsActivity ?? false
);
}
}

View file

@ -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<ReloadedWorkflowDefinition> ReloadedWorkflowDefinitions) : INotification;

View file

@ -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<WorkflowDefinition> WorkflowDefinitions);
public record RefreshWorkflowDefinitionsResponse(ICollection<string> Refreshed, ICollection<string> NotFound);

View file

@ -47,37 +47,47 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP
}
/// <inheritdoc />
public Task PopulateStoreAsync(CancellationToken cancellationToken = default)
public Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(CancellationToken cancellationToken = default)
{
return PopulateStoreAsync(true, cancellationToken);
}
/// <inheritdoc />
public async Task PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default)
public async Task<IEnumerable<WorkflowDefinition>> PopulateStoreAsync(bool indexTriggers, CancellationToken cancellationToken = default)
{
var providers = _workflowDefinitionProviders();
var workflowDefinitions = new List<WorkflowDefinition>();
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;
}
/// <inheritdoc />
public Task AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default)
public Task<WorkflowDefinition> AddAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken = default)
{
return AddAsync(materializedWorkflow, true, cancellationToken);
}
/// <inheritdoc />
public async Task AddAsync(MaterializedWorkflow materializedWorkflow, bool indexTriggers, CancellationToken cancellationToken = default)
public async Task<WorkflowDefinition> 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<WorkflowDefinition> 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<WorkflowDefinition> 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()
{

View file

@ -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;
/// <summary>
/// This implementation saves <see cref="WorkflowExecutionLogRecord"/> directly through the store.
/// </summary>
public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IIdentityGenerator identityGenerator, INotificationSender notificationSender) : IWorkflowExecutionLogSink
public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, IWorkflowExecutionLogRecordExtractor extractor, INotificationSender notificationSender) : IWorkflowExecutionLogSink
{
/// <inheritdoc />
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);
}
}

View file

@ -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<WorkflowDefinition> definitions, CancellationToken cancellationToken)

View file

@ -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;
/// <inheritdoc />
public class WorkflowDefinitionsReloader(IWorkflowDefinitionStorePopulator workflowDefinitionStorePopulator, INotificationSender notificationSender) : IWorkflowDefinitionsReloader
{
/// <inheritdoc />
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);
}
}

View file

@ -0,0 +1,35 @@
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Services;
/// <inheritdoc />
public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGenerator) : IWorkflowExecutionLogRecordExtractor
{
/// <inheritdoc />
public IEnumerable<WorkflowExecutionLogRecord> 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
});
}
}

View file

@ -75,6 +75,12 @@
<None Update="Scenarios\WorkflowDefinitionRefresh\dynamic-endpoint-http-workflow.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\WorkflowDefinitionReload\http-workflow.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\WorkflowDefinitionReload\Workflows\http-workflow.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
</ItemGroup>
</Project>

View file

@ -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<ISignalManager, SignalManager>();
services.AddSingleton<IWorkflowEvents, WorkflowEvents>();
services.AddSingleton<IWorkflowDefinitionEvents, WorkflowDefinitionEvents>();
services.AddSingleton<ITriggerChangeTokenSignalEvents, TriggerChangeTokenSignalEvents>();
services.AddScoped<ITenantResolutionStrategy, TestTenantResolutionStrategy>();
services.AddNotificationHandlersFrom<WorkflowServer>();
services.AddNotificationHandlersFrom<WorkflowEventHandlers>();
services
.AddSingleton<ISignalManager, SignalManager>()
.AddSingleton<IWorkflowEvents, WorkflowEvents>()
.AddSingleton<IWorkflowDefinitionEvents, WorkflowDefinitionEvents>()
.AddSingleton<ITriggerChangeTokenSignalEvents, TriggerChangeTokenSignalEvents>()
.AddScoped<IWorkflowMaterializer, TestWorkflowMaterializer>()
.AddNotificationHandlersFrom<WorkflowServer>()
.AddWorkflowDefinitionProvider<TestWorkflowProvider>()
.AddNotificationHandlersFrom<WorkflowEventHandlers>()
;
});
}

View file

@ -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 <see cref="TestWorkflowProvider"/>.
public class TestWorkflowMaterializer(IEnumerable<IWorkflowProvider> workflowProviders) : IWorkflowMaterializer
{
/// The name of the materializer.
public const string MaterializerName = "Test";
/// <inheritdoc />
public string Name => MaterializerName;
/// <inheritdoc />
public ValueTask<Workflow> 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);
}
}

View file

@ -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<MaterializedWorkflow> MaterializedWorkflows { get; set; } = new List<MaterializedWorkflow>();
public ValueTask<IEnumerable<MaterializedWorkflow>> GetWorkflowsAsync(CancellationToken cancellationToken = default)
{
return new(MaterializedWorkflows);
}
}

View file

@ -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);
}
}

View file

@ -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<IWorkflowDefinitionManager>();
_workflowDefinitionsReloader = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionsReloader>();
_workflowBuilderFactory = Scope.ServiceProvider.GetRequiredService<IWorkflowBuilderFactory>();
_workflowDefinitionService = Scope.ServiceProvider.GetRequiredService<IWorkflowDefinitionService>();
_activityRegistry = Scope.ServiceProvider.GetRequiredService<IActivityRegistry>();
var workflowProviders = Scope.ServiceProvider.GetRequiredService<IEnumerable<IWorkflowProvider>>();
_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<MaterializedWorkflow> 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);
}
}

View file

@ -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": []
}
}