test(workflows): share Memory/EF store conformance scenarios (#8115)

* test(workflows): share Memory/EF store conformance scenarios

Add a shared management/runtime store matrix that runs the same uniqueness
and tenant-isolation assertions against Memory and EF Core (SQLite). Memory
now rejects a second version row for the same (DefinitionId, Version, TenantId)
so that DefinitionId+Version uniqueness can fail closed on both paths.

Closes #8088

Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>

* test(workflows): lock bookmark and log Find semantics

Add the remaining #8088 store-contract row: bookmark, activity-execution,
and execution-log Find/FindMany filters, paging, and same-Id upsert, run
against both Memory and EF/SQLite.

Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>

* test(workflows): correct execution-log ExcludeActivityType assertion

Exclude WriteLine only when another activity type is present, so the
shared Find row asserts a real filter match on both providers.

Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>

* fix(management): reject duplicate definition version keys in SaveMany

MemoryWorkflowDefinitionStore.SaveManyAsync now throws when a batch
contains distinct Ids that share (DefinitionId, Version, TenantId),
matching EF unique-index failure instead of silently dropping rows.

Give the label-filter list fixture distinct versions so two red
definition rows no longer collide on the default Version=1 key.

Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>

---------

Co-authored-by: Cursor Agent <cursoragent@cursor.com>
This commit is contained in:
Sipke Schoorstra 2026-09-13 18:07:09 +02:00 committed by GitHub
parent f477fc8b07
commit f3e77fe3e7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
7 changed files with 827 additions and 3 deletions

View file

@ -521,6 +521,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Labels.UnitTests", "te
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Alterations.Core.UnitTests", "test\unit\Elsa.Alterations.Core.UnitTests\Elsa.Alterations.Core.UnitTests.csproj", "{EAF427BE-7A07-42D7-A223-70FB1FC91BB2}"
EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "integration", "integration", "{823D4020-332D-2C13-F261-6F510F11A57E}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Workflows.Persistence.ConformanceTests", "test\integration\Elsa.Workflows.Persistence.ConformanceTests\Elsa.Workflows.Persistence.ConformanceTests.csproj", "{E099B352-C3E6-499D-B141-67B209125B67}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -2497,6 +2501,18 @@ Global
{EAF427BE-7A07-42D7-A223-70FB1FC91BB2}.Release|x64.Build.0 = Release|Any CPU
{EAF427BE-7A07-42D7-A223-70FB1FC91BB2}.Release|x86.ActiveCfg = Release|Any CPU
{EAF427BE-7A07-42D7-A223-70FB1FC91BB2}.Release|x86.Build.0 = Release|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Debug|Any CPU.Build.0 = Debug|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Debug|x64.ActiveCfg = Debug|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Debug|x64.Build.0 = Debug|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Debug|x86.ActiveCfg = Debug|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Debug|x86.Build.0 = Debug|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Release|Any CPU.ActiveCfg = Release|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Release|Any CPU.Build.0 = Release|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Release|x64.ActiveCfg = Release|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Release|x64.Build.0 = Release|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Release|x86.ActiveCfg = Release|Any CPU
{E099B352-C3E6-499D-B141-67B209125B67}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -2704,6 +2720,7 @@ Global
{0BE4BC01-D02D-4185-A433-E33780584261} = {90031D64-CA0F-46D0-9AF4-8DC023A5FFCD}
{ED8D7C78-154D-4124-8C0A-E63785A862E5} = {18453B51-25EB-4317-A4B3-B10518252E92}
{EAF427BE-7A07-42D7-A223-70FB1FC91BB2} = {18453B51-25EB-4317-A4B3-B10518252E92}
{E099B352-C3E6-499D-B141-67B209125B67} = {823D4020-332D-2C13-F261-6F510F11A57E}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {D4B5CEAA-7D70-4FCB-A68E-B03FBE5E0E5E}

View file

@ -98,7 +98,10 @@ public class MemoryWorkflowDefinitionStore(MemoryStore<WorkflowDefinition> store
public Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
lock (store.Sync)
{
EnsureVersionKeyAvailable(definition);
store.Save(definition, GetId);
}
return Task.CompletedTask;
}
@ -107,7 +110,15 @@ public class MemoryWorkflowDefinitionStore(MemoryStore<WorkflowDefinition> store
public Task SaveManyAsync(IEnumerable<WorkflowDefinition> definitions, CancellationToken cancellationToken = default)
{
lock (store.Sync)
store.SaveMany(definitions, GetId);
{
var definitionList = definitions.ToList();
EnsureBatchVersionKeysUnique(definitionList);
foreach (var definition in definitionList)
EnsureVersionKeyAvailable(definition);
store.SaveMany(definitionList, GetId);
}
return Task.CompletedTask;
}
@ -191,5 +202,43 @@ public class MemoryWorkflowDefinitionStore(MemoryStore<WorkflowDefinition> store
private string CurrentTenantId => tenantAccessor?.TenantId ?? Tenant.DefaultTenantId;
/// <remarks>
/// EF enforces <c>(DefinitionId, Version)</c> globally via
/// <c>IX_WorkflowDefinition_DefinitionId_Version</c>. Memory keeps tenant in the key so
/// same-tenant duplicates fail closed while cross-tenant rows remain distinct until #7539
/// adds <c>TenantId</c> to that index.
/// </remarks>
private void EnsureVersionKeyAvailable(WorkflowDefinition definition)
{
var versionKey = GetVersionKey(definition);
var existing = store.Find(x => x.Id != definition.Id && GetVersionKey(x) == versionKey);
if (existing is not null)
throw CreateVersionKeyConflict(definition);
}
private static void EnsureBatchVersionKeysUnique(IEnumerable<WorkflowDefinition> definitions)
{
var seen = new Dictionary<DefinitionVersionKey, string>();
foreach (var definition in definitions)
{
var versionKey = GetVersionKey(definition);
if (seen.TryGetValue(versionKey, out var existingId) && existingId != definition.Id)
throw CreateVersionKeyConflict(definition);
seen[versionKey] = definition.Id;
}
}
private static InvalidOperationException CreateVersionKeyConflict(WorkflowDefinition definition) =>
new($"A workflow definition already exists for definition '{definition.DefinitionId}' version {definition.Version} tenant '{definition.TenantId}'.");
private static DefinitionVersionKey GetVersionKey(WorkflowDefinition definition) =>
new(definition.DefinitionId, definition.Version, definition.TenantId);
private string GetId(WorkflowDefinition workflowDefinition) => workflowDefinition.Id;
private readonly record struct DefinitionVersionKey(string DefinitionId, int Version, string? TenantId);
}

View file

@ -0,0 +1,20 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<Include>[Elsa.Workflows.Management]*,[Elsa.Workflows.Runtime]*,[Elsa.Persistence.EFCore]*</Include>
<Threshold>0</Threshold>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Data.Sqlite"/>
<PackageReference Include="SQLitePCLRaw.bundle_e_sqlite3"/>
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\src\common\Elsa.Testing.Shared\Elsa.Testing.Shared.csproj"/>
<ProjectReference Include="..\..\..\src\modules\Elsa.Persistence.EFCore.Sqlite\Elsa.Persistence.EFCore.Sqlite.csproj"/>
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Management\Elsa.Workflows.Management.csproj"/>
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj"/>
</ItemGroup>
</Project>

View file

@ -0,0 +1,468 @@
using Elsa.Common.Models;
using Elsa.Common.Multitenancy;
using Elsa.Workflows;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Filters;
namespace Elsa.Workflows.Persistence.ConformanceTests;
/// <summary>
/// Shared Memory / EF Core store-contract assertions for workflow management and runtime ports.
/// </summary>
public abstract class WorkflowStoreConformanceTests
{
protected abstract Task<WorkflowStoreScenario> CreateScenarioAsync();
[Fact]
public async Task TriggerLogicalKeysAreUniqueAndReplaceAsyncSkipsExistingKeys()
{
await using var scenario = await CreateScenarioAsync();
var existing = Trigger("existing-id", hash: "hash-1");
await scenario.Triggers.SaveAsync(existing);
await scenario.AssertUniquenessConflictAsync(() => scenario.Triggers.SaveAsync(Trigger("other-id", hash: "hash-1")).AsTask());
var afterConflict = (await scenario.Triggers.FindManyAsync(new TriggerFilter { TenantAgnostic = true })).ToList();
Assert.Equal("existing-id", Assert.Single(afterConflict).Id);
await scenario.Triggers.SaveAsync(Trigger("existing-id", hash: "hash-2", workflowDefinitionVersionId: "v2"));
var updated = await scenario.Triggers.FindAsync(new TriggerFilter { Id = "existing-id" });
Assert.Equal("hash-2", updated!.Hash);
Assert.Equal("v2", updated.WorkflowDefinitionVersionId);
await scenario.Triggers.ReplaceAsync([], [Trigger("batch-1"), Trigger("batch-2")]);
var afterBatch = (await scenario.Triggers.FindManyAsync(new TriggerFilter { Hash = "hash-1", TenantAgnostic = true })).ToList();
Assert.Equal("batch-1", Assert.Single(afterBatch).Id);
await scenario.Triggers.ReplaceAsync([], [Trigger("skipped-id")]);
var afterSkip = (await scenario.Triggers.FindManyAsync(new TriggerFilter { Hash = "hash-1", TenantAgnostic = true })).ToList();
Assert.Equal("batch-1", Assert.Single(afterSkip).Id);
await scenario.Triggers.ReplaceAsync(afterSkip, [Trigger("replacement-id", workflowDefinitionVersionId: "v3")]);
var afterReplace = (await scenario.Triggers.FindManyAsync(new TriggerFilter { Hash = "hash-1", TenantAgnostic = true })).ToList();
var replacement = Assert.Single(afterReplace);
Assert.Equal("replacement-id", replacement.Id);
Assert.Equal("v3", replacement.WorkflowDefinitionVersionId);
var tenantA = Trigger("id-a", hash: "shared-hash");
tenantA.TenantId = "tenant-a";
var tenantB = Trigger("id-b", hash: "shared-hash");
tenantB.TenantId = "tenant-b";
await scenario.Triggers.ReplaceAsync([], [tenantA, tenantB]);
var bothTenants = (await scenario.Triggers.FindManyAsync(new TriggerFilter { Hash = "shared-hash", TenantAgnostic = true })).ToList();
Assert.Equal(2, bothTenants.Count);
Assert.Contains(bothTenants, x => x.Id == "id-a");
Assert.Contains(bothTenants, x => x.Id == "id-b");
}
[Fact]
public async Task TriggerTenantIsolationHonorsAmbientTenantAndTenantAgnostic()
{
await using var scenario = await CreateScenarioAsync();
await SeedMixedTriggersAsync(scenario);
var visible = (await scenario.Triggers.FindManyAsync(new TriggerFilter())).ToList();
Assert.Equal(2, visible.Count);
Assert.Contains(visible, x => x.Id == "id-a");
Assert.Contains(visible, x => x.Id == "id-star");
Assert.DoesNotContain(visible, x => x.Id == "id-b");
var all = (await scenario.Triggers.FindManyAsync(new TriggerFilter { TenantAgnostic = true })).ToList();
Assert.Equal(3, all.Count);
Assert.Contains(all, x => x.Id == "id-b");
Assert.Null(await scenario.Triggers.FindAsync(new TriggerFilter { Id = "id-b" }));
var deleted = await scenario.Triggers.DeleteManyAsync(new TriggerFilter());
var remaining = (await scenario.Triggers.FindManyAsync(new TriggerFilter { TenantAgnostic = true })).ToList();
Assert.Equal(2, deleted);
Assert.Equal("id-b", Assert.Single(remaining).Id);
using (scenario.UseTenant(Tenant.DefaultTenantId))
{
await scenario.Triggers.SaveAsync(Trigger("id-null", hash: "hash-null", tenantId: null));
await scenario.Triggers.SaveAsync(Trigger("id-named", hash: "hash-named", tenantId: "tenant-a"));
var defaultVisible = (await scenario.Triggers.FindManyAsync(new TriggerFilter())).ToList();
Assert.Contains(defaultVisible, x => x.Id == "id-null");
Assert.DoesNotContain(defaultVisible, x => x.Id == "id-named");
}
using (scenario.UseTenant("tenant-a"))
{
var namedVisible = (await scenario.Triggers.FindManyAsync(new TriggerFilter { Hash = "hash-null" })).ToList();
Assert.Empty(namedVisible);
}
}
[Fact]
public async Task DefinitionVersionsAreUniquePerDefinitionId()
{
await using var scenario = await CreateScenarioAsync();
var first = Definition("def-v1", "order", "tenant-a");
await scenario.Definitions.SaveAsync(first);
await scenario.AssertUniquenessConflictAsync(() => scenario.Definitions.SaveAsync(Definition("def-v1-dup", "order", "tenant-a")));
var afterConflict = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter
{
DefinitionId = "order",
VersionOptions = VersionOptions.SpecificVersion(1),
TenantAgnostic = true
})).ToList();
Assert.Equal("def-v1", Assert.Single(afterConflict).Id);
first.Name = "Order Updated";
await scenario.Definitions.SaveAsync(first);
var updated = await scenario.Definitions.FindAsync(new WorkflowDefinitionFilter { Id = "def-v1" });
Assert.Equal("Order Updated", updated!.Name);
await scenario.Definitions.SaveAsync(Definition("def-v2", "order", "tenant-a", version: 2));
var versions = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter { DefinitionId = "order", TenantAgnostic = true })).ToList();
Assert.Equal(2, versions.Count);
Assert.Contains(versions, x => x.Id == "def-v1" && x.Version == 1);
Assert.Contains(versions, x => x.Id == "def-v2" && x.Version == 2);
}
[Fact]
public async Task DefinitionSaveManyRejectsDuplicateVersionKeysInTheBatch()
{
await using var scenario = await CreateScenarioAsync();
await scenario.Definitions.SaveAsync(Definition("def-kept", "kept", "tenant-a"));
await scenario.AssertUniquenessConflictAsync(() => scenario.Definitions.SaveManyAsync(
[
Definition("def-batch-1", "invoice", "tenant-a"),
Definition("def-batch-2", "invoice", "tenant-a")
]));
var invoices = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter
{
DefinitionId = "invoice",
TenantAgnostic = true
})).ToList();
Assert.Empty(invoices);
Assert.NotNull(await scenario.Definitions.FindAsync(new WorkflowDefinitionFilter { Id = "def-kept", TenantAgnostic = true }));
}
[Fact]
public async Task DefinitionTenantIsolationHonorsAmbientTenantAndTenantAgnostic()
{
await using var scenario = await CreateScenarioAsync();
await SeedMixedDefinitionsAsync(scenario);
var visible = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter())).ToList();
Assert.Equal(2, visible.Count);
Assert.Contains(visible, x => x.Id == "def-a");
Assert.Contains(visible, x => x.Id == "def-star");
Assert.DoesNotContain(visible, x => x.Id == "def-b");
var all = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter { TenantAgnostic = true })).ToList();
Assert.Equal(3, all.Count);
Assert.Contains(all, x => x.Id == "def-b");
Assert.Null(await scenario.Definitions.FindAsync(new WorkflowDefinitionFilter { Id = "def-b" }));
Assert.False(await scenario.Definitions.AnyAsync(new WorkflowDefinitionFilter { Id = "def-b" }));
Assert.Equal(2, await scenario.Definitions.CountDistinctAsync());
var deleted = await scenario.Definitions.DeleteAsync(new WorkflowDefinitionFilter());
var remaining = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter { TenantAgnostic = true })).ToList();
Assert.Equal(2, deleted);
Assert.Equal("def-b", Assert.Single(remaining).Id);
using (scenario.UseTenant(Tenant.DefaultTenantId))
{
await scenario.Definitions.SaveAsync(Definition("def-null", "Null", tenantId: null));
await scenario.Definitions.SaveAsync(Definition("def-named", "Named", "tenant-a"));
var defaultVisible = (await scenario.Definitions.FindManyAsync(new WorkflowDefinitionFilter())).ToList();
Assert.Contains(defaultVisible, x => x.Id == "def-null");
Assert.DoesNotContain(defaultVisible, x => x.Id == "def-named");
}
}
[Fact]
public async Task BookmarkTenantIsolationHonorsAmbientTenantAndTenantAgnostic()
{
await using var scenario = await CreateScenarioAsync();
await SeedMixedBookmarksAsync(scenario);
var visible = (await scenario.Bookmarks.FindManyAsync(new BookmarkFilter())).ToList();
Assert.Equal(2, visible.Count);
Assert.Contains(visible, x => x.Id == "bm-a");
Assert.Contains(visible, x => x.Id == "bm-star");
Assert.DoesNotContain(visible, x => x.Id == "bm-b");
var all = (await scenario.Bookmarks.FindManyAsync(new BookmarkFilter { TenantAgnostic = true })).ToList();
Assert.Equal(3, all.Count);
Assert.Contains(all, x => x.Id == "bm-b");
Assert.Null(await scenario.Bookmarks.FindAsync(new BookmarkFilter { BookmarkId = "bm-b" }));
var deleted = await scenario.Bookmarks.DeleteAsync(new BookmarkFilter());
var remaining = (await scenario.Bookmarks.FindManyAsync(new BookmarkFilter { TenantAgnostic = true })).ToList();
Assert.Equal(2, deleted);
Assert.Equal("bm-b", Assert.Single(remaining).Id);
using (scenario.UseTenant(Tenant.DefaultTenantId))
{
await scenario.Bookmarks.SaveAsync(Bookmark("bm-null", tenantId: null));
await scenario.Bookmarks.SaveAsync(Bookmark("bm-named", tenantId: "tenant-a"));
var defaultVisible = (await scenario.Bookmarks.FindManyAsync(new BookmarkFilter())).ToList();
Assert.Contains(defaultVisible, x => x.Id == "bm-null");
Assert.DoesNotContain(defaultVisible, x => x.Id == "bm-named");
}
}
[Fact]
public async Task DeadLetterOriginalQueueItemIdIsUniqueAndAddOrGetIsIdempotent()
{
await using var scenario = await CreateScenarioAsync();
var first = DeadLetter("dl-1", "queue-1");
await scenario.DeadLetters.SaveAsync(first);
await scenario.AssertUniquenessConflictAsync(() => scenario.DeadLetters.SaveAsync(DeadLetter("dl-2", "queue-1")));
var afterConflict = (await scenario.DeadLetters.FindManyAsync(new BookmarkQueueDeadLetterFilter { OriginalQueueItemId = "queue-1" })).ToList();
Assert.Equal("dl-1", Assert.Single(afterConflict).Id);
var existing = await scenario.DeadLetters.AddOrGetExistingAsync(DeadLetter("dl-3", "queue-1"));
Assert.Equal("dl-1", existing.Id);
var created = await scenario.DeadLetters.AddOrGetExistingAsync(DeadLetter("dl-4", "queue-2"));
Assert.Equal("dl-4", created.Id);
var all = (await scenario.DeadLetters.FindManyAsync(new BookmarkQueueDeadLetterFilter())).ToList();
Assert.Equal(2, all.Count);
Assert.Contains(all, x => x.Id == "dl-1");
Assert.Contains(all, x => x.Id == "dl-4");
}
[Fact]
public async Task BookmarkActivityExecutionAndExecutionLogFindsHonorFiltersAndIdUpsert()
{
await using var scenario = await CreateScenarioAsync();
await scenario.Bookmarks.SaveAsync(Bookmark("bm-1", hash: "hash-http", workflowInstanceId: "instance-1"));
await scenario.Bookmarks.SaveAsync(Bookmark("bm-2", hash: "hash-timer", workflowInstanceId: "instance-1"));
await scenario.Bookmarks.SaveAsync(Bookmark("bm-3", hash: "hash-http", workflowInstanceId: "instance-2"));
Assert.Equal("bm-1", (await scenario.Bookmarks.FindAsync(new BookmarkFilter { BookmarkId = "bm-1" }))!.Id);
Assert.Equal("bm-2", (await scenario.Bookmarks.FindAsync(new BookmarkFilter { Hash = "hash-timer" }))!.Id);
Assert.Null(await scenario.Bookmarks.FindAsync(new BookmarkFilter { BookmarkId = "missing" }));
var instanceBookmarks = (await scenario.Bookmarks.FindManyAsync(new BookmarkFilter { WorkflowInstanceId = "instance-1" })).ToList();
Assert.Equal(2, instanceBookmarks.Count);
Assert.Contains(instanceBookmarks, x => x.Id == "bm-1");
Assert.Contains(instanceBookmarks, x => x.Id == "bm-2");
var hashed = (await scenario.Bookmarks.FindManyAsync(new BookmarkFilter { Hash = "hash-http" })).ToList();
Assert.Equal(2, hashed.Count);
var firstPage = await scenario.Bookmarks.FindManyAsync(new BookmarkFilter { WorkflowInstanceId = "instance-1" }, PageArgs.FromRange(0, 1));
Assert.Equal(2, firstPage.TotalCount);
Assert.Single(firstPage.Items);
await scenario.Bookmarks.SaveAsync(Bookmark("bm-1", hash: "hash-updated", workflowInstanceId: "instance-1"));
Assert.Equal("hash-updated", (await scenario.Bookmarks.FindAsync(new BookmarkFilter { BookmarkId = "bm-1" }))!.Hash);
await scenario.ActivityExecutions.SaveAsync(ActivityExecution("ae-1", "instance-1", "activity-a", ActivityStatus.Running));
await scenario.ActivityExecutions.SaveAsync(ActivityExecution("ae-2", "instance-1", "activity-b", ActivityStatus.Completed, completedAt: StartedAt.AddMinutes(1)));
await scenario.ActivityExecutions.SaveAsync(ActivityExecution("ae-3", "instance-2", "activity-a", ActivityStatus.Running));
Assert.Equal("ae-1", (await scenario.ActivityExecutions.FindAsync(new ActivityExecutionRecordFilter { Id = "ae-1" }))!.Id);
Assert.Null(await scenario.ActivityExecutions.FindAsync(new ActivityExecutionRecordFilter { Id = "missing" }));
var instanceActivities = (await scenario.ActivityExecutions.FindManyAsync(new ActivityExecutionRecordFilter { WorkflowInstanceId = "instance-1" })).ToList();
Assert.Equal(2, instanceActivities.Count);
Assert.Equal(2, await scenario.ActivityExecutions.CountAsync(new ActivityExecutionRecordFilter { WorkflowInstanceId = "instance-1" }));
Assert.Equal("ae-1", Assert.Single(await scenario.ActivityExecutions.FindManyAsync(new ActivityExecutionRecordFilter { ActivityId = "activity-a", WorkflowInstanceId = "instance-1" })).Id);
Assert.Equal("ae-2", Assert.Single(await scenario.ActivityExecutions.FindManyAsync(new ActivityExecutionRecordFilter { Status = ActivityStatus.Completed })).Id);
await scenario.ActivityExecutions.SaveAsync(ActivityExecution("ae-1", "instance-1", "activity-a", ActivityStatus.Completed, completedAt: StartedAt.AddMinutes(2)));
Assert.Equal(ActivityStatus.Completed, (await scenario.ActivityExecutions.FindAsync(new ActivityExecutionRecordFilter { Id = "ae-1" }))!.Status);
Assert.Equal(2, await scenario.ActivityExecutions.DeleteManyAsync(new ActivityExecutionRecordFilter { WorkflowInstanceId = "instance-1" }));
Assert.Equal("ae-3", Assert.Single(await scenario.ActivityExecutions.FindManyAsync(new ActivityExecutionRecordFilter())).Id);
await scenario.ExecutionLogs.SaveAsync(ExecutionLog("el-1", "instance-1", "activity-a", "Started"));
await scenario.ExecutionLogs.SaveAsync(ExecutionLog("el-2", "instance-1", "activity-b", "Completed", activityType: "Elsa.Delay"));
await scenario.ExecutionLogs.SaveAsync(ExecutionLog("el-3", "instance-2", "activity-a", "Started"));
Assert.Equal("el-1", (await scenario.ExecutionLogs.FindAsync(new WorkflowExecutionLogRecordFilter { Id = "el-1" }))!.Id);
Assert.Null(await scenario.ExecutionLogs.FindAsync(new WorkflowExecutionLogRecordFilter { Id = "missing" }));
var instanceLogs = await scenario.ExecutionLogs.FindManyAsync(new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = "instance-1" }, PageArgs.FromRange(0, 10));
Assert.Equal(2, instanceLogs.TotalCount);
Assert.Equal(2, instanceLogs.Items.Count);
var firstLogPage = await scenario.ExecutionLogs.FindManyAsync(new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = "instance-1" }, PageArgs.FromRange(0, 1));
Assert.Equal(2, firstLogPage.TotalCount);
Assert.Single(firstLogPage.Items);
Assert.Equal("el-2", (await scenario.ExecutionLogs.FindAsync(new WorkflowExecutionLogRecordFilter { EventName = "Completed" }))!.Id);
Assert.Equal("el-2", Assert.Single((await scenario.ExecutionLogs.FindManyAsync(new WorkflowExecutionLogRecordFilter { ActivityId = "activity-b" }, PageArgs.All)).Items).Id);
var excluded = await scenario.ExecutionLogs.FindManyAsync(new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = "instance-1", ExcludeActivityType = "Elsa.WriteLine" }, PageArgs.All);
Assert.Equal("el-2", Assert.Single(excluded.Items).Id);
await scenario.ExecutionLogs.SaveAsync(ExecutionLog("el-1", "instance-1", "activity-a", "Resumed"));
Assert.Equal("Resumed", (await scenario.ExecutionLogs.FindAsync(new WorkflowExecutionLogRecordFilter { Id = "el-1" }))!.EventName);
Assert.Equal(2, await scenario.ExecutionLogs.DeleteManyAsync(new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = "instance-1" }));
Assert.Equal("el-3", Assert.Single((await scenario.ExecutionLogs.FindManyAsync(new WorkflowExecutionLogRecordFilter(), PageArgs.All)).Items).Id);
}
private static async Task SeedMixedTriggersAsync(WorkflowStoreScenario scenario)
{
await scenario.Triggers.SaveAsync(Trigger("id-a", hash: "hash-a", tenantId: "tenant-a"));
await scenario.Triggers.SaveAsync(Trigger("id-b", hash: "hash-b", tenantId: "tenant-b"));
await scenario.Triggers.SaveAsync(Trigger("id-star", hash: "hash-star", tenantId: Tenant.AgnosticTenantId));
}
private static async Task SeedMixedDefinitionsAsync(WorkflowStoreScenario scenario)
{
await scenario.Definitions.SaveAsync(Definition("def-a", "A", "tenant-a"));
await scenario.Definitions.SaveAsync(Definition("def-b", "B", "tenant-b"));
await scenario.Definitions.SaveAsync(Definition("def-star", "Star", Tenant.AgnosticTenantId));
}
private static async Task SeedMixedBookmarksAsync(WorkflowStoreScenario scenario)
{
await scenario.Bookmarks.SaveAsync(Bookmark("bm-a", tenantId: "tenant-a"));
await scenario.Bookmarks.SaveAsync(Bookmark("bm-b", tenantId: "tenant-b"));
await scenario.Bookmarks.SaveAsync(Bookmark("bm-star", tenantId: Tenant.AgnosticTenantId));
}
private static StoredTrigger Trigger(
string id,
string workflowDefinitionId = "workflow-1",
string workflowDefinitionVersionId = "v1",
string activityId = "activity-1",
string? hash = "hash-1",
string? tenantId = "tenant-a") =>
new()
{
Id = id,
TenantId = tenantId,
WorkflowDefinitionId = workflowDefinitionId,
WorkflowDefinitionVersionId = workflowDefinitionVersionId,
ActivityId = activityId,
Hash = hash,
Name = "Elsa.HttpEndpoint"
};
private static WorkflowDefinition Definition(string id, string name, string? tenantId, string? definitionId = null, int version = 1) =>
new()
{
Id = id,
DefinitionId = definitionId ?? name.ToLowerInvariant(),
Name = name,
TenantId = tenantId,
Version = version,
IsLatest = version == 1,
MaterializerName = "Json",
CreatedAt = new DateTimeOffset(2026, 9, 13, 12, 0, 0, TimeSpan.Zero)
};
private static readonly DateTimeOffset StartedAt = new(2026, 9, 13, 12, 0, 0, TimeSpan.Zero);
private static StoredBookmark Bookmark(
string id,
string? tenantId = "tenant-a",
string? hash = null,
string workflowInstanceId = "instance-1") =>
new()
{
Id = id,
TenantId = tenantId,
Hash = hash ?? id,
WorkflowInstanceId = workflowInstanceId,
Name = "Elsa.HttpEndpoint",
CreatedAt = StartedAt
};
private static ActivityExecutionRecord ActivityExecution(
string id,
string workflowInstanceId,
string activityId,
ActivityStatus status,
DateTimeOffset? completedAt = null) =>
new()
{
Id = id,
TenantId = "tenant-a",
WorkflowInstanceId = workflowInstanceId,
ActivityId = activityId,
ActivityNodeId = $"node-{activityId}",
ActivityType = "Elsa.WriteLine",
ActivityTypeVersion = 1,
ActivityName = activityId,
Status = status,
StartedAt = StartedAt,
CompletedAt = completedAt
};
private static WorkflowExecutionLogRecord ExecutionLog(
string id,
string workflowInstanceId,
string activityId,
string eventName,
string activityType = "Elsa.WriteLine") =>
new()
{
Id = id,
TenantId = "tenant-a",
WorkflowDefinitionId = "workflow-1",
WorkflowDefinitionVersionId = "v1",
WorkflowInstanceId = workflowInstanceId,
WorkflowVersion = 1,
ActivityInstanceId = $"ai-{id}",
ActivityId = activityId,
ActivityType = activityType,
ActivityTypeVersion = 1,
ActivityNodeId = $"node-{activityId}",
Timestamp = StartedAt,
Sequence = 0,
EventName = eventName
};
private static BookmarkQueueDeadLetterItem DeadLetter(string id, string originalQueueItemId) =>
new()
{
Id = id,
TenantId = "tenant-a",
OriginalQueueItemId = originalQueueItemId,
WorkflowInstanceId = "instance-1",
Reason = "delivery-failed",
OriginalCreatedAt = new DateTimeOffset(2026, 9, 13, 12, 0, 0, TimeSpan.Zero),
DeadLetteredAt = new DateTimeOffset(2026, 9, 13, 12, 5, 0, TimeSpan.Zero)
};
}
[CollectionDefinition(Name)]
public sealed class WorkflowStoreInMemoryConformanceCollection
{
public const string Name = "WorkflowStores:InMemory";
}
[CollectionDefinition(Name)]
public sealed class WorkflowStoreSqliteConformanceCollection
{
public const string Name = "WorkflowStores:EFCore.Sqlite";
}
[Collection(WorkflowStoreInMemoryConformanceCollection.Name)]
public sealed class InMemoryWorkflowStoreConformanceTests : WorkflowStoreConformanceTests
{
protected override Task<WorkflowStoreScenario> CreateScenarioAsync() => WorkflowStoreScenario.CreateInMemoryAsync();
}
[Collection(WorkflowStoreSqliteConformanceCollection.Name)]
public sealed class SqliteWorkflowStoreConformanceTests : WorkflowStoreConformanceTests
{
protected override Task<WorkflowStoreScenario> CreateScenarioAsync() => WorkflowStoreScenario.CreateSqliteAsync();
}

View file

@ -0,0 +1,211 @@
using System.Text.Json;
using Elsa.Common;
using Elsa.Common.Codecs;
using Elsa.Common.Multitenancy;
using Elsa.Common.Services;
using Elsa.Persistence.EFCore;
using Elsa.Persistence.EFCore.EntityHandlers;
using Elsa.Persistence.EFCore.Extensions;
using Elsa.Persistence.EFCore.Modules.Management;
using Elsa.Persistence.EFCore.Modules.Runtime;
using Elsa.Persistence.EFCore.Sqlite;
using Elsa.Tenants.Options;
using Elsa.Testing.Shared.Multitenancy;
using Elsa.Workflows;
using Elsa.Workflows.Management;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Stores;
using Elsa.Workflows.Runtime;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Stores;
using Microsoft.Data.Sqlite;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Persistence.ConformanceTests;
/// <summary>
/// Holds one Memory or EF/SQLite set of management + runtime stores and the ambient tenant they read.
/// </summary>
public sealed class WorkflowStoreScenario(
TestTenantAccessor tenantAccessor,
IWorkflowDefinitionStore definitions,
ITriggerStore triggers,
IBookmarkStore bookmarks,
IBookmarkQueueDeadLetterStore deadLetters,
IActivityExecutionStore activityExecutions,
IWorkflowExecutionLogStore executionLogs,
Func<Func<Task>, Task> assertUniquenessConflictAsync,
Func<ValueTask> disposeAsync) : IAsyncDisposable
{
public TestTenantAccessor TenantAccessor { get; } = tenantAccessor;
public IWorkflowDefinitionStore Definitions { get; } = definitions;
public ITriggerStore Triggers { get; } = triggers;
public IBookmarkStore Bookmarks { get; } = bookmarks;
public IBookmarkQueueDeadLetterStore DeadLetters { get; } = deadLetters;
public IActivityExecutionStore ActivityExecutions { get; } = activityExecutions;
public IWorkflowExecutionLogStore ExecutionLogs { get; } = executionLogs;
public IDisposable UseTenant(string tenantId) =>
TenantAccessor.PushContext(tenantId == Tenant.DefaultTenantId
? Tenant.Default
: new Tenant { Id = tenantId, Name = tenantId });
public Task AssertUniquenessConflictAsync(Func<Task> operation) => assertUniquenessConflictAsync(operation);
public ValueTask DisposeAsync() => disposeAsync();
public static Task<WorkflowStoreScenario> CreateInMemoryAsync()
{
var tenantAccessor = new TestTenantAccessor("tenant-a");
return Task.FromResult(new WorkflowStoreScenario(
tenantAccessor,
new MemoryWorkflowDefinitionStore(new MemoryStore<WorkflowDefinition>(), tenantAccessor),
new MemoryTriggerStore(new MemoryStore<StoredTrigger>(), tenantAccessor),
new MemoryBookmarkStore(new MemoryStore<StoredBookmark>(), tenantAccessor),
new MemoryBookmarkQueueDeadLetterStore(new MemoryStore<BookmarkQueueDeadLetterItem>()),
new MemoryActivityExecutionStore(new MemoryStore<ActivityExecutionRecord>()),
new MemoryWorkflowExecutionLogStore(new MemoryStore<WorkflowExecutionLogRecord>()),
operation => Assert.ThrowsAsync<InvalidOperationException>(operation),
() => ValueTask.CompletedTask));
}
public static async Task<WorkflowStoreScenario> CreateSqliteAsync()
{
var managementPath = Path.Combine(Path.GetTempPath(), $"elsa-workflow-management-conformance-{Guid.NewGuid():N}.db");
var runtimePath = Path.Combine(Path.GetTempPath(), $"elsa-workflow-runtime-conformance-{Guid.NewGuid():N}.db");
var tenantAccessor = new TestTenantAccessor("tenant-a");
ServiceProvider? services = null;
IServiceScope? scope = null;
try
{
var migrationsAssembly = typeof(ManagementDbContextFactory).Assembly;
services = new ServiceCollection()
.AddLogging()
.AddSingleton<ITenantAccessor>(tenantAccessor)
.AddSingleton<IPayloadSerializer, ConformancePayloadSerializer>()
.AddSingleton<ISafeSerializer, ConformanceSafeSerializer>()
.AddSingleton<ICompressionCodecResolver>(_ => new CompressionCodecResolver([new None()]))
.Configure<TenantsOptions>(options => options.IsEnabled = true)
.AddScoped<IEntitySavingHandler, ApplyTenantId>()
.AddScoped<IEntityModelCreatingHandler, SetTenantIdFilter>()
.AddSqliteEntityModelCreatingHandlers()
.AddDbContextFactory<ManagementElsaDbContext>((_, builder) =>
builder.UseElsaSqlite(migrationsAssembly, $"Data Source={managementPath};Default Timeout=30"))
.AddDbContextFactory<RuntimeElsaDbContext>((_, builder) =>
builder.UseElsaSqlite(migrationsAssembly, $"Data Source={runtimePath};Default Timeout=30"))
.Decorate<IDbContextFactory<ManagementElsaDbContext>, TenantAwareDbContextFactory<ManagementElsaDbContext>>()
.Decorate<IDbContextFactory<RuntimeElsaDbContext>, TenantAwareDbContextFactory<RuntimeElsaDbContext>>()
.AddScoped<EntityStore<ManagementElsaDbContext, WorkflowDefinition>>()
.AddScoped<EntityStore<RuntimeElsaDbContext, StoredTrigger>>()
.AddScoped<Store<RuntimeElsaDbContext, StoredBookmark>>()
.AddScoped<Store<RuntimeElsaDbContext, BookmarkQueueDeadLetterItem>>()
.AddScoped<EntityStore<RuntimeElsaDbContext, ActivityExecutionRecord>>()
.AddScoped<EntityStore<RuntimeElsaDbContext, WorkflowExecutionLogRecord>>()
.AddScoped<EFCoreWorkflowDefinitionStore>()
.AddScoped<EFCoreTriggerStore>()
.AddScoped<EFCoreBookmarkStore>()
.AddScoped<EFBookmarkQueueDeadLetterStore>()
.AddScoped<EFCoreActivityExecutionStore>()
.AddScoped<EFCoreWorkflowExecutionLogStore>()
.BuildServiceProvider();
await using (var management = await services.GetRequiredService<IDbContextFactory<ManagementElsaDbContext>>().CreateDbContextAsync())
await management.Database.EnsureCreatedAsync();
await using (var runtime = await services.GetRequiredService<IDbContextFactory<RuntimeElsaDbContext>>().CreateDbContextAsync())
await runtime.Database.EnsureCreatedAsync();
scope = services.CreateScope();
var scoped = scope.ServiceProvider;
return new(
tenantAccessor,
scoped.GetRequiredService<EFCoreWorkflowDefinitionStore>(),
scoped.GetRequiredService<EFCoreTriggerStore>(),
scoped.GetRequiredService<EFCoreBookmarkStore>(),
scoped.GetRequiredService<EFBookmarkQueueDeadLetterStore>(),
scoped.GetRequiredService<EFCoreActivityExecutionStore>(),
scoped.GetRequiredService<EFCoreWorkflowExecutionLogStore>(),
AssertSqliteUniquenessConflictAsync,
async () =>
{
scope.Dispose();
await services.DisposeAsync();
SqliteConnection.ClearAllPools();
File.Delete(managementPath);
File.Delete(runtimePath);
});
}
catch
{
scope?.Dispose();
if (services is not null)
await services.DisposeAsync();
SqliteConnection.ClearAllPools();
File.Delete(managementPath);
File.Delete(runtimePath);
throw;
}
}
private static async Task AssertSqliteUniquenessConflictAsync(Func<Task> operation)
{
var exception = await Record.ExceptionAsync(operation);
var sqliteException = exception switch
{
DbUpdateException { InnerException: SqliteException inner } => inner,
SqliteException direct => direct,
_ => throw new Xunit.Sdk.XunitException($"Expected a SQLite uniqueness violation, received {exception?.GetType().FullName ?? "no exception"}.")
};
Assert.Equal(19, sqliteException.SqliteErrorCode);
Assert.Contains("UNIQUE constraint failed", sqliteException.Message, StringComparison.Ordinal);
}
private sealed class ConformancePayloadSerializer : IPayloadSerializer
{
private static readonly JsonSerializerOptions Options = new()
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
PropertyNameCaseInsensitive = true
};
public string Serialize(object payload) => JsonSerializer.Serialize(payload, Options);
public JsonElement SerializeToElement(object payload) => JsonSerializer.SerializeToElement(payload, Options);
public object Deserialize(string serializedData) => JsonSerializer.Deserialize<object>(serializedData, Options)!;
public object Deserialize(string serializedData, Type type) => JsonSerializer.Deserialize(serializedData, type, Options)!;
public object Deserialize(JsonElement serializedData) => serializedData.Deserialize<object>(Options)!;
public T Deserialize<T>(string serializedData) => JsonSerializer.Deserialize<T>(serializedData, Options)!;
public T Deserialize<T>(JsonElement serializedData) => serializedData.Deserialize<T>(Options)!;
public JsonSerializerOptions GetOptions() => Options;
}
private sealed class ConformanceSafeSerializer : ISafeSerializer
{
private static readonly JsonSerializerOptions Options = new()
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
PropertyNameCaseInsensitive = true
};
public ValueTask<string> SerializeAsync(object? value, CancellationToken cancellationToken = default) =>
ValueTask.FromResult(Serialize(value));
public ValueTask<JsonElement> SerializeToElementAsync(object? value, CancellationToken cancellationToken = default) =>
new(SerializeToElement(value));
public ValueTask<T> DeserializeAsync<T>(string json, CancellationToken cancellationToken = default) =>
new(Deserialize<T>(json));
public ValueTask<T> DeserializeAsync<T>(JsonElement element, CancellationToken cancellationToken = default) =>
new(Deserialize<T>(element));
public string Serialize(object? value) => JsonSerializer.Serialize(value, Options);
public JsonElement SerializeToElement(object? value) => JsonSerializer.SerializeToElement(value, Options);
public T Deserialize<T>(string json) => JsonSerializer.Deserialize<T>(json, Options)!;
public T Deserialize<T>(JsonElement element) => element.Deserialize<T>(Options)!;
public JsonSerializerOptions GetOptions() => Options;
}
}

View file

@ -164,8 +164,8 @@ public class WorkflowDefinitionLabelFilterTests
var workflowDefinitionStore = new MemoryWorkflowDefinitionStore(memoryStore);
await workflowDefinitionStore.SaveManyAsync(
[
new WorkflowDefinition { Id = "red-version", DefinitionId = "red", Name = "Red", MaterializerName = "Json" },
new WorkflowDefinition { Id = "red-unlabeled-version", DefinitionId = "red", Name = "Red", MaterializerName = "Json" },
new WorkflowDefinition { Id = "red-version", DefinitionId = "red", Name = "Red", MaterializerName = "Json", Version = 1 },
new WorkflowDefinition { Id = "red-unlabeled-version", DefinitionId = "red", Name = "Red", MaterializerName = "Json", Version = 2 },
new WorkflowDefinition { Id = "red-second-version", DefinitionId = "red-second", Name = "Red second", MaterializerName = "Json" },
new WorkflowDefinition { Id = "blue-version", DefinitionId = "blue", Name = "Blue", MaterializerName = "Json" },
new WorkflowDefinition { Id = "plain-version", DefinitionId = "plain", Name = "Plain", MaterializerName = "Json" }

View file

@ -0,0 +1,59 @@
using Elsa.Common.Services;
using Elsa.Testing.Shared.Multitenancy;
using Elsa.Workflows.Management.Entities;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Stores;
namespace Elsa.Workflows.Management.UnitTests.Stores;
/// <summary>
/// Memory must reject colliding <c>(DefinitionId, Version, TenantId)</c> keys the same way EF
/// rejects <c>IX_WorkflowDefinition_DefinitionId_Version</c> on Save/SaveMany.
/// </summary>
public class MemoryWorkflowDefinitionStoreUniquenessTests
{
[Fact(DisplayName = "SaveManyAsync rejects two different Ids that share a version key")]
public async Task SaveManyAsync_WhenBatchRepeatsVersionKey_ThrowsAndLeavesStoreUnchanged()
{
var store = CreateStore();
var exception = await Assert.ThrowsAsync<InvalidOperationException>(() =>
store.SaveManyAsync(
[
Definition("def-1", "order"),
Definition("def-2", "order")
]));
Assert.Contains("already exists", exception.Message);
Assert.Empty((await store.FindManyAsync(new WorkflowDefinitionFilter { TenantAgnostic = true })).ToList());
}
[Fact(DisplayName = "SaveManyAsync rejects a batch whose version key is already present under another Id")]
public async Task SaveManyAsync_WhenVersionKeyExistsUnderAnotherId_ThrowsAndLeavesExisting()
{
var store = CreateStore();
await store.SaveAsync(Definition("def-1", "order"));
var exception = await Assert.ThrowsAsync<InvalidOperationException>(() =>
store.SaveManyAsync([Definition("def-2", "order")]));
Assert.Contains("already exists", exception.Message);
var stored = (await store.FindManyAsync(new WorkflowDefinitionFilter { TenantAgnostic = true })).ToList();
Assert.Equal("def-1", Assert.Single(stored).Id);
}
private static MemoryWorkflowDefinitionStore CreateStore() =>
new(new MemoryStore<WorkflowDefinition>(), new TestTenantAccessor("tenant-a"));
private static WorkflowDefinition Definition(string id, string definitionId) =>
new()
{
Id = id,
DefinitionId = definitionId,
Name = definitionId,
TenantId = "tenant-a",
Version = 1,
IsLatest = true,
MaterializerName = "Json"
};
}