diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index 4257d4c5d..0fbb38c00 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -223,7 +223,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/doc/changelogs/3.6.0.md b/doc/changelogs/3.6.0.md index 8fee7186e..53889153e 100644 --- a/doc/changelogs/3.6.0.md +++ b/doc/changelogs/3.6.0.md @@ -8,6 +8,8 @@ Compare: [`3.5.3...3.6.0`](https://github.com/elsa-workflows/elsa-core/compare/3 - **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)) + - **Scripting package names have changed**: - Elsa.JavaScript -> Elsa.Expressions.JavaScript - Elsa.CSharp -> Elsa.Expressions.CSharp @@ -150,9 +152,14 @@ Compare: [`3.5.3...3.6.0`](https://github.com/elsa-workflows/elsa-core/compare/3 ## 📦 Full changelog -- **Trigger serialization fix** (`WorkflowTriggerEqualityComparer`): avoid type discriminator injection during equality comparison. ([6f53c26f24](https://github.com/elsa-workflows/elsa-core/commit/6f53c26f24)) -- **`WorkflowExecutionContext`**: removed unused `ClearCompletionCallbacks` method; `correlationId` parameter made non-optional in the constructor. ([fa798b0a47](https://github.com/elsa-workflows/elsa-core/commit/fa798b0a47), [79a64e90fd](https://github.com/elsa-workflows/elsa-core/commit/79a64e90fd)) -- **`AttributeUsage` targets** updated for `InputAttribute`, `OutputAttribute`, and `ActivityAttribute` to reflect correct usage scenarios. ([3778a14e54](https://github.com/elsa-workflows/elsa-core/commit/3778a14e54)) -- **`ClrWorkflowsProvider`** refactored to remove tenant prefix logic, relying solely on the `TenantId` property. ([41409b156d](https://github.com/elsa-workflows/elsa-core/commit/41409b156d)) -- Removed `ConvertNullTenantIdToEmptyString` migration and its designer file. ([ccd8268413](https://github.com/elsa-workflows/elsa-core/commit/ccd8268413)) -- Tenant isolation enforced in `WorkflowDefinitionActivityProvider`; tenant ID included in activity type cache key. ([558902bb77](https://github.com/elsa-workflows/elsa-core/commit/558902bb77)) +* **Trigger serialization fix** (`WorkflowTriggerEqualityComparer`): avoid type discriminator injection during equality comparison. ([6f53c26f24](https://github.com/elsa-workflows/elsa-core/commit/6f53c26f24)) +* **`WorkflowExecutionContext`**: removed unused `ClearCompletionCallbacks` method; `correlationId` parameter made non-optional in the constructor. ([fa798b0a47](https://github.com/elsa-workflows/elsa-core/commit/fa798b0a47), [79a64e90fd](https://github.com/elsa-workflows/elsa-core/commit/79a64e90fd)) +* **`AttributeUsage` targets** updated for `InputAttribute`, `OutputAttribute`, and `ActivityAttribute` to reflect correct usage scenarios. ([3778a14e54](https://github.com/elsa-workflows/elsa-core/commit/3778a14e54)) +* **`ClrWorkflowsProvider`** refactored to remove tenant prefix logic, relying solely on the `TenantId` property. ([41409b156d](https://github.com/elsa-workflows/elsa-core/commit/41409b156d)) +* Removed `ConvertNullTenantIdToEmptyString` migration and its designer file. ([ccd8268413](https://github.com/elsa-workflows/elsa-core/commit/ccd8268413)) +* Tenant isolation enforced in `WorkflowDefinitionActivityProvider`; tenant ID included in activity type cache key. ([558902bb77](https://github.com/elsa-workflows/elsa-core/commit/558902bb77)) +* Removed legacy API key and service management functionality from server projects. ([e6899ac32c](https://github.com/elsa-workflows/elsa-core/commit/e6899ac32c)) +* 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 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) { 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/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/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 eb5f2b3bd..21ae2edfd 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -273,6 +273,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/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.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/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.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.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)); + } +} 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 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 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