From 8be9b24f2a91b676b3d5661f763ed6048cfea424 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 10 Mar 2026 12:44:18 +0100 Subject: [PATCH 1/7] Revise changelog for version 3.6.0 Updated breaking changes and upgrade notes for version 3.6.0, including package name changes, database migration requirements, and multitenancy ID conventions. --- doc/changelogs/3.6.0.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/doc/changelogs/3.6.0.md b/doc/changelogs/3.6.0.md index 7e8353cfc..0000c6073 100644 --- a/doc/changelogs/3.6.0.md +++ b/doc/changelogs/3.6.0.md @@ -6,6 +6,8 @@ Compare: [`3.5.3...3.6.0`](https://github.com/elsa-workflows/elsa-core/compare/3 ## ⚠️ Breaking changes / upgrade notes +- **EF Core package names have changed**: EF Core persistence packages were renamed from `Elsa.EntityFrameworkCore.*` to `Elsa.Persistence.EFCore.*`. If your application references any of the old package names, you must update them to the new package names when upgrading to 3.6.0. Be sure to review your project files, internal package feeds, CI pipelines, and deployment manifests for old package references. + - **Database migrations required (EF Core — all providers)**: `ActivityNodeId` columns in `ActivityExecutionRecords` and `WorkflowExecutionLogRecords` have been widened to unlimited types (`nvarchar(max)` / `longtext` / `NCLOB`) to support deeply nested workflows. The corresponding B-tree indexes (`IX_ActivityExecutionRecord_ActivityNodeId`, `IX_WorkflowExecutionLogRecord_ActivityNodeId`) are dropped as part of the V3_6 migrations. **Run EF Core migrations before upgrading** any SQL Server, MySQL, or Oracle deployment to 3.6.0. ([71438596f3](https://github.com/elsa-workflows/elsa-core/commit/71438596f3)) ([#7338](https://github.com/elsa-workflows/elsa-core/pull/7338)) - **Multitenancy — tenant ID convention change**: An **empty string (`""`)** is now the canonical default tenant ID for all tenant-aware entities; `null` now means **tenant-agnostic** (visible to all tenants). EF Core query filters and the `ActivityRegistry` have been updated accordingly. If your database contains rows with a `null` `TenantId` that were intended to represent the default tenant, migrate those rows to `""` before upgrading. The new `NormalizeTenantId()` extension method on `string` handles the conversion in code. ([#7217](https://github.com/elsa-workflows/elsa-core/pull/7217), [#7226](https://github.com/elsa-workflows/elsa-core/pull/7226)) @@ -272,4 +274,4 @@ Use `CounterBasedExecution` to keep the pre-3.6.0 behavior if existing workflow * Removed `Elsa.ServerAndStudio.Web` and related sample projects from solution. ([ecf5b390f1](https://github.com/elsa-workflows/elsa-core/commit/ecf5b390f1)) * Updated documentation to reflect .NET 10.0 support and remove deprecated external dependency references. ([7add1030c4](https://github.com/elsa-workflows/elsa-core/commit/7add1030c4)) * Added DeepWiki badge to README. ([b3ad57191a](https://github.com/elsa-workflows/elsa-core/commit/b3ad57191a)) -* Null safety and compiler warning fixes across multiple modules. ([490c8a2c9e](https://github.com/elsa-workflows/elsa-core/commit/490c8a2c9e), [2c0b3da5de](https://github.com/elsa-workflows/elsa-core/commit/2c0b3da5de)) ([#7050](https://github.com/elsa-workflows/elsa-core/pull/7050), [#7051](https://github.com/elsa-workflows/elsa-core/pull/7051)) \ No newline at end of file +* Null safety and compiler warning fixes across multiple modules. ([490c8a2c9e](https://github.com/elsa-workflows/elsa-core/commit/490c8a2c9e), [2c0b3da5de](https://github.com/elsa-workflows/elsa-core/commit/2c0b3da5de)) ([#7050](https://github.com/elsa-workflows/elsa-core/pull/7050), [#7051](https://github.com/elsa-workflows/elsa-core/pull/7051)) From 291d4a81c1df83fbea5ae25b16555641f2d0c4ce Mon Sep 17 00:00:00 2001 From: j03y-nxxbz <113447825+j03y-nxxbz@users.noreply.github.com> Date: Wed, 11 Mar 2026 09:44:21 +0100 Subject: [PATCH 2/7] Update packages.yml: set base version to 3.6.1 --- .github/workflows/packages.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index 16e623e86..a46ddb82b 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -18,7 +18,7 @@ on: types: [prereleased, published] env: - base_version: '3.6.0' + base_version: '3.6.1' feedz_feed_source: 'https://f.feedz.io/elsa-workflows/elsa-3/nuget/index.json' nuget_feed_source: 'https://api.nuget.org/v3/index.json' From a0f0a3c2f5206bf06d29e59aade8c199acae31f1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 11 Mar 2026 17:20:07 +0100 Subject: [PATCH 3/7] Add workflow export functionality and update changelog (#7357) * Add `IWorkflowDefinitionExporter` for workflow export functionality Introduced `IWorkflowDefinitionExporter` interface and its implementation to export workflow definitions as JSON or ZIP archives. Simplified `Export` endpoint logic by utilizing the new exporter service. Updated package version to 3.6.1. * Remove duplicate `ExportAsync` method from `IWorkflowDefinitionExporter` and its implementation in `WorkflowDefinitionExporter`. * Remove deprecated test from `WorkflowDefinitionExporterTests`, regression tests are covered in `WorkflowReferenceGraphBuilderTests`. * Update GitHub Actions workflow to support version 3.6.1 deployment * Update src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionExporter.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Move `WorkflowDefinitionExporterRegressionTests` to separate test file for clarity * Add `DefaultFileNameSanitizer` for sanitizing file names and update `WorkflowDefinitionExporter` to use it. Implement unit tests for the sanitizer. * Update test assertion for exported workflow file name in `WorkflowDefinitionExporterRegressionTests`. * Simplify file naming in `WorkflowDefinitionExporter` by removing duplicate ID from JSON file names. * Update test/unit/Elsa.Workflows.Management.UnitTests/Services/DefaultFileNameSanitizerTests.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .github/workflows/packages.yml | 2 +- .../WorkflowDefinitions/Export/Endpoint.cs | 203 ++---------------- .../Contracts/IFileNameSanitizer.cs | 13 ++ .../Contracts/IWorkflowDefinitionExporter.cs | 46 ++++ .../Features/WorkflowManagementFeature.cs | 2 + .../Models/WorkflowDefinitionExportResult.cs | 8 + .../Services/DefaultFileNameSanitizer.cs | 33 +++ .../Services/WorkflowDefinitionExporter.cs | 176 +++++++++++++++ .../Services/DefaultFileNameSanitizerTests.cs | 35 +++ ...rkflowDefinitionExporterRegressionTests.cs | 80 +++++++ .../WorkflowReferenceGraphBuilderTests.cs | 2 +- 11 files changed, 418 insertions(+), 182 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IFileNameSanitizer.cs create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionExporter.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionExportResult.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/DefaultFileNameSanitizer.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionExporter.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Services/DefaultFileNameSanitizerTests.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionExporterRegressionTests.cs diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index a46ddb82b..5c931e7a7 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -202,7 +202,7 @@ jobs: name: Deploy coverage to GitHub Pages needs: test runs-on: ubuntu-latest - if: github.ref == 'refs/heads/main' || github.ref == 'refs/heads/develop/3.6.0' || github.ref == 'refs/heads/release/3.6.0' + if: github.ref == 'refs/heads/main' || github.ref == 'refs/heads/develop/3.6.1' || github.ref == 'refs/heads/release/3.6.1' permissions: pages: write id-token: write diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Export/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Export/Endpoint.cs index 14cf00ffd..e20a416a3 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Export/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Export/Endpoint.cs @@ -1,13 +1,6 @@ -using System.IO.Compression; -using System.Text.Json; using Elsa.Abstractions; using Elsa.Common.Models; using Elsa.Workflows.Management; -using Elsa.Workflows.Management.Entities; -using Elsa.Workflows.Management.Filters; -using Elsa.Workflows.Management.Mappers; -using Elsa.Workflows.Management.Models; -using Humanizer; using JetBrains.Annotations; namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Export; @@ -16,26 +9,8 @@ namespace Elsa.Workflows.Api.Endpoints.WorkflowDefinitions.Export; /// Exports the specified workflow definition as JSON download. /// [UsedImplicitly] -internal class Export : ElsaEndpoint +internal class Export(IWorkflowDefinitionExporter exporter) : ElsaEndpoint { - private readonly IApiSerializer _serializer; - private readonly IWorkflowDefinitionStore _store; - private readonly IWorkflowReferenceGraphBuilder _workflowReferenceGraphBuilder; - private readonly WorkflowDefinitionMapper _workflowDefinitionMapper; - - /// - public Export( - IWorkflowDefinitionStore store, - IApiSerializer serializer, - WorkflowDefinitionMapper workflowDefinitionMapper, - IWorkflowReferenceGraphBuilder workflowReferenceGraphBuilder) - { - _store = store; - _serializer = serializer; - _workflowDefinitionMapper = workflowDefinitionMapper; - _workflowReferenceGraphBuilder = workflowReferenceGraphBuilder; - } - /// public override void Configure() { @@ -48,165 +23,33 @@ internal class Export : ElsaEndpoint public override async Task HandleAsync(Request request, CancellationToken cancellationToken) { if (request.DefinitionId != null) - await DownloadSingleWorkflowAsync(request.DefinitionId, request.VersionOptions, request.IncludeConsumingWorkflows, cancellationToken); - else if (request.Ids != null) - await DownloadMultipleWorkflowsAsync(request.Ids, request.IncludeConsumingWorkflows, cancellationToken); - else await Send.NoContentAsync(cancellationToken); - } - - private async Task DownloadMultipleWorkflowsAsync(ICollection ids, bool includeConsumingWorkflows, CancellationToken cancellationToken) - { - var definitions = (await _store.FindManyAsync(new() { - Ids = ids - }, cancellationToken)).ToList(); + var versionOptions = string.IsNullOrEmpty(request.VersionOptions) ? VersionOptions.Latest : VersionOptions.FromString(request.VersionOptions); + var result = await exporter.ExportAsync(request.DefinitionId, versionOptions, request.IncludeConsumingWorkflows, cancellationToken); - if (includeConsumingWorkflows) - definitions = await IncludeConsumersAsync(definitions, cancellationToken); + if (result == null) + { + await Send.NotFoundAsync(cancellationToken); + return; + } - if (!definitions.Any()) + await Send.BytesAsync(result.Data, result.FileName, cancellation: cancellationToken); + } + else if (request.Ids != null) + { + var result = await exporter.ExportManyAsync(request.Ids, request.IncludeConsumingWorkflows, cancellationToken); + + if (result == null) + { + await Send.NoContentAsync(cancellationToken); + return; + } + + await Send.BytesAsync(result.Data, result.FileName, cancellation: cancellationToken); + } + else { await Send.NoContentAsync(cancellationToken); - return; } - - await WriteZipResponseAsync(definitions, cancellationToken); - } - - private async Task DownloadSingleWorkflowAsync(string definitionId, string? versionOptions, bool includeConsumingWorkflows, CancellationToken cancellationToken) - { - var parsedVersionOptions = string.IsNullOrEmpty(versionOptions) ? VersionOptions.Latest : VersionOptions.FromString(versionOptions); - var definition = (await _store.FindManyAsync(new() - { - DefinitionId = definitionId, - VersionOptions = parsedVersionOptions - }, cancellationToken)).FirstOrDefault(); - - if (definition == null) - { - await Send.NotFoundAsync(cancellationToken); - return; - } - - if (includeConsumingWorkflows) - { - var definitions = await IncludeConsumersAsync([definition], cancellationToken); - await WriteZipResponseAsync(definitions, cancellationToken); - return; - } - - var model = await CreateWorkflowModelAsync(definition, cancellationToken); - var binaryJson = await SerializeWorkflowDefinitionAsync(model, cancellationToken); - var fileName = GetFileName(model); - - await Send.BytesAsync(binaryJson, fileName, cancellation: cancellationToken); - } - - /// - /// Recursively discovers all consuming workflow definitions and includes them. - /// Consumers are always resolved at , regardless of the version used for the initial definitions. - /// - private async Task> IncludeConsumersAsync(List definitions, CancellationToken cancellationToken) - { - var initialDefinitionIds = definitions.Select(d => d.DefinitionId).ToList(); - var graph = await _workflowReferenceGraphBuilder.BuildGraphAsync(initialDefinitionIds, cancellationToken); - - // Find any consumer definitions not already in our list. - var newDefinitionIds = graph.ConsumerDefinitionIds.Except(initialDefinitionIds).ToList(); - - if (newDefinitionIds.Count > 0) - { - var consumerDefinitions = await _store.FindManyAsync(new WorkflowDefinitionFilter - { - DefinitionIds = newDefinitionIds.ToArray(), - VersionOptions = VersionOptions.Latest - }, cancellationToken); - - definitions = definitions.Concat(consumerDefinitions).ToList(); - } - - return definitions; - } - - private async Task WriteZipResponseAsync(List definitions, CancellationToken cancellationToken) - { - var zipStream = new MemoryStream(); - var sortedDefinitions = definitions.OrderBy(d => d.DefinitionId).ToList(); - - // NOTE: - // - ZIP timestamps cannot be earlier than 1980-01-01 (the ZIP format's minimum). - // - We intentionally use a fixed timestamp (instead of DateTimeOffset.UtcNow) to keep exports deterministic. - // This avoids producing different ZIP bytes for identical exports, which helps tests, caching, and diffing. - var zipEpoch = new DateTimeOffset(1980, 1, 1, 0, 0, 0, TimeSpan.Zero); - -#if NET10_0_OR_GREATER - await using (var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Create, true)) - { - // Create a JSON file for each workflow definition: - foreach (var definition in sortedDefinitions) - { - var model = await CreateWorkflowModelAsync(definition, cancellationToken); - var binaryJson = await SerializeWorkflowDefinitionAsync(model, cancellationToken); - var fileName = GetFileName(model); - var entry = zipArchive.CreateEntry(fileName, CompressionLevel.Optimal); - entry.LastWriteTime = zipEpoch; - await using var entryStream = await entry.OpenAsync(cancellationToken); - await entryStream.WriteAsync(binaryJson, cancellationToken); - } - } -#else - using (var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Create, true)) - { - // Create a JSON file for each workflow definition: - foreach (var definition in sortedDefinitions) - { - var model = await CreateWorkflowModelAsync(definition, cancellationToken); - var binaryJson = await SerializeWorkflowDefinitionAsync(model, cancellationToken); - var fileName = GetFileName(model); - var entry = zipArchive.CreateEntry(fileName, CompressionLevel.Optimal); - entry.LastWriteTime = zipEpoch; - await using var entryStream = entry.Open(); - await entryStream.WriteAsync(binaryJson, cancellationToken); - } - } -#endif - - // Send the zip file to the client: - zipStream.Position = 0; - await Send.BytesAsync(zipStream.ToArray(), "workflow-definitions.zip", cancellation: cancellationToken); - } - - private string GetFileName(WorkflowDefinitionModel definition) - { - var hasWorkflowName = !string.IsNullOrWhiteSpace(definition.Name); - var workflowName = hasWorkflowName ? definition.Name!.Trim() : definition.DefinitionId; - var fileName = $"workflow-definition-{workflowName.Underscore().Dasherize().ToLowerInvariant()}-{definition.DefinitionId}.json"; - return fileName; - } - - private async Task SerializeWorkflowDefinitionAsync(WorkflowDefinitionModel model, CancellationToken cancellationToken) - { - var serializerOptions = _serializer.GetOptions(); - var document = JsonSerializer.SerializeToDocument(model, serializerOptions); - var rootElement = document.RootElement; - - using var output = new MemoryStream(); - await using var writer = new Utf8JsonWriter(output); - - writer.WriteStartObject(); - writer.WriteString("$schema", "https://elsaworkflows.io/schemas/workflow-definition/v3.0.0/schema.json"); - - foreach (var property in rootElement.EnumerateObject()) - property.WriteTo(writer); - - writer.WriteEndObject(); - - await writer.FlushAsync(cancellationToken); - return output.ToArray(); - } - - private async Task CreateWorkflowModelAsync(WorkflowDefinition definition, CancellationToken cancellationToken) - { - return await _workflowDefinitionMapper.MapAsync(definition, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IFileNameSanitizer.cs b/src/modules/Elsa.Workflows.Management/Contracts/IFileNameSanitizer.cs new file mode 100644 index 000000000..052bc854a --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IFileNameSanitizer.cs @@ -0,0 +1,13 @@ +namespace Elsa.Workflows.Management; + +/// +/// Sanitizes file names by replacing invalid characters with safe alternatives. +/// +public interface IFileNameSanitizer +{ + /// + /// Replaces invalid file name characters in the specified value. + /// + string Sanitize(string value); +} + diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionExporter.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionExporter.cs new file mode 100644 index 000000000..686413092 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionExporter.cs @@ -0,0 +1,46 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Models; + +namespace Elsa.Workflows.Management; + +/// +/// Exports workflow definitions as serialized JSON or ZIP archives. +/// +public interface IWorkflowDefinitionExporter +{ + /// + /// Exports a single workflow definition as a JSON byte array, optionally including consuming workflows as a ZIP archive. + /// + /// The definition ID. + /// The version options. Defaults to . + /// When true, includes all consuming workflow definitions in the export as a ZIP archive. + /// The cancellation token. + /// An export result containing the binary content and a suggested file name, or null if the definition was not found. + Task ExportAsync(string definitionId, VersionOptions? versionOptions = null, bool includeConsumingWorkflows = false, CancellationToken cancellationToken = default); + + /// + /// Exports multiple workflow definitions as a ZIP archive. + /// + /// A list of workflow definition version IDs. + /// When true, includes all consuming workflow definitions in the export. + /// The cancellation token. + /// An export result containing the ZIP binary content and a suggested file name, or null if no definitions were found. + Task ExportManyAsync(ICollection ids, bool includeConsumingWorkflows = false, CancellationToken cancellationToken = default); + + /// + /// Serializes a single workflow definition entity to a JSON byte array (with $schema header). + /// + /// The workflow definition entity. + /// The cancellation token. + /// The serialized JSON bytes and suggested file name. + Task ExportDefinitionAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); + + /// + /// Exports multiple workflow definition entities as a ZIP archive. + /// + /// The workflow definitions to export. + /// The cancellation token. + /// The ZIP archive bytes and suggested file name. + Task ExportDefinitionsAsync(ICollection definitions, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 0164a6b03..932d923a1 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -272,6 +272,8 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module) .AddScoped() .AddScoped(_workflowDefinitionPublisher) .AddScoped() + .AddScoped() + .AddSingleton() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionExportResult.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionExportResult.cs new file mode 100644 index 000000000..8f7642298 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionExportResult.cs @@ -0,0 +1,8 @@ +namespace Elsa.Workflows.Management.Models; + +/// +/// Represents the result of exporting one or more workflow definitions. +/// +/// The binary content (JSON or ZIP). +/// The suggested file name. +public record WorkflowDefinitionExportResult(byte[] Data, string FileName); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/DefaultFileNameSanitizer.cs b/src/modules/Elsa.Workflows.Management/Services/DefaultFileNameSanitizer.cs new file mode 100644 index 000000000..44d5c45d3 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/DefaultFileNameSanitizer.cs @@ -0,0 +1,33 @@ +namespace Elsa.Workflows.Management.Services; + +/// +public class DefaultFileNameSanitizer : IFileNameSanitizer +{ + private static readonly char[] InvalidFileNameCharacters = Path.GetInvalidFileNameChars(); + + /// + public string Sanitize(string value) + { + for (var i = 0; i < value.Length; i++) + { + if (!IsInvalidFileNameCharacter(value[i])) + continue; + + return string.Create(value.Length, (value, i), static (buffer, state) => + { + state.value.AsSpan(0, state.i).CopyTo(buffer); + + for (var j = state.i; j < state.value.Length; j++) + { + var character = state.value[j]; + buffer[j] = IsInvalidFileNameCharacter(character) ? '-' : character; + } + }); + } + + return value; + } + + private static bool IsInvalidFileNameCharacter(char character) => character is '/' or '\\' || Array.IndexOf(InvalidFileNameCharacters, character) >= 0; +} + diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionExporter.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionExporter.cs new file mode 100644 index 000000000..e369f97de --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionExporter.cs @@ -0,0 +1,176 @@ +using System.IO.Compression; +using System.Text.Json; +using Elsa.Common.Models; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Mappers; +using Elsa.Workflows.Management.Models; +using Humanizer; + +namespace Elsa.Workflows.Management.Services; + +/// +public class WorkflowDefinitionExporter( + IWorkflowDefinitionStore store, + IApiSerializer serializer, + WorkflowDefinitionMapper workflowDefinitionMapper, + IWorkflowReferenceGraphBuilder workflowReferenceGraphBuilder, + IFileNameSanitizer fileNameSanitizer) : IWorkflowDefinitionExporter +{ + /// + public async Task ExportAsync(string definitionId, VersionOptions? versionOptions = null, bool includeConsumingWorkflows = false, CancellationToken cancellationToken = default) + { + var parsedVersionOptions = versionOptions ?? VersionOptions.Latest; + var definition = (await store.FindManyAsync(new() + { + DefinitionId = definitionId, + VersionOptions = parsedVersionOptions + }, cancellationToken)).FirstOrDefault(); + + if (definition == null) + return null; + + if (includeConsumingWorkflows) + { + var definitions = await IncludeConsumersAsync([definition], cancellationToken); + return await ExportDefinitionsAsync(definitions, cancellationToken); + } + + return await ExportDefinitionAsync(definition, cancellationToken); + } + + /// + public async Task ExportManyAsync(ICollection ids, bool includeConsumingWorkflows = false, CancellationToken cancellationToken = default) + { + var definitions = (await store.FindManyAsync(new() + { + Ids = ids + }, cancellationToken)).ToList(); + + if (includeConsumingWorkflows) + definitions = await IncludeConsumersAsync(definitions, cancellationToken); + + if (definitions.Count == 0) + return null; + + return await ExportDefinitionsAsync(definitions, cancellationToken); + } + + /// + public async Task ExportDefinitionAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + { + var model = await workflowDefinitionMapper.MapAsync(definition, cancellationToken); + var binaryJson = await SerializeWorkflowDefinitionAsync(model, cancellationToken); + var fileName = GetFileName(model); + return new(binaryJson, fileName); + } + + /// + public async Task ExportDefinitionsAsync(ICollection definitions, CancellationToken cancellationToken = default) + { + var zipBytes = await CreateZipArchiveAsync(definitions, cancellationToken); + return new(zipBytes, "workflow-definitions.zip"); + } + + /// + /// Recursively discovers all consuming workflow definitions and includes them. + /// Consumers are always resolved at , regardless of the version used for the initial definitions. + /// + private async Task> IncludeConsumersAsync(List definitions, CancellationToken cancellationToken) + { + var initialDefinitionIds = definitions.Select(d => d.DefinitionId).ToList(); + var graph = await workflowReferenceGraphBuilder.BuildGraphAsync(initialDefinitionIds, cancellationToken); + + // Find any consumer definitions not already in our list. + var newDefinitionIds = graph.ConsumerDefinitionIds.Except(initialDefinitionIds).ToList(); + + if (newDefinitionIds.Count > 0) + { + var consumerDefinitions = await store.FindManyAsync(new() + { + DefinitionIds = newDefinitionIds.ToArray(), + VersionOptions = VersionOptions.Latest + }, cancellationToken); + + definitions = definitions.Concat(consumerDefinitions).ToList(); + } + + return definitions; + } + + private async Task CreateZipArchiveAsync(ICollection definitions, CancellationToken cancellationToken) + { + var zipStream = new MemoryStream(); + var sortedDefinitions = definitions.OrderBy(d => d.DefinitionId).ToList(); + + // NOTE: + // - ZIP timestamps cannot be earlier than 1980-01-01 (the ZIP format's minimum). + // - We intentionally use a fixed timestamp (instead of DateTimeOffset.UtcNow) to keep exports deterministic. + // This avoids producing different ZIP bytes for identical exports, which helps tests, caching, and diffing. + var zipEpoch = new DateTimeOffset(1980, 1, 1, 0, 0, 0, TimeSpan.Zero); + +#if NET10_0_OR_GREATER + await using (var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Create, true)) + { + foreach (var definition in sortedDefinitions) + { + var model = await workflowDefinitionMapper.MapAsync(definition, cancellationToken); + var binaryJson = await SerializeWorkflowDefinitionAsync(model, cancellationToken); + var fileName = GetFileName(model); + var entry = zipArchive.CreateEntry(fileName, CompressionLevel.Optimal); + entry.LastWriteTime = zipEpoch; + await using var entryStream = await entry.OpenAsync(cancellationToken); + await entryStream.WriteAsync(binaryJson, cancellationToken); + } + } +#else + using (var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Create, true)) + { + foreach (var definition in sortedDefinitions) + { + var model = await workflowDefinitionMapper.MapAsync(definition, cancellationToken); + var binaryJson = await SerializeWorkflowDefinitionAsync(model, cancellationToken); + var fileName = GetFileName(model); + var entry = zipArchive.CreateEntry(fileName, CompressionLevel.Optimal); + entry.LastWriteTime = zipEpoch; + await using var entryStream = entry.Open(); + await entryStream.WriteAsync(binaryJson, cancellationToken); + } + } +#endif + + zipStream.Position = 0; + return zipStream.ToArray(); + } + + private string GetFileName(WorkflowDefinitionModel definition) + { + var hasWorkflowName = !string.IsNullOrWhiteSpace(definition.Name); + var workflowName = hasWorkflowName ? definition.Name!.Trim() : definition.DefinitionId; + var workflowSlug = workflowName.Underscore().Dasherize().ToLowerInvariant(); + var dynamicFileNamePart = $"{workflowSlug}-{definition.DefinitionId}-{definition.Id}"; + var sanitizedDynamicFileNamePart = fileNameSanitizer.Sanitize(dynamicFileNamePart); + + return $"workflow-definition-{sanitizedDynamicFileNamePart}.json"; + } + + private async Task SerializeWorkflowDefinitionAsync(WorkflowDefinitionModel model, CancellationToken cancellationToken) + { + var serializerOptions = serializer.GetOptions(); + using var document = JsonSerializer.SerializeToDocument(model, serializerOptions); + var rootElement = document.RootElement; + + using var output = new MemoryStream(); + await using var writer = new Utf8JsonWriter(output); + + writer.WriteStartObject(); + writer.WriteString("$schema", "https://elsaworkflows.io/schemas/workflow-definition/v3.0.0/schema.json"); + + foreach (var property in rootElement.EnumerateObject()) + property.WriteTo(writer); + + writer.WriteEndObject(); + + await writer.FlushAsync(cancellationToken); + return output.ToArray(); + } +} diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Services/DefaultFileNameSanitizerTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Services/DefaultFileNameSanitizerTests.cs new file mode 100644 index 000000000..4fe3d9be1 --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Services/DefaultFileNameSanitizerTests.cs @@ -0,0 +1,35 @@ +using Elsa.Workflows.Management.Services; + +namespace Elsa.Workflows.Management.UnitTests.Services; + +public class DefaultFileNameSanitizerTests +{ + private readonly DefaultFileNameSanitizer _sut = new(); + + [Theory] + [InlineData("folder/child", "folder-child")] + [InlineData("folder\\child", "folder-child")] + [InlineData("already-safe", "already-safe")] + public void Sanitize_Replaces_Path_Separators_And_Preserves_Safe_Names(string input, string expected) + { + var result = _sut.Sanitize(input); + + Assert.Equal(expected, result); + } + + [Fact] + public void Sanitize_Replaces_Runtime_Invalid_File_Name_Characters() + { + var invalidChars = Path.GetInvalidFileNameChars(); + var hasNonSeparatorInvalidChar = invalidChars.Any(c => c is not '/' and not '\\'); + + if (!hasNonSeparatorInvalidChar) + return; + + var invalidCharacter = invalidChars.First(c => c is not '/' and not '\\'); + var result = _sut.Sanitize($"before{invalidCharacter}after"); + + Assert.Equal("before-after", result); + } +} + diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionExporterRegressionTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionExporterRegressionTests.cs new file mode 100644 index 000000000..f9748dc49 --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowDefinitionExporterRegressionTests.cs @@ -0,0 +1,80 @@ +using System.IO.Compression; +using System.Text.Json; +using Elsa.Expressions.Contracts; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Mappers; +using Elsa.Workflows.Management.Services; +using Elsa.Workflows.Models; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using NSubstitute; + +namespace Elsa.Workflows.Management.UnitTests.Services; + +public class WorkflowDefinitionExporterRegressionTests +{ + private readonly IWorkflowDefinitionStore _store = Substitute.For(); + private readonly IApiSerializer _serializer = Substitute.For(); + private readonly IWorkflowDefinitionService _workflowDefinitionService = Substitute.For(); + private readonly IWorkflowReferenceGraphBuilder _workflowReferenceGraphBuilder = Substitute.For(); + + [Fact] + public async Task ExportManyAsync_WithSlashInWorkflowName_CreatesFlatZipEntry() + { + var definition = new WorkflowDefinition + { + Id = "slash-name-workflow-v1", + DefinitionId = "slash-name-workflow", + Name = "folder/child", + CreatedAt = new(2025, 1, 1, 0, 0, 0, TimeSpan.Zero), + Version = 1, + IsLatest = true, + IsPublished = true, + MaterializerName = "Json" + }; + + _serializer.GetOptions().Returns(new JsonSerializerOptions(JsonSerializerDefaults.Web)); + _store.FindManyAsync(Arg.Any(), Arg.Any()) + .Returns(_ => Task.FromResult>([definition])); + _workflowDefinitionService.MaterializeWorkflowAsync(definition, Arg.Any()) + .Returns(Task.FromResult(CreateWorkflowGraph())); + + var sut = CreateExporter(); + var result = await sut.ExportManyAsync([definition.Id]); + + Assert.NotNull(result); + Assert.Equal("workflow-definitions.zip", result.FileName); + + using var zipStream = new MemoryStream(result.Data); + await using var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Read); + var entry = Assert.Single(zipArchive.Entries); + + Assert.Equal(entry.Name, entry.FullName); + Assert.DoesNotContain('/', entry.FullName); + Assert.DoesNotContain('\\', entry.FullName); + Assert.Equal("workflow-definition-folder-child-slash-name-workflow-slash-name-workflow-v1.json", entry.FullName); + } + + private WorkflowDefinitionExporter CreateExporter() + { + var activitySerializer = Substitute.For(); + var wellKnownTypeRegistry = Substitute.For(); + var scopeFactory = Substitute.For(); + var variableDefinitionMapper = new VariableDefinitionMapper(wellKnownTypeRegistry, scopeFactory, NullLogger.Instance); + var workflowDefinitionMapper = new WorkflowDefinitionMapper(activitySerializer, _workflowDefinitionService, variableDefinitionMapper); + var fileNameSanitizer = new DefaultFileNameSanitizer(); + + return new(_store, _serializer, workflowDefinitionMapper, _workflowReferenceGraphBuilder, fileNameSanitizer); + } + + private static WorkflowGraph CreateWorkflowGraph() + { + var root = new Sequence { Id = "root" }; + var rootNode = new ActivityNode(root, "Root"); + var workflow = new Workflow(root); + + return new(workflow, rootNode, [rootNode]); + } +} \ No newline at end of file diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowReferenceGraphBuilderTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowReferenceGraphBuilderTests.cs index df5332080..97682fe02 100644 --- a/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowReferenceGraphBuilderTests.cs +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Services/WorkflowReferenceGraphBuilderTests.cs @@ -232,4 +232,4 @@ public class WorkflowReferenceGraphBuilderTests foreach (var consumerId in consumerIds) Assert.Contains(consumerId, graph.ConsumerDefinitionIds); } -} +} \ No newline at end of file From 6bc0cf23ff75f123237a09d2979b3595b6810cee Mon Sep 17 00:00:00 2001 From: Cristian Gintili <48093965+cristiandolf@users.noreply.github.com> Date: Thu, 9 Apr 2026 11:44:33 +0200 Subject: [PATCH 4/7] fix: decouple background command execution from caller cancellation token (#7224) (#7283) * fix: decouple background command execution from caller cancellation token (fixes #7224) * Apply suggestion from @Copilot Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../BackgroundCommandSenderHostedService.cs | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs index 26916c831..d8aa66be5 100644 --- a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs @@ -78,17 +78,15 @@ public class BackgroundCommandSenderHostedService : BackgroundService using var scope = _scopeFactory.CreateScope(); var commandSender = scope.ServiceProvider.GetRequiredService(); - // Link the service cancellation token with the command's token to ensure proper cancellation - using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource( - cancellationToken, - commandContext.CancellationToken); - - // Process the command using the command sender service with the linked token + // Decouple from the caller's CancellationToken. + // We use only the service's cancellationToken (the background worker's lifetime) + // to ensure that dispatched workflows are processed even if the original + // HTTP request or triggering context has timed out. await commandSender.SendAsync( commandContext.Command, CommandStrategy.Default, commandContext.Headers, - linkedTokenSource.Token); + cancellationToken); } catch (Exception e) { From c0e37e38d08014ee8aaaf27aaabcc11f262638fb Mon Sep 17 00:00:00 2001 From: RalfvandenBurg Date: Sat, 11 Apr 2026 11:00:12 +0200 Subject: [PATCH 5/7] Fix/startuptask activate default tenant when empty (#7305) * fix: activate Tenant.Default when no tenants configured (Option C) Ensures IStartupTask implementations (e.g., PopulateRegistriesStartupTask, RunMigrationsStartupTask) run when multitenancy is enabled but the tenant provider returns an empty list. - In DefaultTenantService: treat empty provider response as [Tenant.Default] in GetTenantsDictionaryAsync (initial load) and RefreshAsync - Logic is internal to tenant service; no explicit call required Co-authored-by: Cursor * test: add DefaultTenantService tests for empty-provider fallback - ActivateTenantsAsync_WhenProviderReturnsEmpty_ActivatesDefaultTenant - ListAsync_WhenProviderReturnsEmpty_ReturnsDefaultTenant - ActivateTenantsAsync_WhenProviderReturnsTenants_ReturnsThoseTenants - RefreshAsync_WhenProviderChangesFromTenantsToEmpty_KeepsDefaultTenant Co-authored-by: Cursor * PR feedback disposes serviceprovider also --------- Co-authored-by: Ralf Co-authored-by: Cursor --- .../Multitenancy/Contracts/ITenantService.cs | 1 + .../Implementations/DefaultTenantService.cs | 9 +- .../Multitenancy/DefaultTenantServiceTests.cs | 145 ++++++++++++++++++ 3 files changed, 153 insertions(+), 2 deletions(-) create mode 100644 test/unit/Elsa.Common.UnitTests/Multitenancy/DefaultTenantServiceTests.cs diff --git a/src/modules/Elsa.Common/Multitenancy/Contracts/ITenantService.cs b/src/modules/Elsa.Common/Multitenancy/Contracts/ITenantService.cs index 88b1e44fb..7abc3f042 100644 --- a/src/modules/Elsa.Common/Multitenancy/Contracts/ITenantService.cs +++ b/src/modules/Elsa.Common/Multitenancy/Contracts/ITenantService.cs @@ -63,6 +63,7 @@ public interface ITenantService /// /// Invokes the and caches the result. /// When new tenants are added, lifecycle events are triggered to ensure background tasks are updated. + /// When the provider returns an empty list, is activated so that startup tasks run. /// Task RefreshAsync(CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Common/Multitenancy/Implementations/DefaultTenantService.cs b/src/modules/Elsa.Common/Multitenancy/Implementations/DefaultTenantService.cs index 6a1f9cb98..5dc7f0c05 100644 --- a/src/modules/Elsa.Common/Multitenancy/Implementations/DefaultTenantService.cs +++ b/src/modules/Elsa.Common/Multitenancy/Implementations/DefaultTenantService.cs @@ -77,7 +77,10 @@ public class DefaultTenantService(IServiceScopeFactory scopeFactory, ITenantScop var tenantsProvider = scope.ServiceProvider.GetRequiredService(); var currentTenants = await GetTenantsDictionaryAsync(cancellationToken); var currentTenantIds = currentTenants.Keys; - var newTenants = (await tenantsProvider.ListAsync(cancellationToken)).ToDictionary(x => x.Id.EmptyIfNull()); + var tenantsFromProvider = (await tenantsProvider.ListAsync(cancellationToken)).ToList(); + var newTenants = tenantsFromProvider.Count == 0 + ? new Dictionary { [Tenant.DefaultTenantId] = Tenant.Default } + : tenantsFromProvider.ToDictionary(x => x.Id.EmptyIfNull()); var newTenantIds = newTenants.Keys; var removedTenantIds = currentTenantIds.Except(newTenantIds).ToArray(); var addedTenantIds = newTenantIds.Except(currentTenantIds).ToArray(); @@ -112,7 +115,9 @@ public class DefaultTenantService(IServiceScopeFactory scopeFactory, ITenantScop _tenantsDictionary = new Dictionary(); _tenantScopesDictionary = new Dictionary(); var tenantsProvider = _serviceScope.ServiceProvider.GetRequiredService(); - var tenants = await tenantsProvider.ListAsync(cancellationToken); + var tenants = (await tenantsProvider.ListAsync(cancellationToken)).ToList(); + if (tenants.Count == 0) + tenants = [Tenant.Default]; foreach (var tenant in tenants) await RegisterTenantAsync(tenant, cancellationToken); diff --git a/test/unit/Elsa.Common.UnitTests/Multitenancy/DefaultTenantServiceTests.cs b/test/unit/Elsa.Common.UnitTests/Multitenancy/DefaultTenantServiceTests.cs new file mode 100644 index 000000000..89685a0d5 --- /dev/null +++ b/test/unit/Elsa.Common.UnitTests/Multitenancy/DefaultTenantServiceTests.cs @@ -0,0 +1,145 @@ +using Elsa.Common.Multitenancy; +using Elsa.Common.Multitenancy.EventHandlers; +using Elsa.Common.RecurringTasks; +using Microsoft.Extensions.DependencyInjection; +using NSubstitute; + +namespace Elsa.Common.UnitTests.Multitenancy; + +/// +/// Tests for , including the fallback to activate +/// when the tenant provider returns an empty list. +/// +public class DefaultTenantServiceTests +{ + [Fact] + public async Task ActivateTenantsAsync_WhenProviderReturnsEmpty_ActivatesDefaultTenant() + { + // Arrange - provider returns no tenants + var (tenantService, serviceProvider) = await CreateTenantServiceAsync(Array.Empty()); + + try + { + // Act + await tenantService.ActivateTenantsAsync(); + + // Assert + var tenants = (await tenantService.ListAsync()).ToList(); + Assert.Single(tenants); + Assert.Same(Tenant.Default, tenants[0]); + Assert.Equal(Tenant.DefaultTenantId, tenants[0].Id); + } + finally + { + if (tenantService is IAsyncDisposable disposable) + await disposable.DisposeAsync(); + await serviceProvider.DisposeAsync(); + } + } + + [Fact] + public async Task ListAsync_WhenProviderReturnsEmpty_ReturnsDefaultTenant() + { + // Arrange - ListAsync triggers initialization when provider returns empty + var (tenantService, serviceProvider) = await CreateTenantServiceAsync(Array.Empty()); + + try + { + // Act + var tenants = (await tenantService.ListAsync()).ToList(); + + // Assert + Assert.Single(tenants); + Assert.Same(Tenant.Default, tenants[0]); + } + finally + { + if (tenantService is IAsyncDisposable disposable) + await disposable.DisposeAsync(); + await serviceProvider.DisposeAsync(); + } + } + + [Fact] + public async Task ActivateTenantsAsync_WhenProviderReturnsTenants_ReturnsThoseTenants() + { + // Arrange - provider returns specific tenants + var tenant1 = new Tenant { Id = "tenant-1", Name = "Tenant 1" }; + var tenant2 = new Tenant { Id = "tenant-2", Name = "Tenant 2" }; + var (tenantService, serviceProvider) = await CreateTenantServiceAsync([tenant1, tenant2]); + + try + { + // Act + await tenantService.ActivateTenantsAsync(); + + // Assert - should not use Tenant.Default fallback + var tenants = (await tenantService.ListAsync()).ToList(); + Assert.Equal(2, tenants.Count); + Assert.Contains(tenants, t => t.Id == "tenant-1"); + Assert.Contains(tenants, t => t.Id == "tenant-2"); + } + finally + { + if (tenantService is IAsyncDisposable disposable) + await disposable.DisposeAsync(); + await serviceProvider.DisposeAsync(); + } + } + + [Fact] + public async Task RefreshAsync_WhenProviderChangesFromTenantsToEmpty_KeepsDefaultTenant() + { + // Arrange - start with tenants, then provider returns empty (simulating config change) + var tenant1 = new Tenant { Id = "tenant-1", Name = "Tenant 1" }; + var providerReturns = new List { tenant1 }; + var (tenantService, serviceProvider) = await CreateTenantServiceAsync(providerReturns, () => providerReturns); + + try + { + await tenantService.ActivateTenantsAsync(); + Assert.Single(await tenantService.ListAsync()); + + // Simulate provider now returning empty (e.g., config removed all tenants) + providerReturns.Clear(); + + // Act + await tenantService.RefreshAsync(); + + // Assert - should fall back to Tenant.Default instead of having zero tenants + var tenants = (await tenantService.ListAsync()).ToList(); + Assert.Single(tenants); + Assert.Same(Tenant.Default, tenants[0]); + } + finally + { + if (tenantService is IAsyncDisposable disposable) + await disposable.DisposeAsync(); + await serviceProvider.DisposeAsync(); + } + } + + private static Task<(ITenantService TenantService, ServiceProvider ServiceProvider)> CreateTenantServiceAsync(IEnumerable tenants, Func>? tenantsFactory = null) + { + var tenantList = tenants.ToList(); + var getTenants = tenantsFactory ?? (() => tenantList); + + var tenantsProvider = Substitute.For(); + tenantsProvider.ListAsync(Arg.Any()).Returns(_ => getTenants()); + + var services = new ServiceCollection(); + services.AddSingleton(_ => tenantsProvider); + services.AddSingleton(); + services.AddSingleton(); + services.AddSingleton(Substitute.For()); + services.AddSingleton(Substitute.For()); + services.AddSingleton(Substitute.For()); + services.AddSingleton(); + services.AddSingleton(); + services.AddSingleton(); + services.AddLogging(); + + var serviceProvider = services.BuildServiceProvider(); + return Task.FromResult((serviceProvider.GetRequiredService(), serviceProvider)); + } +} From 3a3766861964b685f4eab23086138836fd775272 Mon Sep 17 00:00:00 2001 From: RalfvandenBurg Date: Wed, 15 Apr 2026 14:06:53 +0200 Subject: [PATCH 6/7] fix: correct workflow summary pagination on release 3.6.1 (#7392) Co-authored-by: Ralf --- src/modules/Elsa.Common/Models/PageArgs.cs | 2 +- .../Models/PageArgsTests.cs | 39 ++++++++++++ .../WorkflowInstanceStoreExtensionsTests.cs | 63 +++++++++++++++++++ 3 files changed, 103 insertions(+), 1 deletion(-) create mode 100644 test/unit/Elsa.Common.UnitTests/Models/PageArgsTests.cs create mode 100644 test/unit/Elsa.Workflows.Runtime.UnitTests/Extensions/WorkflowInstanceStoreExtensionsTests.cs diff --git a/src/modules/Elsa.Common/Models/PageArgs.cs b/src/modules/Elsa.Common/Models/PageArgs.cs index 054c0dbfe..b1d60a09d 100644 --- a/src/modules/Elsa.Common/Models/PageArgs.cs +++ b/src/modules/Elsa.Common/Models/PageArgs.cs @@ -67,5 +67,5 @@ public record PageArgs /// Returns pagination arguments for the next page. /// /// The arguments for the next page. - public PageArgs Next() => this with { Offset = Page + 1 }; + public PageArgs Next() => this with { Offset = Offset + Limit }; } \ No newline at end of file diff --git a/test/unit/Elsa.Common.UnitTests/Models/PageArgsTests.cs b/test/unit/Elsa.Common.UnitTests/Models/PageArgsTests.cs new file mode 100644 index 000000000..c76fc05aa --- /dev/null +++ b/test/unit/Elsa.Common.UnitTests/Models/PageArgsTests.cs @@ -0,0 +1,39 @@ +using Elsa.Common.Models; + +namespace Elsa.Common.UnitTests.Models; + +public class PageArgsTests +{ + [Fact] + public void Next_FromPage_AdvancesByPageSize() + { + // Arrange + var pageArgs = PageArgs.FromPage(0, 10); + + // Act + var nextPage = pageArgs.Next(); + var thirdPage = nextPage.Next(); + + // Assert + Assert.Equal(10, nextPage.Offset); + Assert.Equal(10, nextPage.Limit); + Assert.Equal(1, nextPage.Page); + Assert.Equal(20, thirdPage.Offset); + Assert.Equal(2, thirdPage.Page); + } + + [Fact] + public void Next_AllPages_RemainsUnbounded() + { + // Arrange + var pageArgs = PageArgs.All; + + // Act + var nextPage = pageArgs.Next(); + + // Assert + Assert.Null(nextPage.Offset); + Assert.Null(nextPage.Limit); + Assert.Null(nextPage.Page); + } +} \ No newline at end of file diff --git a/test/unit/Elsa.Workflows.Runtime.UnitTests/Extensions/WorkflowInstanceStoreExtensionsTests.cs b/test/unit/Elsa.Workflows.Runtime.UnitTests/Extensions/WorkflowInstanceStoreExtensionsTests.cs new file mode 100644 index 000000000..d585bd8ef --- /dev/null +++ b/test/unit/Elsa.Workflows.Runtime.UnitTests/Extensions/WorkflowInstanceStoreExtensionsTests.cs @@ -0,0 +1,63 @@ +using Elsa.Common.Models; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using NSubstitute; + +namespace Elsa.Workflows.Runtime.UnitTests.Extensions; + +public class WorkflowInstanceStoreExtensionsTests +{ + [Fact] + public async Task EnumerateSummariesAsync_AdvancesAcrossPagesWithoutRepeatingOffsets() + { + // Arrange + var store = Substitute.For(); + var filter = new WorkflowInstanceFilter(); + var requestedPages = new List(); + var page1 = new List + { + CreateSummary("workflow-1"), + CreateSummary("workflow-2") + }; + var page2 = new List + { + CreateSummary("workflow-3"), + CreateSummary("workflow-4") + }; + + store + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(callInfo => + { + var pageArgs = callInfo.ArgAt(1); + requestedPages.Add(pageArgs with { }); + + var items = pageArgs.Offset switch + { + 0 => page1, + 2 => page2, + _ => new List() + }; + + return new ValueTask>(Page.Of(items, page1.Count + page2.Count)); + }); + + // Act + var results = new List(); + await foreach (var workflowInstance in store.EnumerateSummariesAsync(filter, 2, CancellationToken.None)) + results.Add(workflowInstance); + + // Assert + Assert.Equal(["workflow-1", "workflow-2", "workflow-3", "workflow-4"], results.Select(x => x.Id).ToArray()); + Assert.Equal([(int?)0, 2, 4], requestedPages.Select(x => x.Offset).ToArray()); + Assert.All(requestedPages, x => Assert.Equal(2, x.Limit)); + } + + private static WorkflowInstanceSummary CreateSummary(string id) => new() + { + Id = id, + DefinitionId = "definition", + DefinitionVersionId = "definition:1" + }; +} \ No newline at end of file From 793ae15bf737c50fdb096e093d646fecd95b725d Mon Sep 17 00:00:00 2001 From: RalfvandenBurg Date: Wed, 15 Apr 2026 14:11:26 +0200 Subject: [PATCH 7/7] Fix tenant scope for interrupted workflow restarts (#7389) * Fix tenant scope for interrupted workflow restarts * Update src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> * Handle per-instance restart failures * fix: make restart interrupted workflows tenant services optional --------- Co-authored-by: Ralf Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> --- .../Models/WorkflowInstanceSummary.cs | 5 + .../Tasks/RestartInterruptedWorkflowsTask.cs | 28 ++- .../RestartInterruptedWorkflowsTaskTests.cs | 204 ++++++++++++++++++ 3 files changed, 234 insertions(+), 3 deletions(-) create mode 100644 test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs index ef673127a..a9ba30ba1 100644 --- a/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowInstanceSummary.cs @@ -16,6 +16,7 @@ public class WorkflowInstanceSummary return new() { Id = workflowInstance.Id, + TenantId = workflowInstance.TenantId, DefinitionId = workflowInstance.DefinitionId, DefinitionVersionId = workflowInstance.DefinitionVersionId, Version = workflowInstance.Version, @@ -36,6 +37,7 @@ public class WorkflowInstanceSummary public static Expression> FromInstanceExpression() => workflowInstance => new() { Id = workflowInstance.Id, + TenantId = workflowInstance.TenantId, DefinitionId = workflowInstance.DefinitionId, DefinitionVersionId = workflowInstance.DefinitionVersionId, Version = workflowInstance.Version, @@ -53,6 +55,9 @@ public class WorkflowInstanceSummary /// The ID of the workflow instance. public string Id { get; set; } = null!; + /// The ID of the tenant that owns the workflow instance. + public string? TenantId { get; set; } + /// The ID of the workflow definition. public string DefinitionId { get; set; } = null!; diff --git a/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs b/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs index 5cb2cc2ba..288edb611 100644 --- a/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs +++ b/src/modules/Elsa.Workflows.Runtime/Tasks/RestartInterruptedWorkflowsTask.cs @@ -1,4 +1,5 @@ using Elsa.Common; +using Elsa.Common.Multitenancy; using Elsa.Common.RecurringTasks; using Elsa.Workflows.Management; using Elsa.Workflows.Management.Filters; @@ -12,11 +13,13 @@ namespace Elsa.Workflows.Runtime.Tasks; [SingleNodeTask] [UsedImplicitly] public class RestartInterruptedWorkflowsTask( - IWorkflowInstanceStore workflowInstanceStore, IWorkflowRestarter workflowRestarter, + IWorkflowInstanceStore workflowInstanceStore, + ILogger logger, IOptions options, ISystemClock systemClock, - ILogger logger) : RecurringTask + ITenantService? tenantService = null, + ITenantAccessor? tenantAccessor = null) : RecurringTask { /// public override async Task ExecuteAsync(CancellationToken cancellationToken) @@ -28,7 +31,26 @@ public class RestartInterruptedWorkflowsTask( logger.LogInformation("Restarting interrupted workflows."); await foreach (var workflowInstance in workflowInstances) { - await workflowRestarter.RestartWorkflowAsync(workflowInstance.Id, cancellationToken: cancellationToken); + try + { + var tenantId = workflowInstance.TenantId ?? string.Empty; + + if (tenantService != null && tenantAccessor != null && !string.IsNullOrWhiteSpace(tenantId) && tenantId != Tenant.AgnosticTenantId) + { + var tenant = await tenantService.FindAsync(tenantId, cancellationToken) ?? new Tenant { Id = tenantId, Name = tenantId }; + + using (tenantAccessor.PushContext(tenant)) + await workflowRestarter.RestartWorkflowAsync(workflowInstance.Id, cancellationToken: cancellationToken); + + continue; + } + + await workflowRestarter.RestartWorkflowAsync(workflowInstance.Id, cancellationToken: cancellationToken); + } + catch (Exception ex) + { + logger.LogError(ex, "Failed to restart interrupted workflow {WorkflowInstanceId}", workflowInstance.Id); + } } logger.LogInformation("Finished restarting interrupted workflows."); } diff --git a/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs new file mode 100644 index 000000000..2b23f2c70 --- /dev/null +++ b/test/unit/Elsa.Workflows.Runtime.UnitTests/Services/RestartInterruptedWorkflowsTaskTests.cs @@ -0,0 +1,204 @@ +using Elsa.Common; +using Elsa.Common.Models; +using Elsa.Common.Multitenancy; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Runtime.Options; +using Elsa.Workflows.Runtime.Tasks; +using Microsoft.Extensions.Logging; +using NSubstitute; + +namespace Elsa.Workflows.Runtime.UnitTests.Services; + +public class RestartInterruptedWorkflowsTaskTests +{ + [Fact(DisplayName = "ExecuteAsync restarts each workflow within its tenant context")] + public async Task ExecuteAsync_RestartsWithinWorkflowTenantContext() + { + var workflowInstanceStore = Substitute.For(); + var tenantAccessor = new DefaultTenantAccessor(); + var tenantService = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", "tenant-a", now), + CreateWorkflowInstance("workflow-2", "tenant-b", now) + }; + var observedTenantIds = new List(); + + clock.UtcNow.Returns(now); + tenantService.FindAsync("tenant-a", Arg.Any()).Returns(new Tenant { Id = "tenant-a", Name = "Tenant A" }); + tenantService.FindAsync("tenant-b", Arg.Any()).Returns(new Tenant { Id = "tenant-b", Name = "Tenant B" }); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + workflowRestarter + .When(x => x.RestartWorkflowAsync(Arg.Any(), Arg.Any())) + .Do(_ => observedTenantIds.Add(tenantAccessor.Tenant?.Id)); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock, tenantService, tenantAccessor); + + await task.ExecuteAsync(CancellationToken.None); + + Assert.Equal(new[] { "tenant-a", "tenant-b" }, observedTenantIds); + Assert.Null(tenantAccessor.Tenant); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any()); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-2", Arg.Any()); + } + + [Theory(DisplayName = "ExecuteAsync does not push tenant context for default or agnostic tenant instances")] + [InlineData(null)] + [InlineData("")] + [InlineData("*")] + public async Task ExecuteAsync_DefaultOrAgnosticTenant_DoesNotPushContext(string? tenantId) + { + var workflowInstanceStore = Substitute.For(); + var tenantAccessor = new DefaultTenantAccessor(); + var tenantService = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", tenantId, now) + }; + var observedTenantIds = new List(); + + clock.UtcNow.Returns(now); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + workflowRestarter + .When(x => x.RestartWorkflowAsync(Arg.Any(), Arg.Any())) + .Do(_ => observedTenantIds.Add(tenantAccessor.Tenant?.Id)); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock, tenantService, tenantAccessor); + + await task.ExecuteAsync(CancellationToken.None); + + Assert.Equal(new string?[] { null }, observedTenantIds); + Assert.Null(tenantAccessor.Tenant); + await tenantService.DidNotReceive().FindAsync(Arg.Any(), Arg.Any()); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any()); + } + + [Fact(DisplayName = "ExecuteAsync continues restarting later workflows after a failure")] + public async Task ExecuteAsync_ContinuesAfterFailure() + { + var workflowInstanceStore = Substitute.For(); + var tenantAccessor = new DefaultTenantAccessor(); + var tenantService = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", "tenant-a", now), + CreateWorkflowInstance("workflow-2", "tenant-b", now) + }; + + clock.UtcNow.Returns(now); + tenantService.FindAsync("tenant-a", Arg.Any()).Returns(_ => Task.FromException(new InvalidOperationException("Transient tenant lookup failure"))); + tenantService.FindAsync("tenant-b", Arg.Any()).Returns(new Tenant { Id = "tenant-b", Name = "Tenant B" }); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock, tenantService, tenantAccessor); + + await task.ExecuteAsync(CancellationToken.None); + + await workflowRestarter.DidNotReceive().RestartWorkflowAsync("workflow-1", Arg.Any()); + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-2", Arg.Any()); + } + + [Fact(DisplayName = "ExecuteAsync restarts tenant-specific workflows when tenant services are unavailable")] + public async Task ExecuteAsync_WithoutTenantServices_RestartsWorkflow() + { + var workflowInstanceStore = Substitute.For(); + var workflowRestarter = Substitute.For(); + var clock = Substitute.For(); + var logger = Substitute.For>(); + var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z"); + var workflowInstances = new List + { + CreateWorkflowInstance("workflow-1", "tenant-a", now) + }; + + clock.UtcNow.Returns(now); + workflowInstanceStore + .SummarizeManyAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns( + callInfo => + { + var pageArgs = callInfo.ArgAt(1); + var items = pageArgs.Offset == 0 ? workflowInstances : new List(); + return new ValueTask>(Page.Of(items, workflowInstances.Count)); + }); + + var options = Microsoft.Extensions.Options.Options.Create(new RuntimeOptions + { + RestartInterruptedWorkflowsBatchSize = 10, + InactivityThreshold = TimeSpan.FromMinutes(5) + }); + var task = new RestartInterruptedWorkflowsTask(workflowRestarter, workflowInstanceStore, logger, options, clock); + + await task.ExecuteAsync(CancellationToken.None); + + await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any()); + } + + private static WorkflowInstanceSummary CreateWorkflowInstance(string id, string? tenantId, DateTimeOffset updatedAt) + { + return new WorkflowInstanceSummary + { + Id = id, + TenantId = tenantId, + DefinitionId = "definition", + DefinitionVersionId = "definition:1", + Status = WorkflowStatus.Running, + SubStatus = WorkflowSubStatus.Pending, + CreatedAt = updatedAt.AddMinutes(-10), + UpdatedAt = updatedAt.AddMinutes(-10) + }; + } +} \ No newline at end of file