fix(alterations): atomic EF tenant ownership on Save/SaveMany (#8128)

Closes #8125.

Alterations-local CAS via ExecuteUpdate+INSERT with ambient ownership gate,
same-tenant race retry, multi-TFM setters, persistence compatibility ctors.
This commit is contained in:
Sipke Schoorstra 2026-09-14 01:38:40 +02:00 committed by GitHub
parent 815d7b7f70
commit 54209e24c2
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
15 changed files with 1256 additions and 54 deletions

View file

@ -0,0 +1,13 @@
namespace Elsa.Alterations.Core.Stores;
/// <summary>
/// Shared refuse-to-overwrite exceptions for Memory and EF alteration stores.
/// </summary>
public static class AlterationStoreConflict
{
public static InvalidOperationException HiddenPlanId(string id) =>
new($"An alteration plan with ID '{id}' already exists and is not visible to the current tenant.");
public static InvalidOperationException HiddenJobId(string id) =>
new($"An alteration job with ID '{id}' already exists and is not visible to the current tenant.");
}

View file

@ -98,27 +98,19 @@ public class MemoryAlterationJobStore : IAlterationJobStore
private string CurrentTenantId => _tenantAccessor?.TenantId ?? Tenant.DefaultTenantId;
private bool IsVisible(Entity entity) => TenantVisibility.IsVisible(entity.TenantId, CurrentTenantId);
private void EnsureIdAvailable(AlterationJob job)
{
var existing = _store.Find(x => x.Id == job.Id);
if (existing is not null && !CanReplace(existing))
{
throw new InvalidOperationException(
$"An alteration job with ID '{job.Id}' already exists and is not visible to the current tenant.");
}
throw AlterationStoreConflict.HiddenJobId(job.Id);
}
/// <summary>
/// <c>*</c> is visible to every tenant, but only an agnostic writer may replace it.
/// Named tenants may upsert their own visible rows.
/// </summary>
private bool CanReplace(Entity existing) =>
existing.TenantId == Tenant.AgnosticTenantId
? CurrentTenantId == Tenant.AgnosticTenantId
: IsVisible(existing);
private bool CanReplace(Entity existing) => TenantVisibility.CanReplace(existing.TenantId, CurrentTenantId);
private void ApplyCurrentTenant(Entity entity)
{

View file

@ -65,27 +65,19 @@ public class MemoryAlterationPlanStore : IAlterationPlanStore
private string CurrentTenantId => _tenantAccessor?.TenantId ?? Tenant.DefaultTenantId;
private bool IsVisible(Entity entity) => TenantVisibility.IsVisible(entity.TenantId, CurrentTenantId);
private void EnsureIdAvailable(AlterationPlan plan)
{
var existing = _store.Find(x => x.Id == plan.Id);
if (existing is not null && !CanReplace(existing))
{
throw new InvalidOperationException(
$"An alteration plan with ID '{plan.Id}' already exists and is not visible to the current tenant.");
}
throw AlterationStoreConflict.HiddenPlanId(plan.Id);
}
/// <summary>
/// <c>*</c> is visible to every tenant, but only an agnostic writer may replace it.
/// Named tenants may upsert their own visible rows.
/// </summary>
private bool CanReplace(Entity existing) =>
existing.TenantId == Tenant.AgnosticTenantId
? CurrentTenantId == Tenant.AgnosticTenantId
: IsVisible(existing);
private bool CanReplace(Entity existing) => TenantVisibility.CanReplace(existing.TenantId, CurrentTenantId);
private void ApplyCurrentTenant(Entity entity)
{

View file

@ -17,6 +17,28 @@ public static class TenantVisibility
|| entityTenantId == Tenant.AgnosticTenantId
|| entityTenantId is null && ambientTenantId == Tenant.DefaultTenantId;
/// <summary>
/// Write counterpart of <see cref="IsVisible"/>. <c>*</c> is visible to every tenant, but only an
/// agnostic writer may replace it. Named tenants may replace their own rows (and the default tenant
/// may replace a null <see cref="Entity.TenantId"/>).
/// </summary>
public static bool CanReplace(string? existingTenantId, string writerTenantId) =>
existingTenantId == writerTenantId
|| existingTenantId is null && writerTenantId == Tenant.DefaultTenantId;
/// <summary>
/// CAS match for an upsert. Named rows are owned by the <paramref name="ambientTenantId"/>
/// (Memory <see cref="CanReplace"/>), not by a forged source <c>TenantId</c>. A <c>*</c> row
/// is replaceable only when both the incoming entity and the ambient writer are <c>*</c>.
/// </summary>
public static bool CanReplaceOwnedRow(string? existingTenantId, string? sourceTenantId, string ambientTenantId)
{
if (existingTenantId == Tenant.AgnosticTenantId)
return sourceTenantId == Tenant.AgnosticTenantId && ambientTenantId == Tenant.AgnosticTenantId;
return CanReplace(existingTenantId, ambientTenantId);
}
/// <summary>
/// Restricts <paramref name="queryable"/> to rows visible to <paramref name="ambientTenantId"/>,
/// unless <paramref name="tenantAgnostic"/> is set (EF <c>IgnoreQueryFilters</c>).

View file

@ -54,8 +54,16 @@ internal static class DbExceptionClassifier
if (typeName.Contains("MySql", StringComparison.OrdinalIgnoreCase) && errorNumbers.Contains(1062))
return true;
if (typeName.Contains("Sqlite", StringComparison.OrdinalIgnoreCase) && errorNumbers.Any(number => number is 19 or 1555 or 2067))
return true;
if (typeName.Contains("Sqlite", StringComparison.OrdinalIgnoreCase))
{
// Microsoft.Data.Sqlite exposes the extended result code when it is
// available. It must take precedence over the base code: a base 19
// can describe a non-duplicate constraint such as NOT NULL.
if (HasProperty(exception, "SqliteExtendedErrorCode"))
return GetIntProperty(exception, "SqliteExtendedErrorCode") is 1555 or 2067;
return GetIntProperty(exception, "SqliteErrorCode") == 19;
}
if (typeName.Contains("Oracle", StringComparison.OrdinalIgnoreCase) && errorNumbers.Contains(1))
return true;
@ -114,6 +122,8 @@ internal static class DbExceptionClassifier
};
}
private static bool HasProperty(object source, string name) => source.GetType().GetProperty(name) is not null;
private static string? GetStringProperty(object source, string name)
{
return source.GetType().GetProperty(name)?.GetValue(source) as string;

View file

@ -19,8 +19,8 @@ namespace Elsa.Persistence.EFCore;
[PublicAPI]
public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextFactory, IServiceProvider serviceProvider) where TDbContext : DbContext where TEntity : class, new()
{
private const int SqlServerBulkWriteMaxRetryCount = 3;
private static readonly TimeSpan SqlServerBulkWriteBaseDelay = TimeSpan.FromMilliseconds(50);
private const int SqlServerWriteMaxRetryCount = 3;
private static readonly TimeSpan SqlServerWriteBaseDelay = TimeSpan.FromMilliseconds(50);
// ReSharper disable once StaticMemberInGenericType
// Justification: This is a static member that is used to ensure that only one thread can access the database for TEntity at a time.
@ -93,7 +93,7 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
if (entityList.Count == 0)
return;
await ExecuteBulkWriteWithSqlServerRetryAsync(async (dbContext, ct) =>
await ExecuteSqlServerWriteWithRetryAsync(async (dbContext, ct) =>
{
if (onSaving != null)
{
@ -190,7 +190,7 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
var tenantId = serviceProvider.GetRequiredService<ITenantAccessor>().TenantId;
await ExecuteBulkWriteWithSqlServerRetryAsync(async (dbContext, ct) =>
await ExecuteSqlServerWriteWithRetryAsync(async (dbContext, ct) =>
{
if (onSaving != null)
{
@ -238,7 +238,32 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
await handler.HandleAsync(context);
}
private async Task ExecuteBulkWriteWithSqlServerRetryAsync(
/// <summary>
/// Executes a database operation and passes failures through the configured database
/// exception handler.
/// </summary>
/// <typeparam name="TResult">The operation result type.</typeparam>
/// <param name="operation">The database operation to execute.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <param name="shouldHandle">A predicate that excludes exceptions which are not database failures.</param>
/// <returns>The result returned by <paramref name="operation"/>.</returns>
internal async Task<TResult> ExecuteWithDbExceptionHandlingAsync<TResult>(
Func<Task<TResult>> operation,
CancellationToken cancellationToken = default,
Func<Exception, bool>? shouldHandle = null)
{
try
{
return await operation();
}
catch (Exception exception) when (shouldHandle is null || shouldHandle(exception))
{
await HandleDbExceptionAsync(exception, cancellationToken);
throw;
}
}
internal async Task ExecuteSqlServerWriteWithRetryAsync(
Func<TDbContext, CancellationToken, Task> operation,
CancellationToken cancellationToken)
{
@ -255,9 +280,9 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
}
catch (Exception ex)
{
if (ShouldRetrySqlServerBulkWrite(providerName, ex, attempt, cancellationToken))
if (ShouldRetrySqlServerWrite(providerName, ex, attempt, cancellationToken))
{
await Task.Delay(GetSqlServerBulkWriteRetryDelay(attempt), cancellationToken);
await Task.Delay(GetSqlServerWriteRetryDelay(attempt), cancellationToken);
continue;
}
@ -266,15 +291,15 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
}
}
private static bool ShouldRetrySqlServerBulkWrite(string providerName, Exception exception, int attempt, CancellationToken cancellationToken)
private static bool ShouldRetrySqlServerWrite(string providerName, Exception exception, int attempt, CancellationToken cancellationToken)
{
return attempt < SqlServerBulkWriteMaxRetryCount
return attempt < SqlServerWriteMaxRetryCount
&& !cancellationToken.IsCancellationRequested
&& exception is not OperationCanceledException
&& DbExceptionClassifier.IsSqlServerTransient(providerName, exception);
}
private static TimeSpan GetSqlServerBulkWriteRetryDelay(int attempt) => TimeSpan.FromMilliseconds(SqlServerBulkWriteBaseDelay.TotalMilliseconds * (attempt + 1));
private static TimeSpan GetSqlServerWriteRetryDelay(int attempt) => TimeSpan.FromMilliseconds(SqlServerWriteBaseDelay.TotalMilliseconds * (attempt + 1));
/// <summary>
/// Updates the entity.

View file

@ -1,37 +1,87 @@
using System.Data.Common;
using System.Text.Json;
using Elsa.Alterations.Core.Contracts;
using Elsa.Alterations.Core.Entities;
using Elsa.Alterations.Core.Filters;
using Elsa.Alterations.Core.Models;
using Elsa.Alterations.Core.Stores;
using Elsa.Tenants.Options;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Open.Linq.AsyncExtensions;
namespace Elsa.Persistence.EFCore.Modules.Alterations;
/// <summary>
/// An EF Core implementation of <see cref="IAlterationPlanStore"/>.
/// An EF Core implementation of <see cref="IAlterationJobStore"/>.
/// </summary>
public class EFCoreAlterationJobStore : IAlterationJobStore
{
private readonly EntityStore<AlterationsElsaDbContext, AlterationJob> _store;
private readonly bool _tenantEnabled;
/// <summary>
/// Constructor.
/// </summary>
public EFCoreAlterationJobStore(EntityStore<AlterationsElsaDbContext, AlterationJob> store)
: this(store, Options.Create(new TenantsOptions()))
{
}
/// <summary>
/// Constructor used by dependency injection. Direct construction through the legacy
/// overload keeps tenancy-aware upsert disabled for compatibility.
/// </summary>
[ActivatorUtilitiesConstructor]
public EFCoreAlterationJobStore(EntityStore<AlterationsElsaDbContext, AlterationJob> store, IOptions<TenantsOptions> tenantsOptions)
{
_store = store;
_tenantEnabled = tenantsOptions.Value.IsEnabled;
}
/// <inheritdoc />
public async Task SaveAsync(AlterationJob record, CancellationToken cancellationToken = default)
{
await _store.SaveAsync(record, OnSaveAsync, cancellationToken);
if (!_tenantEnabled)
{
await _store.SaveAsync(record, OnSaveAsync, cancellationToken);
return;
}
await using var dbContext = await _store.CreateDbContextAsync(cancellationToken);
await UpsertAsync(dbContext, record, cancellationToken);
}
/// <inheritdoc />
public async Task SaveManyAsync(IEnumerable<AlterationJob> jobs, CancellationToken cancellationToken = default)
{
await _store.SaveManyAsync(jobs, OnSaveAsync, cancellationToken);
if (!_tenantEnabled)
{
await _store.SaveManyAsync(jobs, OnSaveAsync, cancellationToken);
return;
}
var list = jobs.OrderBy(job => job.Id, StringComparer.Ordinal).ToList();
if (list.Count == 0)
return;
await _store.ExecuteWithDbExceptionHandlingAsync(
async () =>
{
await _store.ExecuteSqlServerWriteWithRetryAsync(async (dbContext, ct) =>
{
await using var transaction = await dbContext.Database.BeginTransactionAsync(ct);
foreach (var job in list)
await UpsertAsync(dbContext, job, ct, handleDbExceptions: false);
await transaction.CommitAsync(ct);
}, cancellationToken);
return true;
},
cancellationToken,
IsDatabaseException);
}
/// <inheritdoc />
@ -58,9 +108,69 @@ public class EFCoreAlterationJobStore : IAlterationJobStore
return await _store.CountAsync(queryable => Filter(queryable, filter), cancellationToken);
}
private static ValueTask OnSaveAsync(AlterationsElsaDbContext elsaDbContext, AlterationJob entity, CancellationToken cancellationToken)
private async Task UpsertAsync(
AlterationsElsaDbContext dbContext,
AlterationJob record,
CancellationToken cancellationToken,
bool handleDbExceptions = true)
{
var ambientTenantId = AlterationTenantOwnedUpsert.AmbientTenantId(dbContext);
AlterationTenantOwnedUpsert.StampTenantId(record, ambientTenantId);
OnSave(dbContext, record);
var planId = record.PlanId;
var workflowInstanceId = record.WorkflowInstanceId;
var status = record.Status;
var createdAt = record.CreatedAt;
var startedAt = record.StartedAt;
var completedAt = record.CompletedAt;
var serializedLog = dbContext.Entry(record).Property<string>("SerializedLog").CurrentValue;
var query = dbContext.Set<AlterationJob>()
.IgnoreQueryFilters()
.Where(AlterationTenantOwnedUpsert.OwnedId<AlterationJob>(record.Id, record.TenantId, ambientTenantId));
// Inline lambda so net8/net9 bind SetPropertyCalls and net10 binds UpdateSettersBuilder.
Task<int> UpdateOwnedAsync() => query.ExecuteUpdateAsync(
setters => setters
.SetProperty(job => job.PlanId, planId)
.SetProperty(job => job.WorkflowInstanceId, workflowInstanceId)
.SetProperty(job => job.Status, status)
.SetProperty(job => job.CreatedAt, createdAt)
.SetProperty(job => job.StartedAt, startedAt)
.SetProperty(job => job.CompletedAt, completedAt)
.SetProperty(job => EF.Property<string>(job, "SerializedLog"), serializedLog),
cancellationToken);
Task<TResult> ExecuteWriteAsync<TResult>(Func<Task<TResult>> operation) =>
handleDbExceptions
? _store.ExecuteWithDbExceptionHandlingAsync(operation, cancellationToken)
: operation();
var updated = await ExecuteWriteAsync(UpdateOwnedAsync);
if (updated == 0)
{
var inserted = await ExecuteWriteAsync(
() => AlterationTenantOwnedUpsert.InsertIfAbsentAsync(dbContext, record, cancellationToken));
if (!inserted)
{
var retried = await ExecuteWriteAsync(UpdateOwnedAsync);
if (retried == 0)
throw AlterationStoreConflict.HiddenJobId(record.Id);
}
}
}
private static void OnSave(AlterationsElsaDbContext elsaDbContext, AlterationJob entity)
{
elsaDbContext.Entry(entity).Property("SerializedLog").CurrentValue = JsonSerializer.Serialize(entity.Log);
}
private static ValueTask OnSaveAsync(AlterationsElsaDbContext dbContext, AlterationJob entity, CancellationToken cancellationToken)
{
OnSave(dbContext, entity);
return default;
}
@ -74,6 +184,9 @@ public class EFCoreAlterationJobStore : IAlterationJobStore
return default;
}
private static IQueryable<AlterationJob> Filter(IQueryable<AlterationJob> queryable, AlterationJobFilter filter) => filter.Apply(queryable);
}
private static bool IsDatabaseException(Exception exception) =>
exception is DbException or DbUpdateException || exception.InnerException is DbException;
}

View file

@ -4,6 +4,11 @@ using Elsa.Alterations.Core.Contracts;
using Elsa.Alterations.Core.Entities;
using Elsa.Alterations.Core.Filters;
using Elsa.Alterations.Core.Models;
using Elsa.Alterations.Core.Stores;
using Elsa.Tenants.Options;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
namespace Elsa.Persistence.EFCore.Modules.Alterations;
@ -14,20 +19,84 @@ public class EFCoreAlterationPlanStore : IAlterationPlanStore
{
private readonly EntityStore<AlterationsElsaDbContext, AlterationPlan> _store;
private readonly IAlterationSerializer _alterationSerializer;
private readonly bool _tenantEnabled;
/// <summary>
/// Constructor.
/// </summary>
public EFCoreAlterationPlanStore(EntityStore<AlterationsElsaDbContext, AlterationPlan> store, IAlterationSerializer alterationSerializer)
public EFCoreAlterationPlanStore(
EntityStore<AlterationsElsaDbContext, AlterationPlan> store,
IAlterationSerializer alterationSerializer)
: this(store, alterationSerializer, Options.Create(new TenantsOptions()))
{
}
/// <summary>
/// Constructor used by dependency injection. Direct construction through the legacy
/// overload keeps tenancy-aware upsert disabled for compatibility.
/// </summary>
[ActivatorUtilitiesConstructor]
public EFCoreAlterationPlanStore(
EntityStore<AlterationsElsaDbContext, AlterationPlan> store,
IAlterationSerializer alterationSerializer,
IOptions<TenantsOptions> tenantsOptions)
{
_store = store;
_alterationSerializer = alterationSerializer;
_tenantEnabled = tenantsOptions.Value.IsEnabled;
}
/// <inheritdoc />
public async Task SaveAsync(AlterationPlan record, CancellationToken cancellationToken = default)
{
await _store.SaveAsync(record, OnSaveAsync, cancellationToken);
if (!_tenantEnabled)
{
await _store.SaveAsync(record, OnSaveAsync, cancellationToken);
return;
}
await using var dbContext = await _store.CreateDbContextAsync(cancellationToken);
var ambientTenantId = AlterationTenantOwnedUpsert.AmbientTenantId(dbContext);
AlterationTenantOwnedUpsert.StampTenantId(record, ambientTenantId);
OnSave(dbContext, record);
var status = record.Status;
var createdAt = record.CreatedAt;
var startedAt = record.StartedAt;
var completedAt = record.CompletedAt;
var serializedAlterations = dbContext.Entry(record).Property<string>("SerializedAlterations").CurrentValue;
var serializedFilter = dbContext.Entry(record).Property<string>("SerializedWorkflowInstanceFilter").CurrentValue;
var query = dbContext.Set<AlterationPlan>()
.IgnoreQueryFilters()
.Where(AlterationTenantOwnedUpsert.OwnedId<AlterationPlan>(record.Id, record.TenantId, ambientTenantId));
// Inline lambda so net8/net9 bind SetPropertyCalls and net10 binds UpdateSettersBuilder.
Task<int> UpdateOwnedAsync() => query.ExecuteUpdateAsync(
setters => setters
.SetProperty(plan => plan.Status, status)
.SetProperty(plan => plan.CreatedAt, createdAt)
.SetProperty(plan => plan.StartedAt, startedAt)
.SetProperty(plan => plan.CompletedAt, completedAt)
.SetProperty(plan => EF.Property<string>(plan, "SerializedAlterations"), serializedAlterations)
.SetProperty(plan => EF.Property<string>(plan, "SerializedWorkflowInstanceFilter"), serializedFilter),
cancellationToken);
var updated = await _store.ExecuteWithDbExceptionHandlingAsync(UpdateOwnedAsync, cancellationToken);
if (updated == 0)
{
var inserted = await _store.ExecuteWithDbExceptionHandlingAsync(
() => AlterationTenantOwnedUpsert.InsertIfAbsentAsync(dbContext, record, cancellationToken),
cancellationToken);
if (!inserted)
{
var retried = await _store.ExecuteWithDbExceptionHandlingAsync(UpdateOwnedAsync, cancellationToken);
if (retried == 0)
throw AlterationStoreConflict.HiddenPlanId(record.Id);
}
}
}
/// <inheritdoc />
@ -43,10 +112,15 @@ public class EFCoreAlterationPlanStore : IAlterationPlanStore
return await _store.CountAsync(filter.Apply, cancellationToken);
}
private ValueTask OnSaveAsync(AlterationsElsaDbContext elsaDbContext, AlterationPlan entity, CancellationToken cancellationToken)
private void OnSave(AlterationsElsaDbContext elsaDbContext, AlterationPlan entity)
{
elsaDbContext.Entry(entity).Property("SerializedAlterations").CurrentValue = _alterationSerializer.SerializeMany(entity.Alterations);
elsaDbContext.Entry(entity).Property("SerializedWorkflowInstanceFilter").CurrentValue = JsonSerializer.Serialize(entity.WorkflowInstanceFilter);
}
private ValueTask OnSaveAsync(AlterationsElsaDbContext dbContext, AlterationPlan entity, CancellationToken cancellationToken)
{
OnSave(dbContext, entity);
return default;
}
@ -63,4 +137,4 @@ public class EFCoreAlterationPlanStore : IAlterationPlanStore
return default;
}
}
}

View file

@ -0,0 +1,62 @@
using System.Linq.Expressions;
using Elsa.Common.Entities;
using Elsa.Common.Multitenancy;
using Microsoft.EntityFrameworkCore;
namespace Elsa.Persistence.EFCore.Modules.Alterations;
/// <summary>
/// Compare-and-swap upsert for Alterations EF writes. Ownership is decided in the UPDATE
/// predicate itself (no pre-read guard). A 0-row update falls through to INSERT; a PK
/// collision is the unowned-row refuse path.
/// </summary>
internal static class AlterationTenantOwnedUpsert
{
public static void StampTenantId(Entity entity, string? ambientTenantId)
{
if (entity.TenantId == Tenant.AgnosticTenantId)
return;
if (entity.TenantId == null && ambientTenantId != null)
entity.TenantId = ambientTenantId;
}
public static string AmbientTenantId(AlterationsElsaDbContext dbContext) =>
dbContext.TenantId ?? Tenant.DefaultTenantId;
/// <summary>
/// EF translation of <see cref="TenantVisibility.CanReplaceOwnedRow"/>. Named UPDATE
/// matches <c>Target.TenantId</c> to the ambient writer, not a forged source
/// <c>TenantId</c>. <c>*</c> updates still require ambient and source <c>*</c>.
/// </summary>
public static Expression<Func<TEntity, bool>> OwnedId<TEntity>(string id, string? sourceTenantId, string ambientTenantId)
where TEntity : Entity
{
return entity => entity.Id == id && (
(entity.TenantId == Tenant.AgnosticTenantId && sourceTenantId == Tenant.AgnosticTenantId && ambientTenantId == Tenant.AgnosticTenantId)
|| (entity.TenantId != Tenant.AgnosticTenantId && (
entity.TenantId == ambientTenantId
|| (entity.TenantId == null && ambientTenantId == Tenant.DefaultTenantId))));
}
public static async Task<bool> InsertIfAbsentAsync<TEntity>(
AlterationsElsaDbContext dbContext,
TEntity entity,
CancellationToken cancellationToken)
where TEntity : Entity
{
dbContext.Entry(entity).State = EntityState.Added;
try
{
await dbContext.SaveChangesAsync(cancellationToken);
}
catch (DbUpdateException exception) when (DbExceptionClassifier.IsDuplicateKey(exception))
{
dbContext.Entry(entity).State = EntityState.Detached;
return false;
}
return true;
}
}

View file

@ -13,6 +13,7 @@ using Elsa.Tenants.Options;
using Elsa.Testing.Shared.Multitenancy;
using Microsoft.Data.Sqlite;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Alterations.Persistence.ConformanceTests;
@ -50,32 +51,114 @@ public sealed class AlterationStoreScenario(
() => ValueTask.CompletedTask));
}
public static async Task<AlterationStoreScenario> CreateSqliteAsync()
public static Task<AlterationStoreScenario> CreateSqliteAsync() =>
CreateSqliteAsync("tenant-a");
public static Task<AlterationStoreScenario> CreateSqliteAsync(string tenantId) =>
CreateSqliteAsync(tenantId, Path.Join(Path.GetTempPath(), $"elsa-alterations-conformance-{Guid.NewGuid():N}.db"), ownsDatabaseFile: true);
public static Task<AlterationStoreScenario> CreateSqliteAsync(
string tenantId,
bool tenantsEnabled,
DbCommandInterceptor? commandInterceptor = null,
IDbExceptionHandler? dbExceptionHandler = null,
DbTransactionInterceptor? transactionInterceptor = null) =>
CreateSqliteAsync(
tenantId,
Path.Join(Path.GetTempPath(), $"elsa-alterations-conformance-{Guid.NewGuid():N}.db"),
ownsDatabaseFile: true,
tenantsEnabled,
commandInterceptor,
dbExceptionHandler,
transactionInterceptor);
/// <summary>
/// Two EF/SQLite hosts that share a file and keep separate ambient tenants so concurrent
/// Save/SaveMany calls do not mutate a single <see cref="ITenantAccessor"/>.
/// </summary>
public static async Task<SqliteOwnershipPair> CreateSqlitePairAsync(
string firstTenantId,
string secondTenantId,
DbCommandInterceptor? commandInterceptor = null,
DbTransactionInterceptor? transactionInterceptor = null)
{
var databasePath = Path.Join(Path.GetTempPath(), $"elsa-alterations-conformance-{Guid.NewGuid():N}.db");
var tenantAccessor = new TestTenantAccessor("tenant-a");
var databasePath = Path.Join(Path.GetTempPath(), $"elsa-alterations-ownership-{Guid.NewGuid():N}.db");
var first = await CreateSqliteAsync(
firstTenantId,
databasePath,
ownsDatabaseFile: false,
commandInterceptor: commandInterceptor,
transactionInterceptor: transactionInterceptor);
try
{
var second = await CreateSqliteAsync(
secondTenantId,
databasePath,
ownsDatabaseFile: false,
commandInterceptor: commandInterceptor,
transactionInterceptor: transactionInterceptor);
return new SqliteOwnershipPair(first, second, databasePath);
}
catch
{
await first.DisposeAsync();
throw;
}
}
public static Task<AlterationStoreScenario> CreateSqliteAsync(
string tenantId,
string databasePath,
bool ownsDatabaseFile,
bool tenantsEnabled = true,
DbCommandInterceptor? commandInterceptor = null,
IDbExceptionHandler? dbExceptionHandler = null,
DbTransactionInterceptor? transactionInterceptor = null) =>
CreateSqliteHostAsync(tenantId, databasePath, ownsDatabaseFile, tenantsEnabled, commandInterceptor, dbExceptionHandler, transactionInterceptor);
private static async Task<AlterationStoreScenario> CreateSqliteHostAsync(
string tenantId,
string databasePath,
bool ownsDatabaseFile,
bool tenantsEnabled,
DbCommandInterceptor? commandInterceptor,
IDbExceptionHandler? dbExceptionHandler,
DbTransactionInterceptor? transactionInterceptor)
{
var tenantAccessor = new TestTenantAccessor(tenantId);
ServiceProvider? services = null;
IServiceScope? scope = null;
try
{
var migrationsAssembly = typeof(AlterationsDbContextFactories).Assembly;
services = new ServiceCollection()
var serviceCollection = new ServiceCollection()
.AddLogging()
.AddSingleton<ITenantAccessor>(tenantAccessor)
.AddSingleton<IAlterationSerializer, ConformanceAlterationSerializer>()
.Configure<TenantsOptions>(options => options.IsEnabled = true)
.Configure<TenantsOptions>(options => options.IsEnabled = tenantsEnabled)
.AddScoped<IEntitySavingHandler, ApplyTenantId>()
.AddScoped<IEntityModelCreatingHandler, SetTenantIdFilter>()
.AddSqliteEntityModelCreatingHandlers()
.AddDbContextFactory<AlterationsElsaDbContext>((_, builder) =>
builder.UseElsaSqlite(migrationsAssembly, $"Data Source={databasePath};Default Timeout=30"))
{
builder.UseElsaSqlite(migrationsAssembly, $"Data Source={databasePath};Default Timeout=30");
builder.EnableServiceProviderCaching(false);
if (commandInterceptor is not null)
builder.AddInterceptors(commandInterceptor);
if (transactionInterceptor is not null)
builder.AddInterceptors(transactionInterceptor);
})
.Decorate<IDbContextFactory<AlterationsElsaDbContext>, TenantAwareDbContextFactory<AlterationsElsaDbContext>>()
.AddScoped<EntityStore<AlterationsElsaDbContext, AlterationPlan>>()
.AddScoped<EntityStore<AlterationsElsaDbContext, AlterationJob>>()
.AddScoped<EFCoreAlterationPlanStore>()
.AddScoped<EFCoreAlterationJobStore>()
.BuildServiceProvider();
.AddScoped<EFCoreAlterationJobStore>();
if (dbExceptionHandler is not null)
serviceCollection.AddSingleton(dbExceptionHandler);
services = serviceCollection.BuildServiceProvider();
await using (var dbContext = await services.GetRequiredService<IDbContextFactory<AlterationsElsaDbContext>>().CreateDbContextAsync())
await dbContext.Database.EnsureCreatedAsync();
@ -91,8 +174,11 @@ public sealed class AlterationStoreScenario(
{
scope.Dispose();
await services.DisposeAsync();
SqliteConnection.ClearAllPools();
File.Delete(databasePath);
if (ownsDatabaseFile)
{
SqliteConnection.ClearAllPools();
File.Delete(databasePath);
}
});
}
catch
@ -100,8 +186,12 @@ public sealed class AlterationStoreScenario(
scope?.Dispose();
if (services is not null)
await services.DisposeAsync();
SqliteConnection.ClearAllPools();
File.Delete(databasePath);
if (ownsDatabaseFile)
{
SqliteConnection.ClearAllPools();
File.Delete(databasePath);
}
throw;
}
}
@ -123,6 +213,24 @@ public sealed class AlterationStoreScenario(
}
}
/// <summary>
/// Two SQLite alteration-store hosts that share one database file.
/// </summary>
public sealed class SqliteOwnershipPair(AlterationStoreScenario first, AlterationStoreScenario second, string databasePath) : IAsyncDisposable
{
public AlterationStoreScenario First { get; } = first;
public AlterationStoreScenario Second { get; } = second;
public string DatabasePath { get; } = databasePath;
public async ValueTask DisposeAsync()
{
await First.DisposeAsync();
await Second.DisposeAsync();
SqliteConnection.ClearAllPools();
File.Delete(DatabasePath);
}
}
/// <summary>
/// Minimal <see cref="IAlteration"/> used to lock save/load of serialized plan alterations.
/// </summary>

View file

@ -0,0 +1,673 @@
using System.Data.Common;
using Elsa.Alterations.Core.Entities;
using Elsa.Alterations.Core.Enums;
using Elsa.Alterations.Core.Filters;
using Elsa.Alterations.Core.Models;
using Elsa.Common.Multitenancy;
using Elsa.Persistence.EFCore;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
using Microsoft.Extensions.Logging;
namespace Elsa.Alterations.Persistence.ConformanceTests;
/// <summary>
/// EF Alterations Save/SaveMany must refuse ID collisions atomically (no pre-read guard).
/// Kept out of the Memory/EF conformance matrix so that suite stays on read/stamp/query/round-trip.
/// </summary>
[Collection(AlterationStoreSqliteConformanceCollection.Name)]
public sealed class EFCoreAlterationStoreTenantOwnershipTests
{
[Fact]
public async Task SaveAsync_WhenOtherNamedTenantOwnsPlanId_ThrowsAndLeavesOwnerAndPayload()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Plans.SaveAsync(Plan("shared", "tenant-a", AlterationPlanStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() =>
owner.Plans.SaveAsync(Plan("shared", "tenant-b", AlterationPlanStatus.Completed, "stolen")));
Assert.Contains("shared", ex.Message);
}
var remaining = await owner.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" });
AssertUnchangedPlan(remaining, "tenant-a", AlterationPlanStatus.Running, "original");
}
[Fact]
public async Task SaveAsync_WhenAgnosticPlanExists_NamedTenantThrowsAndLeavesOwnerAndPayload()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(Tenant.AgnosticTenantId);
await scenario.Plans.SaveAsync(Plan("shared", Tenant.AgnosticTenantId, AlterationPlanStatus.Running, "original"));
using (scenario.UseTenant("tenant-b"))
{
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() =>
scenario.Plans.SaveAsync(Plan("shared", "tenant-b", AlterationPlanStatus.Completed, "stolen")));
Assert.Contains("shared", ex.Message);
}
using (scenario.UseTenant(Tenant.AgnosticTenantId))
{
var remaining = await scenario.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" });
AssertUnchangedPlan(remaining, Tenant.AgnosticTenantId, AlterationPlanStatus.Running, "original");
}
}
[Fact]
public async Task SaveAsync_WhenAmbientForgesOwnerTenantIdOnPlan_ThrowsAndLeavesOwnerAndPayload()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Plans.SaveAsync(Plan("shared", "tenant-a", AlterationPlanStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() =>
owner.Plans.SaveAsync(Plan("shared", "tenant-a", AlterationPlanStatus.Completed, "stolen")));
Assert.Contains("shared", ex.Message);
}
var remaining = await owner.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" });
AssertUnchangedPlan(remaining, "tenant-a", AlterationPlanStatus.Running, "original");
}
[Fact]
public async Task SaveAsync_WhenSameTenantOwnsPlanId_UpdatesPayload()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await scenario.Plans.SaveAsync(Plan("plan-a", "tenant-a", AlterationPlanStatus.Pending, "before"));
await scenario.Plans.SaveAsync(Plan("plan-a", "tenant-a", AlterationPlanStatus.Completed, "after"));
var found = await scenario.Plans.FindAsync(new AlterationPlanFilter { Id = "plan-a" });
AssertUnchangedPlan(found, "tenant-a", AlterationPlanStatus.Completed, "after");
}
[Fact]
public async Task SaveAsync_WhenIncomingPlanTenantDiffers_PreservesExistingTenantId()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await scenario.Plans.SaveAsync(Plan("plan-a", "tenant-a", AlterationPlanStatus.Pending, "before"));
await scenario.Plans.SaveAsync(Plan("plan-a", "tenant-b", AlterationPlanStatus.Completed, "after"));
var found = await scenario.Plans.FindAsync(new AlterationPlanFilter { Id = "plan-a" });
AssertUnchangedPlan(found, "tenant-a", AlterationPlanStatus.Completed, "after");
}
[Fact]
public async Task SaveAsync_WhenAmbientIsAgnostic_UpdatesAgnosticPlan()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(Tenant.AgnosticTenantId);
await scenario.Plans.SaveAsync(Plan("shared", Tenant.AgnosticTenantId, AlterationPlanStatus.Pending, "before"));
await scenario.Plans.SaveAsync(Plan("shared", Tenant.AgnosticTenantId, AlterationPlanStatus.Completed, "after"));
var found = await scenario.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" });
AssertUnchangedPlan(found, Tenant.AgnosticTenantId, AlterationPlanStatus.Completed, "after");
}
[Fact]
public async Task SaveAsync_WhenOtherNamedTenantOwnsJobId_ThrowsAndLeavesOwnerAndPayload()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Jobs.SaveAsync(Job("shared", "tenant-a", AlterationJobStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() =>
owner.Jobs.SaveAsync(Job("shared", "tenant-b", AlterationJobStatus.Completed, "stolen")));
Assert.Contains("shared", ex.Message);
}
var remaining = await owner.Jobs.FindAsync(new AlterationJobFilter { Id = "shared" });
AssertUnchangedJob(remaining, "tenant-a", AlterationJobStatus.Running, "original");
}
[Fact]
public async Task SaveAsync_WhenAmbientForgesOwnerTenantIdOnJob_ThrowsAndLeavesOwnerAndPayload()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Jobs.SaveAsync(Job("shared", "tenant-a", AlterationJobStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() =>
owner.Jobs.SaveAsync(Job("shared", "tenant-a", AlterationJobStatus.Completed, "stolen")));
Assert.Contains("shared", ex.Message);
}
var remaining = await owner.Jobs.FindAsync(new AlterationJobFilter { Id = "shared" });
AssertUnchangedJob(remaining, "tenant-a", AlterationJobStatus.Running, "original");
}
[Fact]
public async Task SaveManyAsync_WhenOtherNamedTenantOwnsJobId_ThrowsAndLeavesOwnerAndPayload()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Jobs.SaveAsync(Job("shared", "tenant-a", AlterationJobStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
await Assert.ThrowsAsync<InvalidOperationException>(() =>
owner.Jobs.SaveManyAsync([Job("shared", "tenant-b", AlterationJobStatus.Completed, "stolen")]));
}
var remaining = await owner.Jobs.FindAsync(new AlterationJobFilter { Id = "shared" });
AssertUnchangedJob(remaining, "tenant-a", AlterationJobStatus.Running, "original");
}
[Fact]
public async Task SaveManyAsync_WhenAmbientForgesOwnerTenantIdOnJob_ThrowsAndLeavesOwnerAndPayload()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Jobs.SaveAsync(Job("shared", "tenant-a", AlterationJobStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() =>
owner.Jobs.SaveManyAsync([Job("shared", "tenant-a", AlterationJobStatus.Completed, "stolen")]));
Assert.Contains("shared", ex.Message);
}
var remaining = await owner.Jobs.FindAsync(new AlterationJobFilter { Id = "shared" });
AssertUnchangedJob(remaining, "tenant-a", AlterationJobStatus.Running, "original");
}
[Fact]
public async Task SaveManyAsync_WhenAgnosticJobExists_NamedTenantThrowsAndLeavesOwnerAndPayload()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(Tenant.AgnosticTenantId);
await scenario.Jobs.SaveAsync(Job("shared", Tenant.AgnosticTenantId, AlterationJobStatus.Running, "original"));
using (scenario.UseTenant("tenant-b"))
{
await Assert.ThrowsAsync<InvalidOperationException>(() =>
scenario.Jobs.SaveManyAsync([Job("shared", "tenant-b", AlterationJobStatus.Completed, "stolen")]));
}
using (scenario.UseTenant(Tenant.AgnosticTenantId))
{
var remaining = await scenario.Jobs.FindAsync(new AlterationJobFilter { Id = "shared" });
AssertUnchangedJob(remaining, Tenant.AgnosticTenantId, AlterationJobStatus.Running, "original");
}
}
[Fact]
public async Task SaveManyAsync_WhenSameTenantOwnsJobId_UpdatesPayload()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await scenario.Jobs.SaveAsync(Job("job-a", "tenant-a", AlterationJobStatus.Pending, "before"));
await scenario.Jobs.SaveManyAsync([Job("job-a", "tenant-a", AlterationJobStatus.Completed, "after")]);
var found = await scenario.Jobs.FindAsync(new AlterationJobFilter { Id = "job-a" });
AssertUnchangedJob(found, "tenant-a", AlterationJobStatus.Completed, "after");
}
[Fact]
public async Task SaveManyAsync_WhenIncomingJobTenantDiffers_PreservesExistingTenantId()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await scenario.Jobs.SaveAsync(Job("job-a", "tenant-a", AlterationJobStatus.Pending, "before"));
await scenario.Jobs.SaveManyAsync([Job("job-a", "tenant-b", AlterationJobStatus.Completed, "after")]);
var found = await scenario.Jobs.FindAsync(new AlterationJobFilter { Id = "job-a" });
AssertUnchangedJob(found, "tenant-a", AlterationJobStatus.Completed, "after");
}
[Fact]
public async Task SaveAsync_WhenTenancyIsDisabled_UpdatesNamedAndAgnosticPlansThroughLegacyStore()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync("tenant-b", tenantsEnabled: false);
await scenario.Plans.SaveAsync(Plan("named", "tenant-a", AlterationPlanStatus.Pending, "before"));
await scenario.Plans.SaveAsync(Plan("agnostic", Tenant.AgnosticTenantId, AlterationPlanStatus.Pending, "before"));
await scenario.Plans.SaveAsync(Plan("named", "tenant-a", AlterationPlanStatus.Completed, "after"));
await scenario.Plans.SaveAsync(Plan("agnostic", Tenant.AgnosticTenantId, AlterationPlanStatus.Completed, "after"));
AssertUnchangedPlan(await scenario.Plans.FindAsync(new AlterationPlanFilter { Id = "named" }), "tenant-a", AlterationPlanStatus.Completed, "after");
AssertUnchangedPlan(await scenario.Plans.FindAsync(new AlterationPlanFilter { Id = "agnostic" }), Tenant.AgnosticTenantId, AlterationPlanStatus.Completed, "after");
}
[Fact]
public async Task SaveManyAsync_WhenTenancyIsDisabled_UpdatesNamedAndAgnosticJobsThroughLegacyStore()
{
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync("tenant-b", tenantsEnabled: false);
await scenario.Jobs.SaveAsync(Job("named", "tenant-a", AlterationJobStatus.Pending, "before"));
await scenario.Jobs.SaveAsync(Job("agnostic", Tenant.AgnosticTenantId, AlterationJobStatus.Pending, "before"));
await scenario.Jobs.SaveManyAsync(
[
Job("named", "tenant-a", AlterationJobStatus.Completed, "after"),
Job("agnostic", Tenant.AgnosticTenantId, AlterationJobStatus.Completed, "after")
]);
AssertUnchangedJob(await scenario.Jobs.FindAsync(new AlterationJobFilter { Id = "named" }), "tenant-a", AlterationJobStatus.Completed, "after");
AssertUnchangedJob(await scenario.Jobs.FindAsync(new AlterationJobFilter { Id = "agnostic" }), Tenant.AgnosticTenantId, AlterationJobStatus.Completed, "after");
}
[Fact]
public async Task SaveAsync_WhenDirectPlanUpdateFails_InvokesDbExceptionHandler()
{
var handler = new RecordingDbExceptionHandler();
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(
"tenant-a",
tenantsEnabled: true,
commandInterceptor: new ThrowOnAlterationCommand("UPDATE \"AlterationPlans\""),
dbExceptionHandler: handler);
var exception = await Assert.ThrowsAsync<DbUpdateException>(() => scenario.Plans.SaveAsync(Plan(
"handler-plan",
"tenant-a",
AlterationPlanStatus.Pending,
"payload")));
Assert.Same(exception, handler.Exception);
}
[Fact]
public async Task SaveAsync_WhenDirectJobInsertFails_InvokesDbExceptionHandler()
{
var handler = new RecordingDbExceptionHandler();
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(
"tenant-a",
tenantsEnabled: true,
commandInterceptor: new ThrowOnAlterationCommand("INSERT INTO \"AlterationJobs\""),
dbExceptionHandler: handler);
var exception = await Assert.ThrowsAsync<DbUpdateException>(() => scenario.Jobs.SaveAsync(Job(
"handler-job",
"tenant-a",
AlterationJobStatus.Pending,
"payload")));
Assert.Same(exception, handler.Exception);
}
[Fact]
public async Task SaveManyAsync_WhenDirectJobUpdateFails_InvokesDbExceptionHandlerAfterRetryBoundary()
{
var handler = new RecordingDbExceptionHandler();
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(
"tenant-a",
tenantsEnabled: true,
commandInterceptor: new ThrowOnAlterationCommand("UPDATE \"AlterationJobs\""),
dbExceptionHandler: handler);
var exception = await Assert.ThrowsAsync<DbUpdateException>(() => scenario.Jobs.SaveManyAsync([
Job("handler-batch", "tenant-a", AlterationJobStatus.Pending, "payload")
]));
Assert.Same(exception, handler.Exception);
}
[Fact]
public async Task SaveAsync_WhenPlanIdIsHidden_DoesNotSendOwnershipConflictToDbExceptionHandler()
{
var handler = new RecordingDbExceptionHandler();
await using var scenario = await AlterationStoreScenario.CreateSqliteAsync(
"tenant-a",
tenantsEnabled: true,
dbExceptionHandler: handler);
await scenario.Plans.SaveAsync(Plan("handler-conflict", "tenant-a", AlterationPlanStatus.Pending, "owner"));
using (scenario.UseTenant("tenant-b"))
{
await Assert.ThrowsAsync<InvalidOperationException>(() => scenario.Plans.SaveAsync(
Plan("handler-conflict", "tenant-b", AlterationPlanStatus.Completed, "hidden")));
}
Assert.Null(handler.Exception);
}
[Fact]
public async Task SaveAsync_ConcurrentSameTenantPlanId_BothWritersSucceedAndOnePayloadWins()
{
var gate = new GateFirstAlterationUpdates("AlterationPlans");
await using var pair = await AlterationStoreScenario.CreateSqlitePairAsync("tenant-a", "tenant-a", gate);
gate.Arm();
var results = await Task.WhenAll(
Capture(() => pair.First.Plans.SaveAsync(Plan("race", "tenant-a", AlterationPlanStatus.Running, "from-a"))),
Capture(() => pair.Second.Plans.SaveAsync(Plan("race", "tenant-a", AlterationPlanStatus.Completed, "from-b"))));
Assert.All(results, Assert.Null);
Assert.Equal(3, gate.MatchedCommandCount);
using (pair.First.UseTenant("tenant-a"))
{
var found = await pair.First.Plans.FindAsync(new AlterationPlanFilter { Id = "race" });
Assert.NotNull(found);
Assert.Equal("tenant-a", found.TenantId);
Assert.Contains(Payload(found), new[] { "from-a", "from-b" });
}
}
[Fact]
public async Task SaveManyAsync_WhenBatchCollides_RollsBackEarlierInserts()
{
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a");
await owner.Jobs.SaveAsync(Job("owned", "tenant-a", AlterationJobStatus.Running, "original"));
using (owner.UseTenant("tenant-b"))
{
await Assert.ThrowsAsync<InvalidOperationException>(() => owner.Jobs.SaveManyAsync(
[
Job("new-from-b", "tenant-b", AlterationJobStatus.Pending, "should-roll-back"),
Job("owned", "tenant-b", AlterationJobStatus.Completed, "stolen")
]));
Assert.Null(await owner.Jobs.FindAsync(new AlterationJobFilter { Id = "new-from-b" }));
}
var remaining = await owner.Jobs.FindAsync(new AlterationJobFilter { Id = "owned" });
AssertUnchangedJob(remaining, "tenant-a", AlterationJobStatus.Running, "original");
}
[Fact]
public async Task SaveAsync_ConcurrentNamedTenantsOnEmptyPlanId_OneOwnerKeepsPayload()
{
var gate = new GateFirstAlterationUpdates("AlterationPlans");
await using var pair = await AlterationStoreScenario.CreateSqlitePairAsync("tenant-a", "tenant-b", gate);
gate.Arm();
var results = await Task.WhenAll(
Capture(() => pair.First.Plans.SaveAsync(Plan("race", "tenant-a", AlterationPlanStatus.Running, "from-a"))),
Capture(() => pair.Second.Plans.SaveAsync(Plan("race", "tenant-b", AlterationPlanStatus.Completed, "from-b"))));
Assert.Equal(1, results.Count(ex => ex is null));
Assert.Equal(1, results.Count(ex => ex is InvalidOperationException));
Assert.Equal(3, gate.MatchedCommandCount);
var winnerIsA = results[0] is null;
using (pair.First.UseTenant(winnerIsA ? "tenant-a" : "tenant-b"))
{
var remaining = await pair.First.Plans.FindAsync(new AlterationPlanFilter { Id = "race" });
AssertUnchangedPlan(
remaining,
winnerIsA ? "tenant-a" : "tenant-b",
winnerIsA ? AlterationPlanStatus.Running : AlterationPlanStatus.Completed,
winnerIsA ? "from-a" : "from-b");
}
}
[Fact]
public async Task SaveAsync_ConcurrentNamedVersusAgnosticOnExistingStarPlan_PreservesStarPayload()
{
var gate = new GateFirstAlterationUpdates("AlterationPlans");
await using var pair = await AlterationStoreScenario.CreateSqlitePairAsync(Tenant.AgnosticTenantId, "tenant-b", gate);
await pair.First.Plans.SaveAsync(Plan("shared", Tenant.AgnosticTenantId, AlterationPlanStatus.Running, "original"));
gate.Arm();
var results = await Task.WhenAll(
Capture(() => pair.First.Plans.SaveAsync(Plan("shared", Tenant.AgnosticTenantId, AlterationPlanStatus.Failed, "agnostic-update"))),
Capture(() => pair.Second.Plans.SaveAsync(Plan("shared", "tenant-b", AlterationPlanStatus.Completed, "stolen"))));
Assert.IsType<InvalidOperationException>(results[1]);
Assert.Equal(3, gate.MatchedCommandCount);
using (pair.First.UseTenant(Tenant.AgnosticTenantId))
{
var remaining = await pair.First.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" });
Assert.NotNull(remaining);
Assert.Equal(Tenant.AgnosticTenantId, remaining.TenantId);
Assert.NotEqual("stolen", Payload(remaining));
if (results[0] is null)
AssertUnchangedPlan(remaining, Tenant.AgnosticTenantId, AlterationPlanStatus.Failed, "agnostic-update");
else
AssertUnchangedPlan(remaining, Tenant.AgnosticTenantId, AlterationPlanStatus.Running, "original");
}
}
[Fact]
public async Task SaveManyAsync_ConcurrentNamedTenantsOnEmptyJobId_OneOwnerKeepsPayload()
{
var gate = new GateFirstAlterationTransactions();
await using var pair = await AlterationStoreScenario.CreateSqlitePairAsync(
"tenant-a",
"tenant-b",
transactionInterceptor: gate);
gate.Arm();
var results = await Task.WhenAll(
Capture(() => pair.First.Jobs.SaveManyAsync([Job("race", "tenant-a", AlterationJobStatus.Running, "from-a")])),
Capture(() => pair.Second.Jobs.SaveManyAsync([Job("race", "tenant-b", AlterationJobStatus.Completed, "from-b")])));
Assert.Equal(1, results.Count(ex => ex is null));
Assert.Equal(1, results.Count(ex => ex is InvalidOperationException));
Assert.Equal(2, gate.MatchedTransactionCount);
var winnerIsA = results[0] is null;
using (pair.First.UseTenant(winnerIsA ? "tenant-a" : "tenant-b"))
{
var remaining = await pair.First.Jobs.FindAsync(new AlterationJobFilter { Id = "race" });
AssertUnchangedJob(
remaining,
winnerIsA ? "tenant-a" : "tenant-b",
winnerIsA ? AlterationJobStatus.Running : AlterationJobStatus.Completed,
winnerIsA ? "from-a" : "from-b");
}
}
[Fact]
public async Task SaveManyAsync_ConcurrentNamedVersusAgnosticOnExistingStarJob_PreservesStarPayload()
{
await using var pair = await AlterationStoreScenario.CreateSqlitePairAsync(Tenant.AgnosticTenantId, "tenant-b");
await pair.First.Jobs.SaveAsync(Job("shared", Tenant.AgnosticTenantId, AlterationJobStatus.Running, "original"));
var results = await Task.WhenAll(
Capture(() => pair.First.Jobs.SaveManyAsync([Job("shared", Tenant.AgnosticTenantId, AlterationJobStatus.Failed, "agnostic-update")])),
Capture(() => pair.Second.Jobs.SaveManyAsync([Job("shared", "tenant-b", AlterationJobStatus.Completed, "stolen")])));
Assert.NotNull(results[1]);
Assert.IsType<InvalidOperationException>(results[1]);
using (pair.First.UseTenant(Tenant.AgnosticTenantId))
{
var remaining = await pair.First.Jobs.FindAsync(new AlterationJobFilter { Id = "shared" });
Assert.NotNull(remaining);
Assert.Equal(Tenant.AgnosticTenantId, remaining.TenantId);
Assert.NotEqual("stolen", remaining.WorkflowInstanceId);
if (results[0] is null)
AssertUnchangedJob(remaining, Tenant.AgnosticTenantId, AlterationJobStatus.Failed, "agnostic-update");
else
AssertUnchangedJob(remaining, Tenant.AgnosticTenantId, AlterationJobStatus.Running, "original");
}
}
[Fact]
public async Task SaveAsync_ConcurrentNamedTenantsOnExistingNamedPlan_PreservesOriginalOwnerAndPayload()
{
await using var pair = await AlterationStoreScenario.CreateSqlitePairAsync("tenant-b", "tenant-c");
await using var owner = await AlterationStoreScenario.CreateSqliteAsync("tenant-a", pair.DatabasePath, ownsDatabaseFile: false);
await owner.Plans.SaveAsync(Plan("shared", "tenant-a", AlterationPlanStatus.Running, "original"));
var results = await Task.WhenAll(
Capture(() => pair.First.Plans.SaveAsync(Plan("shared", "tenant-b", AlterationPlanStatus.Completed, "from-b"))),
Capture(() => pair.Second.Plans.SaveAsync(Plan("shared", "tenant-c", AlterationPlanStatus.Failed, "from-c"))));
Assert.All(results, ex => Assert.IsType<InvalidOperationException>(ex));
var remaining = await owner.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" });
AssertUnchangedPlan(remaining, "tenant-a", AlterationPlanStatus.Running, "original");
}
private static async Task<Exception?> Capture(Func<Task> action)
{
try
{
await action();
return null;
}
catch (InvalidOperationException ex)
{
return ex;
}
}
private static void AssertUnchangedPlan(AlterationPlan? plan, string tenantId, AlterationPlanStatus status, string payload)
{
Assert.NotNull(plan);
Assert.Equal(tenantId, plan.TenantId);
Assert.Equal(status, plan.Status);
Assert.Equal(payload, Payload(plan));
}
private static void AssertUnchangedJob(AlterationJob? job, string tenantId, AlterationJobStatus status, string payload)
{
Assert.NotNull(job);
Assert.Equal(tenantId, job.TenantId);
Assert.Equal(status, job.Status);
Assert.Equal(payload, job.WorkflowInstanceId);
Assert.Equal(payload, job.Log?.FirstOrDefault()?.Message);
}
private static string Payload(AlterationPlan plan) =>
Assert.IsType<TestAlteration>(Assert.Single(plan.Alterations)).Value;
private static AlterationPlan Plan(string id, string? tenantId, AlterationPlanStatus status, string payload) =>
new()
{
Id = id,
TenantId = tenantId,
Status = status,
CreatedAt = CreatedAt,
Alterations = [new TestAlteration { Kind = "probe", Value = payload }],
WorkflowInstanceFilter = new AlterationWorkflowInstanceFilter { SearchTerm = payload }
};
private static AlterationJob Job(string id, string? tenantId, AlterationJobStatus status, string payload) =>
new()
{
Id = id,
PlanId = "plan-1",
WorkflowInstanceId = payload,
Status = status,
TenantId = tenantId,
CreatedAt = CreatedAt,
Log = [new AlterationLogEntry(payload, LogLevel.Information, CreatedAt, "probe")]
};
private static readonly DateTimeOffset CreatedAt = new(2026, 9, 13, 12, 0, 0, TimeSpan.Zero);
}
/// <summary>
/// Releases the first two relevant alteration UPDATE commands after they complete, so competing
/// Save calls reach their INSERT/retry paths together. The gate is armed explicitly after any
/// setup writes so only the concurrent operation is coordinated.
/// </summary>
public sealed class GateFirstAlterationUpdates(string tableName) : DbCommandInterceptor
{
private readonly TaskCompletionSource<bool> _bothReached = new(TaskCreationOptions.RunContinuationsAsynchronously);
private int _armed;
private int _matchedCommandCount;
public int MatchedCommandCount => Volatile.Read(ref _matchedCommandCount);
public void Arm() => Volatile.Write(ref _armed, 1);
public override async ValueTask<int> NonQueryExecutedAsync(
DbCommand command,
CommandExecutedEventData eventData,
int result,
CancellationToken cancellationToken = default)
{
if (!IsArmedAlterationUpdate(command))
return result;
await CoordinateAsync(cancellationToken);
return result;
}
private bool IsArmedAlterationUpdate(DbCommand command) =>
Volatile.Read(ref _armed) != 0 && command.CommandText.Contains($"UPDATE \"{tableName}\"", StringComparison.OrdinalIgnoreCase);
private async Task CoordinateAsync(CancellationToken cancellationToken)
{
var commandNumber = Interlocked.Increment(ref _matchedCommandCount);
if (commandNumber <= 2)
{
if (commandNumber == 2)
_bothReached.TrySetResult(true);
await _bothReached.Task.WaitAsync(TimeSpan.FromSeconds(15), cancellationToken);
}
}
}
public sealed class GateFirstAlterationTransactions : DbTransactionInterceptor
{
private readonly TaskCompletionSource<bool> _bothReached = new(TaskCreationOptions.RunContinuationsAsynchronously);
private int _armed;
private int _matchedTransactionCount;
public int MatchedTransactionCount => Volatile.Read(ref _matchedTransactionCount);
public void Arm() => Volatile.Write(ref _armed, 1);
public override async ValueTask<InterceptionResult<DbTransaction>> TransactionStartingAsync(
DbConnection connection,
TransactionStartingEventData eventData,
InterceptionResult<DbTransaction> result,
CancellationToken cancellationToken = default)
{
if (Volatile.Read(ref _armed) != 0)
{
var transactionNumber = Interlocked.Increment(ref _matchedTransactionCount);
if (transactionNumber <= 2)
{
if (transactionNumber == 2)
_bothReached.TrySetResult(true);
await _bothReached.Task.WaitAsync(TimeSpan.FromSeconds(15), cancellationToken);
}
}
return result;
}
}
public sealed class ThrowOnAlterationCommand(string commandFragment) : DbCommandInterceptor
{
private void ThrowIfMatched(DbCommand command)
{
if (command.CommandText.Contains(commandFragment, StringComparison.OrdinalIgnoreCase))
throw new DbUpdateException("Forced alteration persistence failure.");
}
public override ValueTask<InterceptionResult<int>> NonQueryExecutingAsync(
DbCommand command,
CommandEventData eventData,
InterceptionResult<int> result,
CancellationToken cancellationToken = default)
{
ThrowIfMatched(command);
return ValueTask.FromResult(result);
}
public override ValueTask<InterceptionResult<DbDataReader>> ReaderExecutingAsync(
DbCommand command,
CommandEventData eventData,
InterceptionResult<DbDataReader> result,
CancellationToken cancellationToken = default)
{
ThrowIfMatched(command);
return ValueTask.FromResult(result);
}
}
public sealed class RecordingDbExceptionHandler : IDbExceptionHandler
{
public Exception? Exception { get; private set; }
public Task HandleAsync(DbUpdateExceptionContext context)
{
Exception = context.Exception;
return Task.CompletedTask;
}
}

View file

@ -95,6 +95,37 @@ public class MemoryAlterationJobStoreTenantIsolationTests
Assert.Equal("tenant-a", remaining.TenantId);
}
[Fact(DisplayName = "SaveAsync refuses a forged owner TenantId from another ambient tenant")]
public async Task SaveAsync_WhenAmbientForgesOwnerTenantId_ThrowsAndLeavesExisting()
{
var backing = new MemoryStore<AlterationJob>();
var tenantA = new MemoryAlterationJobStore(backing, new TestTenantAccessor("tenant-a"));
var tenantB = new MemoryAlterationJobStore(backing, new TestTenantAccessor("tenant-b"));
await tenantA.SaveAsync(Job("shared", "tenant-a"));
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() => tenantB.SaveAsync(Job("shared", "tenant-a")));
var remaining = await tenantA.FindAsync(new AlterationJobFilter { Id = "shared" });
Assert.Contains("shared", ex.Message);
Assert.NotNull(remaining);
Assert.Equal("tenant-a", remaining.TenantId);
}
[Fact(DisplayName = "SaveManyAsync refuses a forged owner TenantId from another ambient tenant")]
public async Task SaveManyAsync_WhenAmbientForgesOwnerTenantId_ThrowsAndLeavesExisting()
{
var backing = new MemoryStore<AlterationJob>();
var tenantA = new MemoryAlterationJobStore(backing, new TestTenantAccessor("tenant-a"));
var tenantB = new MemoryAlterationJobStore(backing, new TestTenantAccessor("tenant-b"));
await tenantA.SaveAsync(Job("shared", "tenant-a"));
await Assert.ThrowsAsync<InvalidOperationException>(() => tenantB.SaveManyAsync([Job("shared", "tenant-a")]));
var remaining = await tenantA.FindAsync(new AlterationJobFilter { Id = "shared" });
Assert.NotNull(remaining);
Assert.Equal("tenant-a", remaining.TenantId);
}
[Fact(DisplayName = "SaveAsync refuses to overwrite a tenant-agnostic row by Id")]
public async Task SaveAsync_WhenAgnosticRowExists_NamedTenantThrowsAndLeavesExisting()
{

View file

@ -69,6 +69,22 @@ public class MemoryAlterationPlanStoreTenantIsolationTests
Assert.Equal("tenant-a", remaining.TenantId);
}
[Fact(DisplayName = "SaveAsync refuses a forged owner TenantId from another ambient tenant")]
public async Task SaveAsync_WhenAmbientForgesOwnerTenantId_ThrowsAndLeavesExisting()
{
var backing = new MemoryStore<AlterationPlan>();
var tenantA = new MemoryAlterationPlanStore(backing, new TestTenantAccessor("tenant-a"));
var tenantB = new MemoryAlterationPlanStore(backing, new TestTenantAccessor("tenant-b"));
await tenantA.SaveAsync(Plan("shared", "tenant-a"));
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() => tenantB.SaveAsync(Plan("shared", "tenant-a")));
var remaining = await tenantA.FindAsync(new AlterationPlanFilter { Id = "shared" });
Assert.Contains("shared", ex.Message);
Assert.NotNull(remaining);
Assert.Equal("tenant-a", remaining.TenantId);
}
[Fact(DisplayName = "SaveAsync refuses to overwrite a tenant-agnostic row by Id")]
public async Task SaveAsync_WhenAgnosticRowExists_NamedTenantThrowsAndLeavesExisting()
{

View file

@ -54,6 +54,40 @@ public class TenantVisibilityTests
Assert.Equal(2, visible.Count);
}
[Theory]
[InlineData("tenant-a", "tenant-a", true)]
[InlineData("tenant-a", "tenant-b", false)]
[InlineData(Tenant.AgnosticTenantId, Tenant.AgnosticTenantId, true)]
[InlineData(Tenant.AgnosticTenantId, "tenant-a", false)]
[InlineData(Tenant.AgnosticTenantId, Tenant.DefaultTenantId, false)]
[InlineData(null, Tenant.DefaultTenantId, true)]
[InlineData(null, "tenant-a", false)]
[InlineData(Tenant.DefaultTenantId, Tenant.DefaultTenantId, true)]
public void CanReplace_MatchesMemoryAlterationOwnership(string? existingTenantId, string writerTenantId, bool expected)
{
Assert.Equal(expected, TenantVisibility.CanReplace(existingTenantId, writerTenantId));
}
[Theory]
[InlineData("tenant-a", "tenant-a", "tenant-a", true)]
[InlineData("tenant-a", "tenant-b", "tenant-b", false)]
[InlineData("tenant-a", "tenant-a", "tenant-b", false)]
[InlineData(Tenant.AgnosticTenantId, Tenant.AgnosticTenantId, Tenant.AgnosticTenantId, true)]
[InlineData(Tenant.AgnosticTenantId, Tenant.AgnosticTenantId, "tenant-a", false)]
[InlineData(Tenant.AgnosticTenantId, "tenant-a", "tenant-a", false)]
[InlineData(Tenant.AgnosticTenantId, "tenant-a", Tenant.AgnosticTenantId, false)]
[InlineData(null, Tenant.DefaultTenantId, Tenant.DefaultTenantId, true)]
[InlineData(null, "tenant-a", "tenant-a", false)]
[InlineData(null, Tenant.DefaultTenantId, "tenant-a", false)]
public void CanReplaceOwnedRow_GatesNamedRowsOnAmbientNotForgedSource(
string? existingTenantId,
string? sourceTenantId,
string ambientTenantId,
bool expected)
{
Assert.Equal(expected, TenantVisibility.CanReplaceOwnedRow(existingTenantId, sourceTenantId, ambientTenantId));
}
[Fact]
public void WhereVisibleToTenant_WhenAmbientIsDefault_IncludesNullTenantId()
{

View file

@ -0,0 +1,37 @@
namespace Elsa.Persistence.EFCore.UnitTests;
public sealed class DbExceptionClassifierTests
{
[Theory]
[InlineData(19, 1555, true)]
[InlineData(19, 2067, true)]
[InlineData(19, 19, false)]
[InlineData(19, 0, false)]
public void IsDuplicateKey_WhenSqliteExtendedCodeExists_UsesOnlyExtendedDuplicateCodes(int baseCode, int extendedCode, bool expected)
{
var exception = new SqliteExceptionWithExtendedCode(baseCode, extendedCode);
Assert.Equal(expected, DbExceptionClassifier.IsDuplicateKey(exception));
}
[Theory]
[InlineData(19, true)]
[InlineData(1, false)]
public void IsDuplicateKey_WhenSqliteExtendedCodeIsUnavailable_FallsBackToBaseCode(int baseCode, bool expected)
{
var exception = new SqliteExceptionWithoutExtendedCode(baseCode);
Assert.Equal(expected, DbExceptionClassifier.IsDuplicateKey(exception));
}
private sealed class SqliteExceptionWithExtendedCode(int sqliteErrorCode, int sqliteExtendedErrorCode) : Exception
{
public int SqliteErrorCode { get; } = sqliteErrorCode;
public int SqliteExtendedErrorCode { get; } = sqliteExtendedErrorCode;
}
private sealed class SqliteExceptionWithoutExtendedCode(int sqliteErrorCode) : Exception
{
public int SqliteErrorCode { get; } = sqliteErrorCode;
}
}