From 319cf89a44ca24ef976da786fa78a5e70a783a66 Mon Sep 17 00:00:00 2001 From: Matt Date: Wed, 4 Jun 2025 06:08:59 +0100 Subject: [PATCH 1/7] =?UTF-8?q?Update=20SqlEvaluator=20expression=20keywor?= =?UTF-8?q?d=20resolver=20for=20using=20variables.=20=E2=80=A6=20(#6705)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Update SqlEvaluator expression keyword resolver for using variables. Changes from 'Variables.' to 'Variable.' to match other properties singular naming convention. * Fix substring index for variable retrieval. --- src/modules/Elsa.Sql/Services/SqlEvaluator.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/modules/Elsa.Sql/Services/SqlEvaluator.cs b/src/modules/Elsa.Sql/Services/SqlEvaluator.cs index 10ee5a0df..3f0e4d724 100644 --- a/src/modules/Elsa.Sql/Services/SqlEvaluator.cs +++ b/src/modules/Elsa.Sql/Services/SqlEvaluator.cs @@ -85,7 +85,7 @@ public class SqlEvaluator() : ISqlEvaluator "LastResult" => expressionContext.GetLastResult(), var i when i.StartsWith("Input.") => executionContext.Input.TryGetValue(i.Substring(6), out var v) ? v : null, var o when o.StartsWith("Output.") => executionContext.Output.TryGetValue(o.Substring(7), out var v) ? v : null, - var v when v.StartsWith("Variables.") => executionContext.Variables.FirstOrDefault(x => x.Name == v.Substring(10), null)?.Value ?? null, + var v when v.StartsWith("Variable.") => expressionContext.GetVariableInScope(v.Substring(9)) ?? null, _ => throw new NullReferenceException($"No matching property found for {{{{{key}}}}}.") }; } From c817320d66de59609c8a43d8b362c0ea9e5c30ff Mon Sep 17 00:00:00 2001 From: Matt Date: Wed, 4 Jun 2025 19:57:31 +0100 Subject: [PATCH 2/7] Adds back support for 'Variables.' within SqlEvaluator with an obsolete comment. (#6710) --- src/modules/Elsa.Sql/Services/SqlEvaluator.cs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/modules/Elsa.Sql/Services/SqlEvaluator.cs b/src/modules/Elsa.Sql/Services/SqlEvaluator.cs index 3f0e4d724..ed1fc46a2 100644 --- a/src/modules/Elsa.Sql/Services/SqlEvaluator.cs +++ b/src/modules/Elsa.Sql/Services/SqlEvaluator.cs @@ -86,6 +86,8 @@ public class SqlEvaluator() : ISqlEvaluator var i when i.StartsWith("Input.") => executionContext.Input.TryGetValue(i.Substring(6), out var v) ? v : null, var o when o.StartsWith("Output.") => executionContext.Output.TryGetValue(o.Substring(7), out var v) ? v : null, var v when v.StartsWith("Variable.") => expressionContext.GetVariableInScope(v.Substring(9)) ?? null, + // OBSOLETE: This is deprecated and will be removed in a future version. Use 'Variable.' instead. + var v when v.StartsWith("Variables.") => expressionContext.GetVariableInScope(v.Substring(10)) ?? null, _ => throw new NullReferenceException($"No matching property found for {{{{{key}}}}}.") }; } From 33fecf4a0e53e60e84b1eaff33f4f94b727c87d0 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 6 Jun 2025 18:51:08 +0200 Subject: [PATCH 3/7] Update Docker image tags to v3-5-0-preview in GitHub workflows --- .github/workflows/elsa-server-and-studio.yml | 2 +- .github/workflows/elsa-server.yml | 2 +- .github/workflows/elsa-studio.yml | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/elsa-server-and-studio.yml b/.github/workflows/elsa-server-and-studio.yml index 9c23cc7c9..ba0db7875 100644 --- a/.github/workflows/elsa-server-and-studio.yml +++ b/.github/workflows/elsa-server-and-studio.yml @@ -29,7 +29,7 @@ jobs: with: # list of Docker images to use as base name for tags images: | - elsaworkflows/elsa-server-and-studio-v3-4-0-preview + elsaworkflows/elsa-server-and-studio-v3-5-0-preview flavor: | latest=true # generate Docker tags based on the following events/attributes diff --git a/.github/workflows/elsa-server.yml b/.github/workflows/elsa-server.yml index 289332ba4..6dfe94700 100644 --- a/.github/workflows/elsa-server.yml +++ b/.github/workflows/elsa-server.yml @@ -29,7 +29,7 @@ jobs: with: # list of Docker images to use as base name for tags images: | - elsaworkflows/elsa-server-v3-4-0-preview + elsaworkflows/elsa-server-v3-5-0-preview flavor: | latest=true # generate Docker tags based on the following events/attributes diff --git a/.github/workflows/elsa-studio.yml b/.github/workflows/elsa-studio.yml index f2fe2bff2..d3c8ac075 100644 --- a/.github/workflows/elsa-studio.yml +++ b/.github/workflows/elsa-studio.yml @@ -29,7 +29,7 @@ jobs: with: # list of Docker images to use as base name for tags images: | - elsaworkflows/elsa-studio-v3-4-0-preview + elsaworkflows/elsa-studio-v3-5-0-preview flavor: | latest=true # generate Docker tags based on the following events/attributes From 6eccb3540ac07d1710d5bd28e476f8b16bae61c1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 11 Jun 2025 09:48:00 +0200 Subject: [PATCH 4/7] Update package versions (#6722) * Update package versions Updated multiple package versions to their latest release, including `Polly`, `Elastic.Clients.Elasticsearch`, `System.Text.Json`, and others. * Remove duplicate package entry * Fix case mismatch in package name for SQL Server migration dependency. * Add `Proto.Cluster.AzureContainerApps` package version to dependencies * Update Directory.Packages.props Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- Directory.Packages.props | 111 +++++++++++++++++++-------------------- 1 file changed, 55 insertions(+), 56 deletions(-) diff --git a/Directory.Packages.props b/Directory.Packages.props index fd43a70aa..8579d1f11 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -8,8 +8,9 @@ - + + @@ -17,7 +18,7 @@ - + @@ -30,7 +31,7 @@ - + @@ -45,7 +46,7 @@ - + @@ -63,10 +64,10 @@ - + - + @@ -79,17 +80,18 @@ - + - - + + + @@ -104,68 +106,65 @@ - + - + - - - - - - + + + + + + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - - + + + + \ No newline at end of file From 623b75ed2c29682492ab2c29d56d07d19988d323 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 12 Jun 2025 09:25:43 +0200 Subject: [PATCH 5/7] Fix order of Order and Pagination (#6727) * Refactor query composition to ensure consistent ordering and pagination logic. Reordered method calls for `OrderBy` and `Paginate` across multiple stores to enhance readability and maintain consistent execution. Simplified redundant query operations for improved clarity and performance. * Updates Elsa Studio version to 3.4.0 Updates the Elsa Studio version to the stable release. Removes the preview tag from the version number. * Move `ElsaStudioVersion` property to `Directory.Packages.props` for centralized management. --- Directory.Build.props | 3 --- Directory.Packages.props | 3 +++ .../Modules/Management/WorkflowDefinitionStore.cs | 4 ++-- .../Modules/Runtime/BookmarkQueueStore.cs | 6 +++--- .../Modules/Runtime/WorkflowExecutionLogStore.cs | 5 +++-- .../Elsa.Secrets.Management/Services/MemorySecretStore.cs | 4 ++-- .../EFCoreSecretStore.cs | 2 +- 7 files changed, 14 insertions(+), 13 deletions(-) diff --git a/Directory.Build.props b/Directory.Build.props index fddc11126..79b1b1640 100644 --- a/Directory.Build.props +++ b/Directory.Build.props @@ -36,7 +36,4 @@ $(NoWarn);IL2026;IL2046;IL2057;IL2067;IL2070;IL2072;IL2075;IL2087;IL2091 - - 3.4.0-preview.1025 - \ No newline at end of file diff --git a/Directory.Packages.props b/Directory.Packages.props index 21873199a..3758682dd 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -3,6 +3,9 @@ true true + + 3.4.0 + diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs index cc0a459be..f3b7c6996 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowDefinitionStore.cs @@ -45,8 +45,8 @@ public class EFCoreWorkflowDefinitionStore(EntityStore public async Task> FindManyAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, PageArgs pageArgs, CancellationToken cancellationToken = default) { - var count = await store.QueryAsync(queryable => Filter(queryable, filter).OrderBy(order), cancellationToken).LongCount(); - var results = await store.QueryAsync(queryable => Paginate(Filter(queryable, filter), pageArgs), OnLoadAsync, filter.TenantAgnostic, cancellationToken).ToList(); + var count = await store.QueryAsync(queryable => Filter(queryable, filter), cancellationToken).LongCount(); + var results = await store.QueryAsync(queryable => Filter(queryable, filter).OrderBy(order).Paginate(pageArgs), OnLoadAsync, filter.TenantAgnostic, cancellationToken).ToList(); return new(results, count); } diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs index 0e7f37812..be232b47b 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs @@ -43,14 +43,14 @@ public class EFBookmarkQueueStore(Store public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) { - var count = await store.QueryAsync(queryable => queryable.OrderBy(orderBy), cancellationToken).LongCount(); - var results = await store.QueryAsync(queryable => queryable.Paginate(pageArgs), OnLoadAsync, cancellationToken).ToList(); + var count = await store.QueryAsync(queryable => queryable, cancellationToken).LongCount(); + var results = await store.QueryAsync(queryable => queryable.OrderBy(orderBy).Paginate(pageArgs), OnLoadAsync, cancellationToken).ToList(); return new(results, count); } public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueFilter filter, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) { - var count = await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(orderBy), cancellationToken).LongCount(); + var count = await store.QueryAsync(filter.Apply, cancellationToken).LongCount(); var results = await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(orderBy).Paginate(pageArgs), OnLoadAsync, cancellationToken).ToList(); return new(results, count); } diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs index 03b9088cb..2569f479e 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs @@ -61,7 +61,7 @@ public class EFCoreWorkflowExecutionLogStore(EntityStore> FindManyAsync(WorkflowExecutionLogRecordFilter filter, PageArgs pageArgs, WorkflowExecutionLogRecordOrder order, CancellationToken cancellationToken = default) { var count = await store.QueryAsync(queryable => Filter(queryable, filter), cancellationToken).LongCount(); - var results = await store.QueryAsync(queryable => Filter(queryable, filter).Paginate(pageArgs).OrderBy(order), OnLoadAsync, cancellationToken).ToList(); + var results = await store.QueryAsync(queryable => Filter(queryable, filter).OrderBy(order).Paginate(pageArgs), OnLoadAsync, cancellationToken).ToList(); return new(results, count); } @@ -71,10 +71,11 @@ public class EFCoreWorkflowExecutionLogStore(EntityStore Filter(queryable, filter), cancellationToken); } - private async ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, WorkflowExecutionLogRecord entity, CancellationToken cancellationToken) + private ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, WorkflowExecutionLogRecord entity, CancellationToken cancellationToken) { entity = entity.SanitizeLogMessage(); dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = ShouldSerializePayload(entity) ? safeSerializer.Serialize(entity.Payload) : null; + return ValueTask.CompletedTask; } private async ValueTask OnLoadAsync(RuntimeElsaDbContext dbContext, WorkflowExecutionLogRecord? entity, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Secrets.Management/Services/MemorySecretStore.cs b/src/modules/Elsa.Secrets.Management/Services/MemorySecretStore.cs index 30440998f..d30b3d692 100644 --- a/src/modules/Elsa.Secrets.Management/Services/MemorySecretStore.cs +++ b/src/modules/Elsa.Secrets.Management/Services/MemorySecretStore.cs @@ -10,8 +10,8 @@ public class MemorySecretStore(MemoryStore memoryStore) : ISecretStore { public Task> FindManyAsync(SecretFilter filter, SecretOrder order, PageArgs pageArgs, CancellationToken cancellationToken = default) { - var count = memoryStore.Query(query => Filter(query, filter).OrderBy(order)).LongCount(); - var result = memoryStore.Query(query => Filter(query, filter).Paginate(pageArgs)).ToList(); + var count = memoryStore.Query(query => Filter(query, filter)).LongCount(); + var result = memoryStore.Query(query => Filter(query, filter).OrderBy(order).Paginate(pageArgs)).ToList(); return Task.FromResult(Page.Of(result, count)); } diff --git a/src/modules/Elsa.Secrets.Persistence.EntityFrameworkCore/EFCoreSecretStore.cs b/src/modules/Elsa.Secrets.Persistence.EntityFrameworkCore/EFCoreSecretStore.cs index 8dcebdc8c..87ae79004 100644 --- a/src/modules/Elsa.Secrets.Persistence.EntityFrameworkCore/EFCoreSecretStore.cs +++ b/src/modules/Elsa.Secrets.Persistence.EntityFrameworkCore/EFCoreSecretStore.cs @@ -16,7 +16,7 @@ public class EFCoreSecretStore(EntityStore store) : IS public async Task> FindManyAsync(SecretFilter filter, SecretOrder order, PageArgs pageArgs, CancellationToken cancellationToken = default) { var count = await store.QueryAsync(query => Filter(query, filter), cancellationToken).LongCount(); - var secrets = await store.QueryAsync(query => Filter(query, filter).Paginate(pageArgs).OrderBy(order), cancellationToken).ToList(); + var secrets = await store.QueryAsync(query => Filter(query, filter).OrderBy(order).Paginate(pageArgs), cancellationToken).ToList(); return new(secrets, count); } From 935d6987d552697331da8ba1c678d8d89756034c Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 13 Jun 2025 08:30:02 +0200 Subject: [PATCH 6/7] Commits workflow state during alteration (#6736) * Reduce default logging verbosity in appsettings.json * Commit workflow state during alteration execution Added `ICommitStateHandler` dependency and implemented workflow state commitment in `DefaultAlterationRunner` to ensure state persistence during alteration execution. --- src/apps/Elsa.Server.Web/appsettings.json | 6 ++---- .../Elsa.Alterations/Services/DefaultAlterationRunner.cs | 7 +++++++ 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/src/apps/Elsa.Server.Web/appsettings.json b/src/apps/Elsa.Server.Web/appsettings.json index 2f6eccad2..39aa50c97 100644 --- a/src/apps/Elsa.Server.Web/appsettings.json +++ b/src/apps/Elsa.Server.Web/appsettings.json @@ -1,10 +1,8 @@ { "Logging": { "LogLevel": { - "Default": "Information", - "Microsoft.Hosting.Lifetime": "Information", - "Elsa": "Warning", - "Elsa.Workflows.Runtime.Middleware.Workflows.WorkflowHeartbeatMiddleware": "Warning" + "Default": "Warning", + "Microsoft.Hosting.Lifetime": "Information" } }, "HostBuilder": { diff --git a/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs b/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs index 71317a188..37ae0e87f 100644 --- a/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs +++ b/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs @@ -4,6 +4,7 @@ using Elsa.Alterations.Core.Results; using Elsa.Alterations.Middleware.Workflows; using Elsa.Common; using Elsa.Workflows; +using Elsa.Workflows.CommitStates; using Elsa.Workflows.Management; using Elsa.Workflows.Pipelines.WorkflowExecution; using Elsa.Workflows.Runtime; @@ -17,6 +18,7 @@ public class DefaultAlterationRunner( IWorkflowExecutionPipeline workflowExecutionPipeline, IWorkflowDefinitionService workflowDefinitionService, IWorkflowStateExtractor workflowStateExtractor, + ICommitStateHandler commitStateHandler, ISystemClock systemClock, IServiceProvider serviceProvider) : IAlterationRunner @@ -81,8 +83,13 @@ public class DefaultAlterationRunner( // Extract workflow state. workflowState = workflowStateExtractor.Extract(workflowExecutionContext); + + // Commit workflow state. + await commitStateHandler.CommitAsync(workflowExecutionContext, workflowState, cancellationToken); // Apply updated workflow state. + // TODO: Importing back into the workflow runtime makes sense, but this also causes another SAVE ction of the workflow instance in the DB, which also happens in the previous step during the commit action. + // Can we avoid this? Perhaps we need more granular control over when we purge and when we save to DB. await workflowClient.ImportStateAsync(workflowState, cancellationToken); // Check if the workflow has scheduled work. From c694a18c13bb252b8cc37bc1b829bdb882a57dfb Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 13 Jun 2025 14:26:51 +0200 Subject: [PATCH 7/7] Enhances Mediator with Tenant Context Propagation (#6738) * Update package versions in Directory.Packages.props Upgraded multiple package dependencies to latest versions, ensuring compatibility, security, and access to the newest features. * Refactor mediator pipeline to support tenant context propagation - Introduced `TenantPropagatingMiddleware` to handle tenant context propagation during command execution. - Added `SetupMediatorPipelines` hosted service for configuring mediator pipelines. - Enhanced `CommandPipeline` and builder to allow middleware insertion, removal, and reordering. - Updated `CommandContext` and related components to support headers for tenant context handling. - Improved logging and refactored `BackgroundWorkflowDispatcher` to include tenant headers during command dispatch. * Fix typos in XML documentation and improve middleware extension clarity - Corrected duplicated slashes in XML doc comments in `ICommandSender.cs`. - Refined phrasing in `MiddlewareExtensions.cs` to clarify method parameters and improve readability. * Update src/common/Elsa.Mediator/Middleware/Command/Components/CommandLoggingMiddleware.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- Directory.Packages.props | 74 +++++++++---------- src/apps/Elsa.Server.Web/Program.cs | 4 +- src/apps/Elsa.Server.Web/appsettings.json | 2 +- .../Elsa.Mediator/Channels/CommandsChannel.cs | 3 +- .../CommandStrategies/BackgroundStrategy.cs | 2 +- .../CommandStrategies/DefaultStrategy.cs | 3 +- .../Contexts/CommandStrategyContext.cs | 5 +- .../Elsa.Mediator/Contracts/ICommandSender.cs | 32 +++++++- .../Contracts/ICommandsChannel.cs | 5 +- .../DependencyInjectionExtensions.cs | 8 +- .../BackgroundCommandSenderHostedService.cs | 21 +++--- .../Middleware/Command/CommandContext.cs | 14 +++- .../Middleware/Command/CommandPipeline.cs | 25 ++++--- .../Command/CommandPipelineBuilder.cs | 36 +++++++-- .../CommandHandlerInvokerMiddleware.cs | 31 +++----- .../Components/CommandLoggingMiddleware.cs | 24 ++---- .../Command/Contracts/ICommandPipeline.cs | 2 +- .../Contracts/ICommandPipelineBuilder.cs | 24 +++++- .../Command/MiddlewareExtensions.cs | 37 +++++++--- .../NotificationHandlerInvokerMiddleware.cs | 36 +++------ .../Notification/MiddlewareExtensions.cs | 2 +- .../Notification/NotificationContext.cs | 11 ++- .../Middleware/Request/RequestContext.cs | 30 +++----- .../Elsa.Mediator/Services/DefaultMediator.cs | 69 +++++++++++------ .../Elsa.Tenants/Features/TenantsFeature.cs | 6 ++ .../Middleware/TenantPropagatingMiddleware.cs | 28 +++++++ .../Mediator/Tasks/SetupMediatorPipelines.cs | 19 +++++ .../Elsa.Tenants/Mediator/TenantHeaders.cs | 13 ++++ .../Services/BackgroundWorkflowDispatcher.cs | 29 ++++---- 29 files changed, 375 insertions(+), 220 deletions(-) create mode 100644 src/modules/Elsa.Tenants/Mediator/Middleware/TenantPropagatingMiddleware.cs create mode 100644 src/modules/Elsa.Tenants/Mediator/Tasks/SetupMediatorPipelines.cs create mode 100644 src/modules/Elsa.Tenants/Mediator/TenantHeaders.cs diff --git a/Directory.Packages.props b/Directory.Packages.props index 3758682dd..0f9791acc 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -126,47 +126,47 @@ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + - - + + - - - - + + + + \ No newline at end of file diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index 1f2617642..aee6fb0c9 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -89,7 +89,7 @@ const PersistenceProvider persistenceProvider = PersistenceProvider.EntityFramew const bool useDbContextPooling = false; const bool useHangfire = false; const bool useQuartz = true; -const bool useMassTransit = true; +const bool useMassTransit = false; const bool useZipCompression = false; const bool runEFCoreMigrations = true; const bool useMemoryStores = false; @@ -98,7 +98,7 @@ const bool useKafka = false; const bool useReadOnlyMode = false; const bool useSignalR = false; // Disabled until Elsa Studio sends authenticated requests. const WorkflowRuntime workflowRuntime = WorkflowRuntime.Distributed; -const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit; +const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.Memory; const MassTransitBroker massTransitBroker = MassTransitBroker.Memory; const bool useMultitenancy = false; const bool useTenantsFromConfiguration = true; diff --git a/src/apps/Elsa.Server.Web/appsettings.json b/src/apps/Elsa.Server.Web/appsettings.json index 91735a81b..a617ec190 100644 --- a/src/apps/Elsa.Server.Web/appsettings.json +++ b/src/apps/Elsa.Server.Web/appsettings.json @@ -44,7 +44,7 @@ }, { "Id": "tenant-1", - "Name": "Tenant 12 dd", + "Name": "Tenant 1", "Configuration": { "Http": { "Prefix": "/tenant-1", diff --git a/src/common/Elsa.Mediator/Channels/CommandsChannel.cs b/src/common/Elsa.Mediator/Channels/CommandsChannel.cs index 43a8b98b0..f80bbdbaa 100644 --- a/src/common/Elsa.Mediator/Channels/CommandsChannel.cs +++ b/src/common/Elsa.Mediator/Channels/CommandsChannel.cs @@ -1,9 +1,10 @@ using Elsa.Mediator.Abstractions; using Elsa.Mediator.Contracts; +using Elsa.Mediator.Middleware.Command; namespace Elsa.Mediator.Channels; /// -public class CommandsChannel : ChannelBase, ICommandsChannel +public class CommandsChannel : ChannelBase, ICommandsChannel { } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/CommandStrategies/BackgroundStrategy.cs b/src/common/Elsa.Mediator/CommandStrategies/BackgroundStrategy.cs index e99a0a1b0..8d4210e05 100644 --- a/src/common/Elsa.Mediator/CommandStrategies/BackgroundStrategy.cs +++ b/src/common/Elsa.Mediator/CommandStrategies/BackgroundStrategy.cs @@ -13,7 +13,7 @@ public class BackgroundStrategy : ICommandStrategy public async Task ExecuteAsync(CommandStrategyContext context) { var commandsChannel = context.ServiceProvider.GetRequiredService(); - await commandsChannel.Writer.WriteAsync(context.Command, context.CancellationToken); + await commandsChannel.Writer.WriteAsync(context.CommandContext, context.CancellationToken); return default!; } } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/CommandStrategies/DefaultStrategy.cs b/src/common/Elsa.Mediator/CommandStrategies/DefaultStrategy.cs index 9b204332f..a1e26aa79 100644 --- a/src/common/Elsa.Mediator/CommandStrategies/DefaultStrategy.cs +++ b/src/common/Elsa.Mediator/CommandStrategies/DefaultStrategy.cs @@ -12,7 +12,8 @@ public class DefaultStrategy : ICommandStrategy /// public async Task ExecuteAsync(CommandStrategyContext context) { - var command = context.Command; + var commandContext = context.CommandContext; + var command = commandContext.Command; var cancellationToken = context.CancellationToken; var commandType = command.GetType(); var handleMethod = commandType.GetCommandHandlerMethod(); diff --git a/src/common/Elsa.Mediator/Contexts/CommandStrategyContext.cs b/src/common/Elsa.Mediator/Contexts/CommandStrategyContext.cs index d0e395063..9c4c9a640 100644 --- a/src/common/Elsa.Mediator/Contexts/CommandStrategyContext.cs +++ b/src/common/Elsa.Mediator/Contexts/CommandStrategyContext.cs @@ -1,12 +1,13 @@ using Elsa.Mediator.Contracts; +using Elsa.Mediator.Middleware.Command; namespace Elsa.Mediator.Contexts; /// /// Represents a context for executing a command. /// -/// The command to execute. +/// The command context to execute. /// The command handler. /// The service provider to resolve services from. /// The cancellation token. -public record CommandStrategyContext(ICommand Command, ICommandHandler Handler, IServiceProvider ServiceProvider, CancellationToken CancellationToken = default); \ No newline at end of file +public record CommandStrategyContext(CommandContext CommandContext, ICommandHandler Handler, IServiceProvider ServiceProvider, CancellationToken CancellationToken = default); \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Contracts/ICommandSender.cs b/src/common/Elsa.Mediator/Contracts/ICommandSender.cs index fb7a44b34..5373afb3b 100644 --- a/src/common/Elsa.Mediator/Contracts/ICommandSender.cs +++ b/src/common/Elsa.Mediator/Contracts/ICommandSender.cs @@ -6,13 +6,19 @@ namespace Elsa.Mediator.Contracts; public interface ICommandSender { /// - /// Sends a command using he default strategy. + /// Sends a command using the default strategy. + /// + Task SendAsync(ICommand command, CancellationToken cancellationToken = default); + + /// + /// Sends a command using the default strategy. /// /// The command to send. + /// Any headers to pass along. /// The cancellation token. /// The type of the result. /// The result. - Task SendAsync(ICommand command, CancellationToken cancellationToken = default); + Task SendAsync(ICommand command, IDictionary headers, CancellationToken cancellationToken = default); /// /// Sends a command using the specified strategy. @@ -24,6 +30,17 @@ public interface ICommandSender /// The result. Task SendAsync(ICommand command, ICommandStrategy strategy, CancellationToken cancellationToken = default); + /// + /// Sends a command using the specified strategy. + /// + /// The command to send. + /// The command strategy to use. + /// Any headers to pass along. + /// The cancellation token. + /// The type of the result. + /// The result. + Task SendAsync(ICommand command, ICommandStrategy strategy, IDictionary headers, CancellationToken cancellationToken = default); + /// /// Sends a command using the default strategy. /// @@ -37,5 +54,14 @@ public interface ICommandSender /// The command to send. /// The command strategy to use. /// The cancellation token. - Task SendAsync(ICommand command, ICommandStrategy? strategy, CancellationToken cancellationToken = default); + Task SendAsync(ICommand command, ICommandStrategy strategy, CancellationToken cancellationToken = default); + + /// + /// Sends a command using the specified strategy. + /// + /// The command to send. + /// The command strategy to use. + /// Any headers to pass along. + /// The cancellation token. + Task SendAsync(ICommand command, ICommandStrategy strategy, IDictionary headers, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Contracts/ICommandsChannel.cs b/src/common/Elsa.Mediator/Contracts/ICommandsChannel.cs index 1ce9b9b14..280df942d 100644 --- a/src/common/Elsa.Mediator/Contracts/ICommandsChannel.cs +++ b/src/common/Elsa.Mediator/Contracts/ICommandsChannel.cs @@ -1,4 +1,5 @@ using System.Threading.Channels; +using Elsa.Mediator.Middleware.Command; namespace Elsa.Mediator.Contracts; @@ -10,10 +11,10 @@ public interface ICommandsChannel /// /// Gets the writer for the commands queue. /// - ChannelWriter Writer { get; } + ChannelWriter Writer { get; } /// /// Gets the reader for the commands queue. /// - ChannelReader Reader { get; } + ChannelReader Reader { get; } } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs index 2ac94c073..68ea123d8 100644 --- a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs +++ b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs @@ -36,9 +36,9 @@ public static class DependencyInjectionExtensions .AddScoped(sp => sp.GetRequiredService()) .AddScoped(sp => sp.GetRequiredService()) .AddScoped(sp => sp.GetRequiredService()) - .AddScoped() - .AddScoped() - .AddScoped() + .AddSingleton() + .AddSingleton() + .AddSingleton() ; } @@ -53,7 +53,7 @@ public static class DependencyInjectionExtensions .AddSingleton() .AddSingleton() .AddSingleton() - .AddSingleton() + .AddSingleton() .AddHostedService() .AddHostedService() .AddHostedService(); diff --git a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs index 8b4a0db0c..47770bab5 100644 --- a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs @@ -1,5 +1,6 @@ using System.Threading.Channels; using Elsa.Mediator.Contracts; +using Elsa.Mediator.Middleware.Command; using Elsa.Mediator.Options; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; @@ -16,7 +17,7 @@ public class BackgroundCommandSenderHostedService : BackgroundService private readonly int _workerCount; private readonly ICommandsChannel _commandsChannel; private readonly IServiceScopeFactory _scopeFactory; - private readonly List> _outputs; + private readonly List> _outputs; private readonly ILogger _logger; /// @@ -26,7 +27,7 @@ public class BackgroundCommandSenderHostedService : BackgroundService _commandsChannel = commandsChannel; _scopeFactory = scopeFactory; _logger = logger; - _outputs = new List>(_workerCount); + _outputs = new(_workerCount); } /// @@ -36,34 +37,32 @@ public class BackgroundCommandSenderHostedService : BackgroundService for (var i = 0; i < _workerCount; i++) { - var output = Channel.CreateUnbounded(); + var output = Channel.CreateUnbounded(); _outputs.Add(output); _ = ReadOutputAsync(output, cancellationToken); } - await foreach (var command in _commandsChannel.Reader.ReadAllAsync(cancellationToken)) + await foreach (var commandContext in _commandsChannel.Reader.ReadAllAsync(cancellationToken)) { var output = _outputs[index]; - await output.Writer.WriteAsync(command, cancellationToken); + await output.Writer.WriteAsync(commandContext, cancellationToken); index = (index + 1) % _workerCount; } - foreach (var output in _outputs) - { + foreach (var output in _outputs) output.Writer.Complete(); - } } - private async Task ReadOutputAsync(Channel output, CancellationToken cancellationToken) + private async Task ReadOutputAsync(Channel output, CancellationToken cancellationToken) { - await foreach (var command in output.Reader.ReadAllAsync(cancellationToken)) + await foreach (var commandContext in output.Reader.ReadAllAsync(cancellationToken)) { try { using var scope = _scopeFactory.CreateScope(); var commandSender = scope.ServiceProvider.GetRequiredService(); - await commandSender.SendAsync(command, CommandStrategy.Default, cancellationToken); + await commandSender.SendAsync(commandContext.Command, CommandStrategy.Default, commandContext.Headers, cancellationToken); } catch (Exception e) { diff --git a/src/common/Elsa.Mediator/Middleware/Command/CommandContext.cs b/src/common/Elsa.Mediator/Middleware/Command/CommandContext.cs index b7e61b836..f04384942 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/CommandContext.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/CommandContext.cs @@ -10,11 +10,13 @@ public class CommandContext /// /// Initializes a new instance of the class. /// - public CommandContext(ICommand command, ICommandStrategy commandStrategy, Type resultType, CancellationToken cancellationToken) + public CommandContext(ICommand command, ICommandStrategy commandStrategy, Type resultType, IDictionary headers, IServiceProvider serviceProvider, CancellationToken cancellationToken) { Command = command; CommandStrategy = commandStrategy; ResultType = resultType; + Headers = headers; + ServiceProvider = serviceProvider; CancellationToken = cancellationToken; } @@ -28,6 +30,16 @@ public class CommandContext /// public ICommandStrategy CommandStrategy { get; } + /// + /// Gets or sets the headers associated with the command context. + /// + public IDictionary Headers { get; } + + /// + /// Gets the service provider used to resolve services for the command context. + /// + public IServiceProvider ServiceProvider { get; } + /// /// Gets the result type. /// diff --git a/src/common/Elsa.Mediator/Middleware/Command/CommandPipeline.cs b/src/common/Elsa.Mediator/Middleware/Command/CommandPipeline.cs index f0e060c9b..3af4df116 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/CommandPipeline.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/CommandPipeline.cs @@ -5,28 +5,29 @@ namespace Elsa.Mediator.Middleware.Command; /// public class CommandPipeline : ICommandPipeline { - private readonly IServiceProvider _serviceProvider; - private CommandMiddlewareDelegate? _pipeline; + private readonly CommandPipelineBuilder _builder; + private CommandMiddlewareDelegate _pipeline = null!; /// /// Constructor. /// - public CommandPipeline(IServiceProvider serviceProvider) => _serviceProvider = serviceProvider; - - /// - public CommandMiddlewareDelegate Pipeline => _pipeline ??= CreateDefaultPipeline(); + public CommandPipeline(IServiceProvider serviceProvider) + { + _builder = new(serviceProvider); + Setup(x => x.UseCommandInvoker().UseCommandLogging()); + } /// - public CommandMiddlewareDelegate Setup(Action? setup = default) + public CommandMiddlewareDelegate Pipeline => _pipeline; + + /// + public CommandMiddlewareDelegate Setup(Action? setup = null) { - var builder = new CommandPipelineBuilder(_serviceProvider); - setup?.Invoke(builder); - _pipeline = builder.Build(); + setup?.Invoke(_builder); + _pipeline = _builder.Build(); return _pipeline; } /// public async Task InvokeAsync(CommandContext context) => await Pipeline(context); - - private CommandMiddlewareDelegate CreateDefaultPipeline() => Setup(x => x.UseCommandInvoker().UseCommandLogging()); } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Middleware/Command/CommandPipelineBuilder.cs b/src/common/Elsa.Mediator/Middleware/Command/CommandPipelineBuilder.cs index 8b98d49a3..94012e682 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/CommandPipelineBuilder.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/CommandPipelineBuilder.cs @@ -33,19 +33,45 @@ public class CommandPipelineBuilder : ICommandPipelineBuilder return this; } + /// + public ICommandPipelineBuilder Use(int index, Func middleware) + { + _components.Insert(index, middleware); + return this; + } + + /// + public ICommandPipelineBuilder Remove(Func middleware) + { + _components.Remove(middleware); + return this; + } + + /// + public ICommandPipelineBuilder RemoveAt(int index) + { + _components.RemoveAt(index); + return this; + } + + /// + public ICommandPipelineBuilder Clear() + { + _components.Clear(); + return this; + } + /// public CommandMiddlewareDelegate Build() { - CommandMiddlewareDelegate pipeline = _ => new ValueTask(); + CommandMiddlewareDelegate pipeline = _ => new(); - for (int i = _components.Count - 1; i >= 0; i--) - { + for (var i = _components.Count - 1; i >= 0; i--) pipeline = _components[i](pipeline); - } return pipeline; } - private T? GetProperty(string key) => Properties.TryGetValue(key, out var value) ? (T?)value : default(T); + private T? GetProperty(string key) => Properties.TryGetValue(key, out var value) ? (T?)value : default; private void SetProperty(string key, T value) => Properties[key] = value; } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Middleware/Command/Components/CommandHandlerInvokerMiddleware.cs b/src/common/Elsa.Mediator/Middleware/Command/Components/CommandHandlerInvokerMiddleware.cs index 7ad8487eb..47cceb0d9 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/Components/CommandHandlerInvokerMiddleware.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/Components/CommandHandlerInvokerMiddleware.cs @@ -1,28 +1,17 @@ using Elsa.Mediator.Contexts; using Elsa.Mediator.Contracts; using Elsa.Mediator.Middleware.Command.Contracts; +using JetBrains.Annotations; +using Microsoft.Extensions.DependencyInjection; namespace Elsa.Mediator.Middleware.Command.Components; /// /// A command middleware that invokes the command. /// -public class CommandHandlerInvokerMiddleware : ICommandMiddleware +[UsedImplicitly] +public class CommandHandlerInvokerMiddleware(CommandMiddlewareDelegate next) : ICommandMiddleware { - private readonly CommandMiddlewareDelegate _next; - private readonly IServiceProvider _serviceProvider; - private readonly IEnumerable _commandHandlers; - - /// - /// Constructor. - /// - public CommandHandlerInvokerMiddleware(CommandMiddlewareDelegate next, IEnumerable commandHandlers, IServiceProvider serviceProvider) - { - _next = next; - _serviceProvider = serviceProvider; - _commandHandlers = commandHandlers.DistinctBy(x => x.GetType()).ToList(); - } - /// public async ValueTask InvokeAsync(CommandContext context) { @@ -31,7 +20,9 @@ public class CommandHandlerInvokerMiddleware : ICommandMiddleware var commandType = command.GetType(); var resultType = context.ResultType; var handlerType = typeof(ICommandHandler<,>).MakeGenericType(commandType, resultType); - var handlers = _commandHandlers.Where(x => handlerType.IsInstanceOfType(x)).ToArray(); + var serviceProvider = context.ServiceProvider; + var commandHandlers = serviceProvider.GetServices(); + var handlers = commandHandlers.DistinctBy(x => x.GetType()).Where(x => handlerType.IsInstanceOfType(x)).ToArray(); if (handlers.Length == 0) throw new InvalidOperationException($"There is no handler to handle the {commandType.FullName} command"); @@ -40,20 +31,20 @@ public class CommandHandlerInvokerMiddleware : ICommandMiddleware throw new InvalidOperationException($"Multiple handlers were found to handle the {commandType.FullName} command"); var handler = handlers.First(); - var strategyContext = new CommandStrategyContext(command, handler, _serviceProvider, context.CancellationToken); + var strategyContext = new CommandStrategyContext(context, handler, serviceProvider, context.CancellationToken); var strategy = context.CommandStrategy; var executeMethod = strategy.GetType().GetMethod(nameof(ICommandStrategy.ExecuteAsync))!; var executeMethodWithReturnType = executeMethod.MakeGenericMethod(resultType); // Execute command. - var task = executeMethodWithReturnType.Invoke(strategy, new object[] { strategyContext }); + var task = executeMethodWithReturnType.Invoke(strategy, [strategyContext]); - // Get result of task. + // Get the result of the task. var taskWithReturnType = typeof(Task<>).MakeGenericType(resultType); var resultProperty = taskWithReturnType.GetProperty(nameof(Task.Result))!; context.Result = resultProperty.GetValue(task); // Invoke next middleware. - await _next(context); + await next(context); } } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Middleware/Command/Components/CommandLoggingMiddleware.cs b/src/common/Elsa.Mediator/Middleware/Command/Components/CommandLoggingMiddleware.cs index a4a51b426..910f03f6c 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/Components/CommandLoggingMiddleware.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/Components/CommandLoggingMiddleware.cs @@ -1,5 +1,6 @@ using Elsa.Mediator.Middleware.Command.Contracts; using Elsa.Mediator.Models; +using JetBrains.Annotations; using Microsoft.Extensions.Logging; namespace Elsa.Mediator.Middleware.Command.Components; @@ -7,31 +8,20 @@ namespace Elsa.Mediator.Middleware.Command.Components; /// /// A command middleware that logs the command being invoked. /// -public class CommandLoggingMiddleware : ICommandMiddleware +[UsedImplicitly] +public class CommandLoggingMiddleware(CommandMiddlewareDelegate next, ILogger logger) : ICommandMiddleware { - private readonly CommandMiddlewareDelegate _next; - private readonly ILogger _logger; - - /// - /// Constructor. - /// - public CommandLoggingMiddleware(CommandMiddlewareDelegate next, ILogger logger) - { - _next = next; - _logger = logger; - } - /// public async ValueTask InvokeAsync(CommandContext context) { var commandType = context.Command.GetType(); - _logger.LogInformation("Invoking {CommandName}", commandType.Name); + logger.LogInformation("Invoking {CommandName}", commandType.Name); - await _next(context); + await next(context); if (context.Result is null or Unit) - _logger.LogInformation("{CommandName} completed with no result", commandType.Name); + logger.LogInformation("{CommandName} completed with no result", commandType.Name); else - _logger.LogInformation("{CommandName} completed wit result {CommandResult}", commandType.Name, context.Result); + logger.LogInformation("{CommandName} completed with result {CommandResult}", commandType.Name, context.Result); } } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipeline.cs b/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipeline.cs index a74048b18..b049f1c86 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipeline.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipeline.cs @@ -1,7 +1,7 @@ namespace Elsa.Mediator.Middleware.Command.Contracts; /// -/// +/// Represents a pipeline for processing commands. The pipeline is responsible for orchestrating the execution of registered middleware in sequence. /// public interface ICommandPipeline { diff --git a/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipelineBuilder.cs b/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipelineBuilder.cs index ca5c763cd..b6ac7c83d 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipelineBuilder.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/Contracts/ICommandPipelineBuilder.cs @@ -16,11 +16,29 @@ public interface ICommandPipelineBuilder IServiceProvider ApplicationServices { get; } /// - /// Adds a middleware component to the pipeline. + /// Appends a middleware component to the pipeline. /// - /// The middleware component. - /// The pipeline builder. ICommandPipelineBuilder Use(Func middleware); + + /// + /// Adds a middleware component at the specified index. + /// + ICommandPipelineBuilder Use(int index, Func middleware); + + /// + /// Removes a middleware component from the pipeline. + /// + ICommandPipelineBuilder Remove(Func middleware); + + /// + /// Removes a middleware component at the specified index from the pipeline. + /// + ICommandPipelineBuilder RemoveAt(int index); + + /// + /// Clears the pipeline. + /// + ICommandPipelineBuilder Clear(); /// /// Builds the pipeline. diff --git a/src/common/Elsa.Mediator/Middleware/Command/MiddlewareExtensions.cs b/src/common/Elsa.Mediator/Middleware/Command/MiddlewareExtensions.cs index 7ab5e20f9..deac86de2 100644 --- a/src/common/Elsa.Mediator/Middleware/Command/MiddlewareExtensions.cs +++ b/src/common/Elsa.Mediator/Middleware/Command/MiddlewareExtensions.cs @@ -11,20 +11,35 @@ public static class MiddlewareExtensions /// /// Adds middleware to the pipeline. /// - /// The pipeline builder. - /// Any arguments to pass to the middleware constructor. - /// The middleware type. - /// The pipeline builder. public static ICommandPipelineBuilder UseMiddleware(this ICommandPipelineBuilder builder, params object[] args) where TMiddleware : ICommandMiddleware { - var middleware = typeof(TMiddleware); + return builder.Use(next => BuildMiddlewareDelegate(builder, next, args)); + } - return builder.Use(next => + /// + /// Inserts middleware at a specific index in the pipeline. + /// + public static ICommandPipelineBuilder UseMiddleware(this ICommandPipelineBuilder builder, int index, params object[] args) where TMiddleware : ICommandMiddleware + { + return builder.Use(index, next => BuildMiddlewareDelegate(builder, next, args)); + } + + /// + /// Builds a delegate for the middleware type. + /// + private static CommandMiddlewareDelegate BuildMiddlewareDelegate( + ICommandPipelineBuilder builder, + CommandMiddlewareDelegate next, + object[] args + ) where TMiddleware : ICommandMiddleware + { + var middleware = typeof(TMiddleware); + var invokeMethod = MiddlewareHelpers.GetInvokeMethod(middleware); + var ctorParams = new[] { - var invokeMethod = MiddlewareHelpers.GetInvokeMethod(middleware); - var ctorParams = new[] { next }.Concat(args).Select(x => x!).ToArray(); - var instance = ActivatorUtilities.CreateInstance(builder.ApplicationServices, middleware, ctorParams); - return (CommandMiddlewareDelegate)invokeMethod.CreateDelegate(typeof(CommandMiddlewareDelegate), instance); - }); + next + }.Concat(args).Select(x => x!).ToArray(); + var instance = ActivatorUtilities.CreateInstance(builder.ApplicationServices, middleware, ctorParams); + return (CommandMiddlewareDelegate)invokeMethod.CreateDelegate(typeof(CommandMiddlewareDelegate), instance); } } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs b/src/common/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs index c555a7f5f..874c64f17 100644 --- a/src/common/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs +++ b/src/common/Elsa.Mediator/Middleware/Notification/Components/NotificationHandlerInvokerMiddleware.cs @@ -1,33 +1,19 @@ using Elsa.Mediator.Contexts; using Elsa.Mediator.Contracts; using Elsa.Mediator.Middleware.Notification.Contracts; +using JetBrains.Annotations; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; namespace Elsa.Mediator.Middleware.Notification.Components; /// -public class NotificationHandlerInvokerMiddleware : INotificationMiddleware +[UsedImplicitly] +public class NotificationHandlerInvokerMiddleware( + NotificationMiddlewareDelegate next, + ILogger logger) + : INotificationMiddleware { - private readonly NotificationMiddlewareDelegate _next; - private readonly ILogger _logger; - private readonly IServiceProvider _serviceProvider; - private readonly IEnumerable _notificationHandlers; - - /// - /// Initializes a new instance of the class. - /// - public NotificationHandlerInvokerMiddleware( - NotificationMiddlewareDelegate next, - ILogger logger, - IServiceProvider serviceProvider, - IEnumerable notificationHandlers) - { - _next = next; - _logger = logger; - _serviceProvider = serviceProvider; - _notificationHandlers = notificationHandlers; - } - /// public async ValueTask InvokeAsync(NotificationContext context) { @@ -35,12 +21,14 @@ public class NotificationHandlerInvokerMiddleware : INotificationMiddleware var notification = context.Notification; var notificationType = notification.GetType(); var handlerType = typeof(INotificationHandler<>).MakeGenericType(notificationType); - var handlers = _notificationHandlers.Where(x => handlerType.IsInstanceOfType(x)).DistinctBy(x => x.GetType()).ToArray(); - var strategyContext = new NotificationStrategyContext(notification, handlers, _logger, _serviceProvider, context.CancellationToken); + var serviceProvider = context.ServiceProvider; + var notificationHandlers = serviceProvider.GetServices(); + var handlers = notificationHandlers.Where(x => handlerType.IsInstanceOfType(x)).DistinctBy(x => x.GetType()).ToArray(); + var strategyContext = new NotificationStrategyContext(notification, handlers, logger, serviceProvider, context.CancellationToken); await context.NotificationStrategy.PublishAsync(strategyContext); // Invoke next middleware. - await _next(context); + await next(context); } } diff --git a/src/common/Elsa.Mediator/Middleware/Notification/MiddlewareExtensions.cs b/src/common/Elsa.Mediator/Middleware/Notification/MiddlewareExtensions.cs index c45fc9f27..5be932c94 100644 --- a/src/common/Elsa.Mediator/Middleware/Notification/MiddlewareExtensions.cs +++ b/src/common/Elsa.Mediator/Middleware/Notification/MiddlewareExtensions.cs @@ -21,7 +21,7 @@ public static class MiddlewareExtensions return builder.Use(next => { var invokeMethod = MiddlewareHelpers.GetInvokeMethod(middleware); - var ctorParams = new[] { next }.Concat(args).Select(x => x!).ToArray(); + var ctorParams = new[] { next }.Concat(args).Select(x => x).ToArray(); var instance = ActivatorUtilities.CreateInstance(builder.ApplicationServices, middleware, ctorParams); return (NotificationMiddlewareDelegate)invokeMethod.CreateDelegate(typeof(NotificationMiddlewareDelegate), instance); }); diff --git a/src/common/Elsa.Mediator/Middleware/Notification/NotificationContext.cs b/src/common/Elsa.Mediator/Middleware/Notification/NotificationContext.cs index 60a6dd1e4..4f4f63048 100644 --- a/src/common/Elsa.Mediator/Middleware/Notification/NotificationContext.cs +++ b/src/common/Elsa.Mediator/Middleware/Notification/NotificationContext.cs @@ -12,11 +12,13 @@ public class NotificationContext /// /// The notification to publish. /// The publishing strategy to use. + /// The service provider to resolve services from. /// The cancellation token. - public NotificationContext(INotification notification, IEventPublishingStrategy notificationStrategy, CancellationToken cancellationToken = default) + public NotificationContext(INotification notification, IEventPublishingStrategy notificationStrategy, IServiceProvider serviceProvider, CancellationToken cancellationToken = default) { Notification = notification; NotificationStrategy = notificationStrategy; + ServiceProvider = serviceProvider; CancellationToken = cancellationToken; } @@ -29,7 +31,12 @@ public class NotificationContext /// Gets the publishing strategy to use. /// public IEventPublishingStrategy NotificationStrategy { get; init; } - + + /// + /// Gets the service provider used for resolving dependencies within the notification context. + /// + public IServiceProvider ServiceProvider { get; } + /// /// Gets the cancellation token. /// diff --git a/src/common/Elsa.Mediator/Middleware/Request/RequestContext.cs b/src/common/Elsa.Mediator/Middleware/Request/RequestContext.cs index b6643a567..422c0e102 100644 --- a/src/common/Elsa.Mediator/Middleware/Request/RequestContext.cs +++ b/src/common/Elsa.Mediator/Middleware/Request/RequestContext.cs @@ -5,35 +5,27 @@ namespace Elsa.Mediator.Middleware.Request; /// /// Provides context to a request handler. /// -public class RequestContext +public class RequestContext(IRequest request, Type responseType, IServiceProvider serviceProvider, CancellationToken cancellationToken) { - /// - /// Initializes a new instance of the class. - /// - /// The request. - /// The response type. - /// The cancellation token. - public RequestContext(IRequest request, Type responseType, CancellationToken cancellationToken) - { - Request = request; - ResponseType = responseType; - CancellationToken = cancellationToken; - } - /// /// Gets the request. /// - public IRequest Request { get; init; } - + public IRequest Request { get; init; } = request; + /// /// Gets the response type. /// - public Type ResponseType { get; init; } - + public Type ResponseType { get; init; } = responseType; + + /// + /// Gets the service provider used for resolving dependencies within the request context. + /// + public IServiceProvider ServiceProvider { get; } = serviceProvider; + /// /// Gets the cancellation token. /// - public CancellationToken CancellationToken { get; init; } + public CancellationToken CancellationToken { get; init; } = cancellationToken; /// /// Gets the response the request handler. diff --git a/src/common/Elsa.Mediator/Services/DefaultMediator.cs b/src/common/Elsa.Mediator/Services/DefaultMediator.cs index 23a68162b..40587c3b7 100644 --- a/src/common/Elsa.Mediator/Services/DefaultMediator.cs +++ b/src/common/Elsa.Mediator/Services/DefaultMediator.cs @@ -17,6 +17,7 @@ public class DefaultMediator : IMediator private readonly IRequestPipeline _requestPipeline; private readonly ICommandPipeline _commandPipeline; private readonly INotificationPipeline _notificationPipeline; + private readonly IServiceProvider _serviceProvider; private readonly IEventPublishingStrategy _defaultPublishingStrategy; private readonly ICommandStrategy _defaultCommandStrategy; @@ -31,11 +32,13 @@ public class DefaultMediator : IMediator IRequestPipeline requestPipeline, ICommandPipeline commandPipeline, INotificationPipeline notificationPipeline, - IOptions options) + IOptions options, + IServiceProvider serviceProvider) { _requestPipeline = requestPipeline; _commandPipeline = commandPipeline; _notificationPipeline = notificationPipeline; + _serviceProvider = serviceProvider; _defaultPublishingStrategy = options.Value.DefaultPublishingStrategy; _defaultCommandStrategy = options.Value.DefaultCommandStrategy; } @@ -44,46 +47,66 @@ public class DefaultMediator : IMediator public async Task SendAsync(IRequest request, CancellationToken cancellationToken = default) { var responseType = typeof(T); - var context = new RequestContext(request, responseType, cancellationToken); + var context = new RequestContext(request, responseType, _serviceProvider, cancellationToken); await _requestPipeline.ExecuteAsync(context); return (T)context.Response; } - /// - public async Task SendAsync(ICommand command, CancellationToken cancellationToken = default) => await SendAsync(command, _defaultCommandStrategy, cancellationToken); - - /// - public async Task SendAsync(ICommand command, ICommandStrategy? strategy = null, CancellationToken cancellationToken = default) - { - var resultType = typeof(Unit); - strategy ??= _defaultCommandStrategy; - var context = new CommandContext(command, strategy, resultType, cancellationToken); - await _commandPipeline.InvokeAsync(context); - } - - /// - public async Task SendAsync(ICommand command, CancellationToken cancellationToken = default) => await SendAsync(command, _defaultCommandStrategy, cancellationToken); - - /// - public async Task SendAsync(ICommand command, ICommandStrategy? strategy, CancellationToken cancellationToken = default) + public async Task SendAsync(ICommand command, ICommandStrategy strategy, IDictionary headers, CancellationToken cancellationToken = default) { var resultType = typeof(T); - strategy ??= _defaultCommandStrategy; - var context = new CommandContext(command, strategy, resultType, cancellationToken); + var context = new CommandContext(command, strategy, resultType, headers, _serviceProvider, cancellationToken); await _commandPipeline.InvokeAsync(context); return (T)context.Result!; } /// - public async Task SendAsync(INotification notification, CancellationToken cancellationToken = default) => await SendAsync(notification, _defaultPublishingStrategy, cancellationToken); + public async Task SendAsync(ICommand command, CancellationToken cancellationToken = default) => await SendAsync(command, _defaultCommandStrategy, cancellationToken); + + /// + public Task SendAsync(ICommand command, ICommandStrategy? strategy = null, CancellationToken cancellationToken = default) + { + return SendAsync(command, strategy, new Dictionary(), cancellationToken); + } + + public async Task SendAsync(ICommand command, ICommandStrategy? strategy, IDictionary headers, CancellationToken cancellationToken = default) + { + var resultType = typeof(Unit); + strategy ??= _defaultCommandStrategy; + var context = new CommandContext(command, strategy, resultType, headers, _serviceProvider, cancellationToken); + await _commandPipeline.InvokeAsync(context); + } + + /// + public Task SendAsync(ICommand command, CancellationToken cancellationToken = default) + { + return SendAsync(command, new Dictionary(), cancellationToken); + } + + public Task SendAsync(ICommand command, IDictionary headers, CancellationToken cancellationToken = default) + { + return SendAsync(command, _defaultCommandStrategy, headers, cancellationToken); + } + + /// + public Task SendAsync(ICommand command, ICommandStrategy strategy, CancellationToken cancellationToken = default) + { + return SendAsync(command, strategy, new Dictionary(), cancellationToken); + } + + /// + public async Task SendAsync(INotification notification, CancellationToken cancellationToken = default) + { + await SendAsync(notification, _defaultPublishingStrategy, cancellationToken); + } /// public async Task SendAsync(INotification notification, IEventPublishingStrategy? strategy = null, CancellationToken cancellationToken = default) { strategy ??= _defaultPublishingStrategy; - var context = new NotificationContext(notification, strategy, cancellationToken); + var context = new NotificationContext(notification, strategy, _serviceProvider, cancellationToken); await _notificationPipeline.ExecuteAsync(context); } } \ No newline at end of file diff --git a/src/modules/Elsa.Tenants/Features/TenantsFeature.cs b/src/modules/Elsa.Tenants/Features/TenantsFeature.cs index ce4a18a9a..e19e248b7 100644 --- a/src/modules/Elsa.Tenants/Features/TenantsFeature.cs +++ b/src/modules/Elsa.Tenants/Features/TenantsFeature.cs @@ -3,6 +3,7 @@ using Elsa.Common.Multitenancy; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.Tenants.Mediator.Tasks; using Elsa.Tenants.Options; using Elsa.Tenants.Providers; using Microsoft.Extensions.DependencyInjection; @@ -45,6 +46,11 @@ public class TenantsFeature(IModule serviceConfiguration) : FeatureBase(serviceC Module.Configure(feature => feature.UseTenantsProvider()); } + public override void ConfigureHostedServices() + { + Module.ConfigureHostedService(); + } + /// public override void Apply() { diff --git a/src/modules/Elsa.Tenants/Mediator/Middleware/TenantPropagatingMiddleware.cs b/src/modules/Elsa.Tenants/Mediator/Middleware/TenantPropagatingMiddleware.cs new file mode 100644 index 000000000..ad18be2a8 --- /dev/null +++ b/src/modules/Elsa.Tenants/Mediator/Middleware/TenantPropagatingMiddleware.cs @@ -0,0 +1,28 @@ +using Elsa.Common.Multitenancy; +using Elsa.Mediator.Middleware.Command; +using Elsa.Mediator.Middleware.Command.Contracts; +using JetBrains.Annotations; + +namespace Elsa.Tenants.Mediator.Middleware; + +/// +/// Middleware that ensures tenant context is propagated through the request pipeline. +/// +[UsedImplicitly] +public class TenantPropagatingMiddleware(CommandMiddlewareDelegate next, ITenantScopeFactory tenantScopeFactory, ITenantService tenantService) : ICommandMiddleware +{ + /// + public async ValueTask InvokeAsync(CommandContext context) + { + if (context.Headers.TryGetValue(TenantHeaders.TenantIdKey, out var tenantIdVal)) + { + var tenantId = (string)tenantIdVal; + var tenant = await tenantService.FindAsync(tenantId); + await using var tenantScope = tenantScopeFactory.CreateScope(tenant); + await next(context); + return; + } + + await next(context); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Tenants/Mediator/Tasks/SetupMediatorPipelines.cs b/src/modules/Elsa.Tenants/Mediator/Tasks/SetupMediatorPipelines.cs new file mode 100644 index 000000000..33613f09a --- /dev/null +++ b/src/modules/Elsa.Tenants/Mediator/Tasks/SetupMediatorPipelines.cs @@ -0,0 +1,19 @@ +using Elsa.Mediator.Middleware.Command; +using Elsa.Mediator.Middleware.Command.Contracts; +using Elsa.Tenants.Mediator.Middleware; +using JetBrains.Annotations; +using Microsoft.Extensions.Hosting; + +namespace Elsa.Tenants.Mediator.Tasks; + +[UsedImplicitly] +public class SetupMediatorPipelines(ICommandPipeline commandPipeline) : IHostedService +{ + public Task StartAsync(CancellationToken cancellationToken) + { + commandPipeline.Setup(pipeline => pipeline.UseMiddleware(0)); + return Task.CompletedTask; + } + + public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; +} \ No newline at end of file diff --git a/src/modules/Elsa.Tenants/Mediator/TenantHeaders.cs b/src/modules/Elsa.Tenants/Mediator/TenantHeaders.cs new file mode 100644 index 000000000..a6aa66334 --- /dev/null +++ b/src/modules/Elsa.Tenants/Mediator/TenantHeaders.cs @@ -0,0 +1,13 @@ +namespace Elsa.Tenants.Mediator; + +public static class TenantHeaders +{ + public static readonly object TenantIdKey = new(); + + public static IDictionary CreateHeaders(string? tenantId) + { + var headers = new Dictionary(); + if (tenantId != null) headers.Add(TenantIdKey, tenantId); + return headers; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs index 09e889628..d0279afcb 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs @@ -1,5 +1,7 @@ +using Elsa.Common.Multitenancy; using Elsa.Mediator; using Elsa.Mediator.Contracts; +using Elsa.Tenants.Mediator; using Elsa.Workflows.Runtime.Commands; using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.Runtime.Responses; @@ -9,18 +11,8 @@ namespace Elsa.Workflows.Runtime; /// /// A simple implementation that queues the specified request for workflow execution on a non-durable background worker. /// -public class BackgroundWorkflowDispatcher : IWorkflowDispatcher +public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher { - private readonly ICommandSender _commandSender; - - /// - /// Constructor. - /// - public BackgroundWorkflowDispatcher(ICommandSender commandSender) - { - _commandSender = commandSender; - } - /// public async Task DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { @@ -32,8 +24,8 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher InstanceId = request.InstanceId, TriggerActivityId = request.TriggerActivityId }; - - await _commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken); + + await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); return DispatchWorkflowResponse.Success(); } @@ -47,7 +39,7 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher Properties = request.Properties, CorrelationId = request.CorrelationId}; - await _commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken); + await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); return DispatchWorkflowResponse.Success(); } @@ -62,7 +54,7 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher Input = request.Input, Properties = request.Properties }; - await _commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken); + await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); return DispatchWorkflowResponse.Success(); } @@ -76,7 +68,12 @@ public class BackgroundWorkflowDispatcher : IWorkflowDispatcher ActivityInstanceId = request.ActivityInstanceId, Input = request.Input }; - await _commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken); + await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); return DispatchWorkflowResponse.Success(); } + + private IDictionary CreateHeaders() + { + return TenantHeaders.CreateHeaders(tenantAccessor.Tenant?.Id); + } } \ No newline at end of file