diff --git a/src/modules/Elsa.Alterations.Core/Stores/AlterationStoreConflict.cs b/src/modules/Elsa.Alterations.Core/Stores/AlterationStoreConflict.cs new file mode 100644 index 000000000..213fe884b --- /dev/null +++ b/src/modules/Elsa.Alterations.Core/Stores/AlterationStoreConflict.cs @@ -0,0 +1,13 @@ +namespace Elsa.Alterations.Core.Stores; + +/// +/// Shared refuse-to-overwrite exceptions for Memory and EF alteration stores. +/// +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."); +} diff --git a/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationJobStore.cs b/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationJobStore.cs index e4d541e76..7f5032616 100644 --- a/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationJobStore.cs +++ b/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationJobStore.cs @@ -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); } /// /// * is visible to every tenant, but only an agnostic writer may replace it. /// Named tenants may upsert their own visible rows. /// - 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) { diff --git a/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationPlanStore.cs b/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationPlanStore.cs index bfde620ea..77646945f 100644 --- a/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationPlanStore.cs +++ b/src/modules/Elsa.Alterations.Core/Stores/MemoryAlterationPlanStore.cs @@ -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); } /// /// * is visible to every tenant, but only an agnostic writer may replace it. /// Named tenants may upsert their own visible rows. /// - 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) { diff --git a/src/modules/Elsa.Common/Multitenancy/TenantVisibility.cs b/src/modules/Elsa.Common/Multitenancy/TenantVisibility.cs index a6b602367..0377ac4a7 100644 --- a/src/modules/Elsa.Common/Multitenancy/TenantVisibility.cs +++ b/src/modules/Elsa.Common/Multitenancy/TenantVisibility.cs @@ -17,6 +17,28 @@ public static class TenantVisibility || entityTenantId == Tenant.AgnosticTenantId || entityTenantId is null && ambientTenantId == Tenant.DefaultTenantId; + /// + /// Write counterpart of . * 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 ). + /// + public static bool CanReplace(string? existingTenantId, string writerTenantId) => + existingTenantId == writerTenantId + || existingTenantId is null && writerTenantId == Tenant.DefaultTenantId; + + /// + /// CAS match for an upsert. Named rows are owned by the + /// (Memory ), not by a forged source TenantId. A * row + /// is replaceable only when both the incoming entity and the ambient writer are *. + /// + 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); + } + /// /// Restricts to rows visible to , /// unless is set (EF IgnoreQueryFilters). diff --git a/src/modules/Elsa.Persistence.EFCore.Common/DbExceptionClassifier.cs b/src/modules/Elsa.Persistence.EFCore.Common/DbExceptionClassifier.cs index 8f8e1e677..0e062f61d 100644 --- a/src/modules/Elsa.Persistence.EFCore.Common/DbExceptionClassifier.cs +++ b/src/modules/Elsa.Persistence.EFCore.Common/DbExceptionClassifier.cs @@ -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; diff --git a/src/modules/Elsa.Persistence.EFCore.Common/Store.cs b/src/modules/Elsa.Persistence.EFCore.Common/Store.cs index 900bb0c95..9ee59a48e 100644 --- a/src/modules/Elsa.Persistence.EFCore.Common/Store.cs +++ b/src/modules/Elsa.Persistence.EFCore.Common/Store.cs @@ -19,8 +19,8 @@ namespace Elsa.Persistence.EFCore; [PublicAPI] public class Store(IDbContextFactory 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(IDbContextFactory 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(IDbContextFactory dbContextF var tenantId = serviceProvider.GetRequiredService().TenantId; - await ExecuteBulkWriteWithSqlServerRetryAsync(async (dbContext, ct) => + await ExecuteSqlServerWriteWithRetryAsync(async (dbContext, ct) => { if (onSaving != null) { @@ -238,7 +238,32 @@ public class Store(IDbContextFactory dbContextF await handler.HandleAsync(context); } - private async Task ExecuteBulkWriteWithSqlServerRetryAsync( + /// + /// Executes a database operation and passes failures through the configured database + /// exception handler. + /// + /// The operation result type. + /// The database operation to execute. + /// The cancellation token. + /// A predicate that excludes exceptions which are not database failures. + /// The result returned by . + internal async Task ExecuteWithDbExceptionHandlingAsync( + Func> operation, + CancellationToken cancellationToken = default, + Func? 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 operation, CancellationToken cancellationToken) { @@ -255,9 +280,9 @@ public class Store(IDbContextFactory 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(IDbContextFactory 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)); /// /// Updates the entity. diff --git a/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationJobStore.cs b/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationJobStore.cs index eea874bee..4db2753ed 100644 --- a/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationJobStore.cs +++ b/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationJobStore.cs @@ -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; /// -/// An EF Core implementation of . +/// An EF Core implementation of . /// public class EFCoreAlterationJobStore : IAlterationJobStore { private readonly EntityStore _store; + private readonly bool _tenantEnabled; /// /// Constructor. /// public EFCoreAlterationJobStore(EntityStore store) + : this(store, Options.Create(new TenantsOptions())) + { + } + + /// + /// Constructor used by dependency injection. Direct construction through the legacy + /// overload keeps tenancy-aware upsert disabled for compatibility. + /// + [ActivatorUtilitiesConstructor] + public EFCoreAlterationJobStore(EntityStore store, IOptions tenantsOptions) { _store = store; + _tenantEnabled = tenantsOptions.Value.IsEnabled; } /// 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); } /// public async Task SaveManyAsync(IEnumerable 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); } /// @@ -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("SerializedLog").CurrentValue; + + var query = dbContext.Set() + .IgnoreQueryFilters() + .Where(AlterationTenantOwnedUpsert.OwnedId(record.Id, record.TenantId, ambientTenantId)); + + // Inline lambda so net8/net9 bind SetPropertyCalls and net10 binds UpdateSettersBuilder. + Task 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(job, "SerializedLog"), serializedLog), + cancellationToken); + + Task ExecuteWriteAsync(Func> 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 Filter(IQueryable queryable, AlterationJobFilter filter) => filter.Apply(queryable); -} \ No newline at end of file + + private static bool IsDatabaseException(Exception exception) => + exception is DbException or DbUpdateException || exception.InnerException is DbException; +} diff --git a/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationPlanStore.cs b/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationPlanStore.cs index 83ed2e59c..fb5effe9c 100644 --- a/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationPlanStore.cs +++ b/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationPlanStore.cs @@ -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 _store; private readonly IAlterationSerializer _alterationSerializer; + private readonly bool _tenantEnabled; /// /// Constructor. /// - public EFCoreAlterationPlanStore(EntityStore store, IAlterationSerializer alterationSerializer) + public EFCoreAlterationPlanStore( + EntityStore store, + IAlterationSerializer alterationSerializer) + : this(store, alterationSerializer, Options.Create(new TenantsOptions())) + { + } + + /// + /// Constructor used by dependency injection. Direct construction through the legacy + /// overload keeps tenancy-aware upsert disabled for compatibility. + /// + [ActivatorUtilitiesConstructor] + public EFCoreAlterationPlanStore( + EntityStore store, + IAlterationSerializer alterationSerializer, + IOptions tenantsOptions) { _store = store; _alterationSerializer = alterationSerializer; + _tenantEnabled = tenantsOptions.Value.IsEnabled; } /// 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("SerializedAlterations").CurrentValue; + var serializedFilter = dbContext.Entry(record).Property("SerializedWorkflowInstanceFilter").CurrentValue; + + var query = dbContext.Set() + .IgnoreQueryFilters() + .Where(AlterationTenantOwnedUpsert.OwnedId(record.Id, record.TenantId, ambientTenantId)); + + // Inline lambda so net8/net9 bind SetPropertyCalls and net10 binds UpdateSettersBuilder. + Task 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(plan, "SerializedAlterations"), serializedAlterations) + .SetProperty(plan => EF.Property(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); + } + } } /// @@ -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; } -} \ No newline at end of file +} diff --git a/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationTenantOwnedUpsert.cs b/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationTenantOwnedUpsert.cs new file mode 100644 index 000000000..302885866 --- /dev/null +++ b/src/modules/Elsa.Persistence.EFCore/Modules/Alterations/AlterationTenantOwnedUpsert.cs @@ -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; + +/// +/// 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. +/// +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; + + /// + /// EF translation of . Named UPDATE + /// matches Target.TenantId to the ambient writer, not a forged source + /// TenantId. * updates still require ambient and source *. + /// + public static Expression> OwnedId(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 InsertIfAbsentAsync( + 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; + } +} diff --git a/test/integration/Elsa.Alterations.Persistence.ConformanceTests/AlterationStoreScenario.cs b/test/integration/Elsa.Alterations.Persistence.ConformanceTests/AlterationStoreScenario.cs index 2ba85e8cf..900d9af31 100644 --- a/test/integration/Elsa.Alterations.Persistence.ConformanceTests/AlterationStoreScenario.cs +++ b/test/integration/Elsa.Alterations.Persistence.ConformanceTests/AlterationStoreScenario.cs @@ -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 CreateSqliteAsync() + public static Task CreateSqliteAsync() => + CreateSqliteAsync("tenant-a"); + + public static Task CreateSqliteAsync(string tenantId) => + CreateSqliteAsync(tenantId, Path.Join(Path.GetTempPath(), $"elsa-alterations-conformance-{Guid.NewGuid():N}.db"), ownsDatabaseFile: true); + + public static Task 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); + + /// + /// Two EF/SQLite hosts that share a file and keep separate ambient tenants so concurrent + /// Save/SaveMany calls do not mutate a single . + /// + public static async Task 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 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 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(tenantAccessor) .AddSingleton() - .Configure(options => options.IsEnabled = true) + .Configure(options => options.IsEnabled = tenantsEnabled) .AddScoped() .AddScoped() .AddSqliteEntityModelCreatingHandlers() .AddDbContextFactory((_, 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, TenantAwareDbContextFactory>() .AddScoped>() .AddScoped>() .AddScoped() - .AddScoped() - .BuildServiceProvider(); + .AddScoped(); + + if (dbExceptionHandler is not null) + serviceCollection.AddSingleton(dbExceptionHandler); + + services = serviceCollection.BuildServiceProvider(); await using (var dbContext = await services.GetRequiredService>().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( } } +/// +/// Two SQLite alteration-store hosts that share one database file. +/// +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); + } +} + /// /// Minimal used to lock save/load of serialized plan alterations. /// diff --git a/test/integration/Elsa.Alterations.Persistence.ConformanceTests/EFCoreAlterationStoreTenantOwnershipTests.cs b/test/integration/Elsa.Alterations.Persistence.ConformanceTests/EFCoreAlterationStoreTenantOwnershipTests.cs new file mode 100644 index 000000000..a821261ff --- /dev/null +++ b/test/integration/Elsa.Alterations.Persistence.ConformanceTests/EFCoreAlterationStoreTenantOwnershipTests.cs @@ -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; + +/// +/// 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. +/// +[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(() => + 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(() => + 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(() => + 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(() => + 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(() => + 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(() => + 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(() => + 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(() => + 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(() => 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(() => 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(() => 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(() => 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(() => 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(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(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(ex)); + + var remaining = await owner.Plans.FindAsync(new AlterationPlanFilter { Id = "shared" }); + AssertUnchangedPlan(remaining, "tenant-a", AlterationPlanStatus.Running, "original"); + } + + private static async Task Capture(Func 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(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); +} + +/// +/// 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. +/// +public sealed class GateFirstAlterationUpdates(string tableName) : DbCommandInterceptor +{ + private readonly TaskCompletionSource _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 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 _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> TransactionStartingAsync( + DbConnection connection, + TransactionStartingEventData eventData, + InterceptionResult 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> NonQueryExecutingAsync( + DbCommand command, + CommandEventData eventData, + InterceptionResult result, + CancellationToken cancellationToken = default) + { + ThrowIfMatched(command); + + return ValueTask.FromResult(result); + } + + public override ValueTask> ReaderExecutingAsync( + DbCommand command, + CommandEventData eventData, + InterceptionResult 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; + } +} diff --git a/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationJobStoreTenantIsolationTests.cs b/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationJobStoreTenantIsolationTests.cs index ab152522f..6bbde4814 100644 --- a/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationJobStoreTenantIsolationTests.cs +++ b/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationJobStoreTenantIsolationTests.cs @@ -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(); + 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(() => 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(); + 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(() => 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() { diff --git a/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationPlanStoreTenantIsolationTests.cs b/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationPlanStoreTenantIsolationTests.cs index 2efc30483..8b275286d 100644 --- a/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationPlanStoreTenantIsolationTests.cs +++ b/test/unit/Elsa.Alterations.Core.UnitTests/Stores/MemoryAlterationPlanStoreTenantIsolationTests.cs @@ -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(); + 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(() => 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() { diff --git a/test/unit/Elsa.Common.UnitTests/Multitenancy/TenantVisibilityTests.cs b/test/unit/Elsa.Common.UnitTests/Multitenancy/TenantVisibilityTests.cs index 07ff5b958..9b9124bef 100644 --- a/test/unit/Elsa.Common.UnitTests/Multitenancy/TenantVisibilityTests.cs +++ b/test/unit/Elsa.Common.UnitTests/Multitenancy/TenantVisibilityTests.cs @@ -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() { diff --git a/test/unit/Elsa.Persistence.EFCore.UnitTests/DbExceptionClassifierTests.cs b/test/unit/Elsa.Persistence.EFCore.UnitTests/DbExceptionClassifierTests.cs new file mode 100644 index 000000000..8f7d6e125 --- /dev/null +++ b/test/unit/Elsa.Persistence.EFCore.UnitTests/DbExceptionClassifierTests.cs @@ -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; + } +}