Merge remote-tracking branch 'origin/release/3.6.1'

This commit is contained in:
Sipke Schoorstra 2026-04-20 15:08:10 +02:00
commit 3d8d3b7de2
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
22 changed files with 926 additions and 201 deletions

View file

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

View file

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

View file

@ -78,17 +78,15 @@ public class BackgroundCommandSenderHostedService : BackgroundService
using var scope = _scopeFactory.CreateScope();
var commandSender = scope.ServiceProvider.GetRequiredService<ICommandSender>();
// 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)
{

View file

@ -67,5 +67,5 @@ public record PageArgs
/// Returns pagination arguments for the next page.
/// </summary>
/// <returns>The arguments for the next page.</returns>
public PageArgs Next() => this with { Offset = Page + 1 };
public PageArgs Next() => this with { Offset = Offset + Limit };
}

View file

@ -63,6 +63,7 @@ public interface ITenantService
/// <summary>
/// Invokes the <see cref="ITenantsProvider"/> 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, <see cref="Tenant.Default"/> is activated so that startup tasks run.
/// </summary>
Task RefreshAsync(CancellationToken cancellationToken = default);
}

View file

@ -77,7 +77,10 @@ public class DefaultTenantService(IServiceScopeFactory scopeFactory, ITenantScop
var tenantsProvider = scope.ServiceProvider.GetRequiredService<ITenantsProvider>();
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<string, Tenant> { [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<string, Tenant>();
_tenantScopesDictionary = new Dictionary<Tenant, TenantScope>();
var tenantsProvider = _serviceScope.ServiceProvider.GetRequiredService<ITenantsProvider>();
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);

View file

@ -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.
/// </summary>
[UsedImplicitly]
internal class Export : ElsaEndpoint<Request>
internal class Export(IWorkflowDefinitionExporter exporter) : ElsaEndpoint<Request>
{
private readonly IApiSerializer _serializer;
private readonly IWorkflowDefinitionStore _store;
private readonly IWorkflowReferenceGraphBuilder _workflowReferenceGraphBuilder;
private readonly WorkflowDefinitionMapper _workflowDefinitionMapper;
/// <inheritdoc />
public Export(
IWorkflowDefinitionStore store,
IApiSerializer serializer,
WorkflowDefinitionMapper workflowDefinitionMapper,
IWorkflowReferenceGraphBuilder workflowReferenceGraphBuilder)
{
_store = store;
_serializer = serializer;
_workflowDefinitionMapper = workflowDefinitionMapper;
_workflowReferenceGraphBuilder = workflowReferenceGraphBuilder;
}
/// <inheritdoc />
public override void Configure()
{
@ -48,165 +23,33 @@ internal class Export : ElsaEndpoint<Request>
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<string> 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);
}
/// <summary>
/// Recursively discovers all consuming workflow definitions and includes them.
/// Consumers are always resolved at <see cref="VersionOptions.Latest"/>, regardless of the version used for the initial definitions.
/// </summary>
private async Task<List<WorkflowDefinition>> IncludeConsumersAsync(List<WorkflowDefinition> 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<WorkflowDefinition> 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<byte[]> 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<WorkflowDefinitionModel> CreateWorkflowModelAsync(WorkflowDefinition definition, CancellationToken cancellationToken)
{
return await _workflowDefinitionMapper.MapAsync(definition, cancellationToken);
}
}

View file

@ -0,0 +1,13 @@
namespace Elsa.Workflows.Management;
/// <summary>
/// Sanitizes file names by replacing invalid characters with safe alternatives.
/// </summary>
public interface IFileNameSanitizer
{
/// <summary>
/// Replaces invalid file name characters in the specified value.
/// </summary>
string Sanitize(string value);
}

View file

@ -0,0 +1,46 @@
using Elsa.Common.Models;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Models;
namespace Elsa.Workflows.Management;
/// <summary>
/// Exports workflow definitions as serialized JSON or ZIP archives.
/// </summary>
public interface IWorkflowDefinitionExporter
{
/// <summary>
/// Exports a single workflow definition as a JSON byte array, optionally including consuming workflows as a ZIP archive.
/// </summary>
/// <param name="definitionId">The definition ID.</param>
/// <param name="versionOptions">The version options. Defaults to <see cref="VersionOptions.Latest"/>.</param>
/// <param name="includeConsumingWorkflows">When true, includes all consuming workflow definitions in the export as a ZIP archive.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>An export result containing the binary content and a suggested file name, or null if the definition was not found.</returns>
Task<WorkflowDefinitionExportResult?> ExportAsync(string definitionId, VersionOptions? versionOptions = null, bool includeConsumingWorkflows = false, CancellationToken cancellationToken = default);
/// <summary>
/// Exports multiple workflow definitions as a ZIP archive.
/// </summary>
/// <param name="ids">A list of workflow definition version IDs.</param>
/// <param name="includeConsumingWorkflows">When true, includes all consuming workflow definitions in the export.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>An export result containing the ZIP binary content and a suggested file name, or null if no definitions were found.</returns>
Task<WorkflowDefinitionExportResult?> ExportManyAsync(ICollection<string> ids, bool includeConsumingWorkflows = false, CancellationToken cancellationToken = default);
/// <summary>
/// Serializes a single workflow definition entity to a JSON byte array (with $schema header).
/// </summary>
/// <param name="definition">The workflow definition entity.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The serialized JSON bytes and suggested file name.</returns>
Task<WorkflowDefinitionExportResult> ExportDefinitionAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default);
/// <summary>
/// Exports multiple workflow definition entities as a ZIP archive.
/// </summary>
/// <param name="definitions">The workflow definitions to export.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The ZIP archive bytes and suggested file name.</returns>
Task<WorkflowDefinitionExportResult> ExportDefinitionsAsync(ICollection<WorkflowDefinition> definitions, CancellationToken cancellationToken = default);
}

View file

@ -273,6 +273,8 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module)
.AddScoped<IWorkflowReferenceGraphBuilder, WorkflowReferenceGraphBuilder>()
.AddScoped(_workflowDefinitionPublisher)
.AddScoped<IWorkflowDefinitionImporter, WorkflowDefinitionImporter>()
.AddScoped<IWorkflowDefinitionExporter, WorkflowDefinitionExporter>()
.AddSingleton<IFileNameSanitizer, DefaultFileNameSanitizer>()
.AddScoped<IWorkflowDefinitionManager, WorkflowDefinitionManager>()
.AddScoped<IWorkflowInstanceManager, WorkflowInstanceManager>()
.AddScoped<IWorkflowReferenceUpdater, WorkflowReferenceUpdater>()

View file

@ -0,0 +1,8 @@
namespace Elsa.Workflows.Management.Models;
/// <summary>
/// Represents the result of exporting one or more workflow definitions.
/// </summary>
/// <param name="Data">The binary content (JSON or ZIP).</param>
/// <param name="FileName">The suggested file name.</param>
public record WorkflowDefinitionExportResult(byte[] Data, string FileName);

View file

@ -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<Func<WorkflowInstance, WorkflowInstanceSummary>> 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
/// <summary>The ID of the workflow instance.</summary>
public string Id { get; set; } = null!;
/// <summary>The ID of the tenant that owns the workflow instance.</summary>
public string? TenantId { get; set; }
/// <summary>The ID of the workflow definition.</summary>
public string DefinitionId { get; set; } = null!;

View file

@ -0,0 +1,33 @@
namespace Elsa.Workflows.Management.Services;
/// <inheritdoc />
public class DefaultFileNameSanitizer : IFileNameSanitizer
{
private static readonly char[] InvalidFileNameCharacters = Path.GetInvalidFileNameChars();
/// <inheritdoc />
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;
}

View file

@ -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;
/// <inheritdoc />
public class WorkflowDefinitionExporter(
IWorkflowDefinitionStore store,
IApiSerializer serializer,
WorkflowDefinitionMapper workflowDefinitionMapper,
IWorkflowReferenceGraphBuilder workflowReferenceGraphBuilder,
IFileNameSanitizer fileNameSanitizer) : IWorkflowDefinitionExporter
{
/// <inheritdoc />
public async Task<WorkflowDefinitionExportResult?> 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);
}
/// <inheritdoc />
public async Task<WorkflowDefinitionExportResult?> ExportManyAsync(ICollection<string> 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);
}
/// <inheritdoc />
public async Task<WorkflowDefinitionExportResult> 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);
}
/// <inheritdoc />
public async Task<WorkflowDefinitionExportResult> ExportDefinitionsAsync(ICollection<WorkflowDefinition> definitions, CancellationToken cancellationToken = default)
{
var zipBytes = await CreateZipArchiveAsync(definitions, cancellationToken);
return new(zipBytes, "workflow-definitions.zip");
}
/// <summary>
/// Recursively discovers all consuming workflow definitions and includes them.
/// Consumers are always resolved at <see cref="VersionOptions.Latest"/>, regardless of the version used for the initial definitions.
/// </summary>
private async Task<List<WorkflowDefinition>> IncludeConsumersAsync(List<WorkflowDefinition> 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<byte[]> CreateZipArchiveAsync(ICollection<WorkflowDefinition> 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<byte[]> 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();
}
}

View file

@ -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<RestartInterruptedWorkflowsTask> logger,
IOptions<RuntimeOptions> options,
ISystemClock systemClock,
ILogger<RestartInterruptedWorkflowsTask> logger) : RecurringTask
ITenantService? tenantService = null,
ITenantAccessor? tenantAccessor = null) : RecurringTask
{
/// <inheritdoc />
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.");
}

View file

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

View file

@ -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;
/// <summary>
/// Tests for <see cref="DefaultTenantService"/>, including the fallback to activate
/// <see cref="Tenant.Default"/> when the tenant provider returns an empty list.
/// </summary>
public class DefaultTenantServiceTests
{
[Fact]
public async Task ActivateTenantsAsync_WhenProviderReturnsEmpty_ActivatesDefaultTenant()
{
// Arrange - provider returns no tenants
var (tenantService, serviceProvider) = await CreateTenantServiceAsync(Array.Empty<Tenant>());
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<Tenant>());
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<Tenant> { 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<Tenant> tenants, Func<List<Tenant>>? tenantsFactory = null)
{
var tenantList = tenants.ToList();
var getTenants = tenantsFactory ?? (() => tenantList);
var tenantsProvider = Substitute.For<ITenantsProvider>();
tenantsProvider.ListAsync(Arg.Any<CancellationToken>()).Returns(_ => getTenants());
var services = new ServiceCollection();
services.AddSingleton(_ => tenantsProvider);
services.AddSingleton<ITenantScopeFactory, DefaultTenantScopeFactory>();
services.AddSingleton<ITenantAccessor, DefaultTenantAccessor>();
services.AddSingleton<ITenantActivatedEvent>(Substitute.For<ITenantActivatedEvent>());
services.AddSingleton<ITenantDeactivatedEvent>(Substitute.For<ITenantDeactivatedEvent>());
services.AddSingleton<ITenantDeletedEvent>(Substitute.For<ITenantDeletedEvent>());
services.AddSingleton<TenantEventsManager>();
services.AddSingleton<RecurringTaskScheduleManager>();
services.AddSingleton<ITenantService, DefaultTenantService>();
services.AddLogging();
var serviceProvider = services.BuildServiceProvider();
return Task.FromResult((serviceProvider.GetRequiredService<ITenantService>(), serviceProvider));
}
}

View file

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

View file

@ -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<IWorkflowDefinitionStore>();
private readonly IApiSerializer _serializer = Substitute.For<IApiSerializer>();
private readonly IWorkflowDefinitionService _workflowDefinitionService = Substitute.For<IWorkflowDefinitionService>();
private readonly IWorkflowReferenceGraphBuilder _workflowReferenceGraphBuilder = Substitute.For<IWorkflowReferenceGraphBuilder>();
[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<WorkflowDefinitionFilter>(), Arg.Any<CancellationToken>())
.Returns(_ => Task.FromResult<IEnumerable<WorkflowDefinition>>([definition]));
_workflowDefinitionService.MaterializeWorkflowAsync(definition, Arg.Any<CancellationToken>())
.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<IActivitySerializer>();
var wellKnownTypeRegistry = Substitute.For<IWellKnownTypeRegistry>();
var scopeFactory = Substitute.For<IServiceScopeFactory>();
var variableDefinitionMapper = new VariableDefinitionMapper(wellKnownTypeRegistry, scopeFactory, NullLogger<VariableDefinitionMapper>.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]);
}
}

View file

@ -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<IWorkflowInstanceStore>();
var filter = new WorkflowInstanceFilter();
var requestedPages = new List<PageArgs>();
var page1 = new List<WorkflowInstanceSummary>
{
CreateSummary("workflow-1"),
CreateSummary("workflow-2")
};
var page2 = new List<WorkflowInstanceSummary>
{
CreateSummary("workflow-3"),
CreateSummary("workflow-4")
};
store
.SummarizeManyAsync(Arg.Any<WorkflowInstanceFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
.Returns(callInfo =>
{
var pageArgs = callInfo.ArgAt<PageArgs>(1);
requestedPages.Add(pageArgs with { });
var items = pageArgs.Offset switch
{
0 => page1,
2 => page2,
_ => new List<WorkflowInstanceSummary>()
};
return new ValueTask<Page<WorkflowInstanceSummary>>(Page.Of(items, page1.Count + page2.Count));
});
// Act
var results = new List<WorkflowInstanceSummary>();
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"
};
}

View file

@ -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<IWorkflowInstanceStore>();
var tenantAccessor = new DefaultTenantAccessor();
var tenantService = Substitute.For<ITenantService>();
var workflowRestarter = Substitute.For<IWorkflowRestarter>();
var clock = Substitute.For<ISystemClock>();
var logger = Substitute.For<ILogger<RestartInterruptedWorkflowsTask>>();
var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z");
var workflowInstances = new List<WorkflowInstanceSummary>
{
CreateWorkflowInstance("workflow-1", "tenant-a", now),
CreateWorkflowInstance("workflow-2", "tenant-b", now)
};
var observedTenantIds = new List<string?>();
clock.UtcNow.Returns(now);
tenantService.FindAsync("tenant-a", Arg.Any<CancellationToken>()).Returns(new Tenant { Id = "tenant-a", Name = "Tenant A" });
tenantService.FindAsync("tenant-b", Arg.Any<CancellationToken>()).Returns(new Tenant { Id = "tenant-b", Name = "Tenant B" });
workflowInstanceStore
.SummarizeManyAsync(Arg.Any<WorkflowInstanceFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
.Returns(
callInfo =>
{
var pageArgs = callInfo.ArgAt<PageArgs>(1);
var items = pageArgs.Offset == 0 ? workflowInstances : new List<WorkflowInstanceSummary>();
return new ValueTask<Page<WorkflowInstanceSummary>>(Page.Of(items, workflowInstances.Count));
});
workflowRestarter
.When(x => x.RestartWorkflowAsync(Arg.Any<string>(), Arg.Any<CancellationToken>()))
.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<CancellationToken>());
await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-2", Arg.Any<CancellationToken>());
}
[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<IWorkflowInstanceStore>();
var tenantAccessor = new DefaultTenantAccessor();
var tenantService = Substitute.For<ITenantService>();
var workflowRestarter = Substitute.For<IWorkflowRestarter>();
var clock = Substitute.For<ISystemClock>();
var logger = Substitute.For<ILogger<RestartInterruptedWorkflowsTask>>();
var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z");
var workflowInstances = new List<WorkflowInstanceSummary>
{
CreateWorkflowInstance("workflow-1", tenantId, now)
};
var observedTenantIds = new List<string?>();
clock.UtcNow.Returns(now);
workflowInstanceStore
.SummarizeManyAsync(Arg.Any<WorkflowInstanceFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
.Returns(
callInfo =>
{
var pageArgs = callInfo.ArgAt<PageArgs>(1);
var items = pageArgs.Offset == 0 ? workflowInstances : new List<WorkflowInstanceSummary>();
return new ValueTask<Page<WorkflowInstanceSummary>>(Page.Of(items, workflowInstances.Count));
});
workflowRestarter
.When(x => x.RestartWorkflowAsync(Arg.Any<string>(), Arg.Any<CancellationToken>()))
.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<string>(), Arg.Any<CancellationToken>());
await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-1", Arg.Any<CancellationToken>());
}
[Fact(DisplayName = "ExecuteAsync continues restarting later workflows after a failure")]
public async Task ExecuteAsync_ContinuesAfterFailure()
{
var workflowInstanceStore = Substitute.For<IWorkflowInstanceStore>();
var tenantAccessor = new DefaultTenantAccessor();
var tenantService = Substitute.For<ITenantService>();
var workflowRestarter = Substitute.For<IWorkflowRestarter>();
var clock = Substitute.For<ISystemClock>();
var logger = Substitute.For<ILogger<RestartInterruptedWorkflowsTask>>();
var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z");
var workflowInstances = new List<WorkflowInstanceSummary>
{
CreateWorkflowInstance("workflow-1", "tenant-a", now),
CreateWorkflowInstance("workflow-2", "tenant-b", now)
};
clock.UtcNow.Returns(now);
tenantService.FindAsync("tenant-a", Arg.Any<CancellationToken>()).Returns(_ => Task.FromException<Tenant?>(new InvalidOperationException("Transient tenant lookup failure")));
tenantService.FindAsync("tenant-b", Arg.Any<CancellationToken>()).Returns(new Tenant { Id = "tenant-b", Name = "Tenant B" });
workflowInstanceStore
.SummarizeManyAsync(Arg.Any<WorkflowInstanceFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
.Returns(
callInfo =>
{
var pageArgs = callInfo.ArgAt<PageArgs>(1);
var items = pageArgs.Offset == 0 ? workflowInstances : new List<WorkflowInstanceSummary>();
return new ValueTask<Page<WorkflowInstanceSummary>>(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<CancellationToken>());
await workflowRestarter.Received(1).RestartWorkflowAsync("workflow-2", Arg.Any<CancellationToken>());
}
[Fact(DisplayName = "ExecuteAsync restarts tenant-specific workflows when tenant services are unavailable")]
public async Task ExecuteAsync_WithoutTenantServices_RestartsWorkflow()
{
var workflowInstanceStore = Substitute.For<IWorkflowInstanceStore>();
var workflowRestarter = Substitute.For<IWorkflowRestarter>();
var clock = Substitute.For<ISystemClock>();
var logger = Substitute.For<ILogger<RestartInterruptedWorkflowsTask>>();
var now = DateTimeOffset.Parse("2026-04-14T12:00:00Z");
var workflowInstances = new List<WorkflowInstanceSummary>
{
CreateWorkflowInstance("workflow-1", "tenant-a", now)
};
clock.UtcNow.Returns(now);
workflowInstanceStore
.SummarizeManyAsync(Arg.Any<WorkflowInstanceFilter>(), Arg.Any<PageArgs>(), Arg.Any<CancellationToken>())
.Returns(
callInfo =>
{
var pageArgs = callInfo.ArgAt<PageArgs>(1);
var items = pageArgs.Offset == 0 ? workflowInstances : new List<WorkflowInstanceSummary>();
return new ValueTask<Page<WorkflowInstanceSummary>>(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<CancellationToken>());
}
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)
};
}
}