diff --git a/framework/src/Volo.Abp.Data/Volo/Abp/Data/AppliedDatabaseMigrationsEto.cs b/framework/src/Volo.Abp.Data/Volo/Abp/Data/AppliedDatabaseMigrationsEto.cs new file mode 100644 index 0000000000..739571e241 --- /dev/null +++ b/framework/src/Volo.Abp.Data/Volo/Abp/Data/AppliedDatabaseMigrationsEto.cs @@ -0,0 +1,12 @@ +using System; +using Volo.Abp.EventBus; + +namespace Volo.Abp.Data; + +[Serializable] +[EventName("abp.data.applied_database_migrations")] +public class AppliedDatabaseMigrationsEto +{ + public string DatabaseName { get; set; } + public Guid? TenantId { get; set; } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/Migrations/DatabaseMigrationEventHandlerBase.cs b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/Migrations/DatabaseMigrationEventHandlerBase.cs new file mode 100644 index 0000000000..ce7e21dc0a --- /dev/null +++ b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/Migrations/DatabaseMigrationEventHandlerBase.cs @@ -0,0 +1,318 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Volo.Abp.Data; +using Volo.Abp.DependencyInjection; +using Volo.Abp.Domain.Entities.Events.Distributed; +using Volo.Abp.EventBus.Distributed; +using Volo.Abp.MultiTenancy; +using Volo.Abp.Uow; + +namespace Volo.Abp.EntityFrameworkCore.Migrations; + +public abstract class DatabaseMigrationEventHandlerBase : + IDistributedEventHandler, + IDistributedEventHandler, + IDistributedEventHandler, + ITransientDependency + where TDbContext : DbContext, IEfCoreDbContext +{ + protected string DatabaseName { get; } + + protected const string TryCountPropertyName = "__TryCount"; + + protected int MaxEventTryCount { get; set; } = 3; + + /// + /// As milliseconds. + /// + protected int MinValueToWaitOnFailure { get; set; } = 5000; + + /// + /// As milliseconds. + /// + protected int MaxValueToWaitOnFailure { get; set; } = 15000; + + protected ICurrentTenant CurrentTenant { get; } + protected IUnitOfWorkManager UnitOfWorkManager { get; } + protected ITenantStore TenantStore { get; } + protected IDistributedEventBus DistributedEventBus { get; } + protected ILogger> Logger { get; } + + protected DatabaseMigrationEventHandlerBase( + string databaseName, + ICurrentTenant currentTenant, + IUnitOfWorkManager unitOfWorkManager, + ITenantStore tenantStore, + IDistributedEventBus distributedEventBus, + ILoggerFactory loggerFactory) + { + CurrentTenant = currentTenant; + UnitOfWorkManager = unitOfWorkManager; + TenantStore = tenantStore; + DatabaseName = databaseName; + DistributedEventBus = distributedEventBus; + + Logger = loggerFactory.CreateLogger>(); + } + + public virtual async Task HandleEventAsync(ApplyDatabaseMigrationsEto eventData) + { + if (eventData.DatabaseName != DatabaseName) + { + return; + } + + var schemaMigrated = false; + try + { + schemaMigrated = await MigrateDatabaseSchemaAsync(eventData.TenantId); + await SeedAsync(eventData.TenantId); + + if (schemaMigrated) + { + await DistributedEventBus.PublishAsync( + new AppliedDatabaseMigrationsEto + { + DatabaseName = DatabaseName, + TenantId = eventData.TenantId + } + ); + } + } + catch (Exception ex) + { + await HandleErrorOnApplyDatabaseMigrationAsync(eventData, ex); + } + + await AfterApplyDatabaseMigrations(eventData, schemaMigrated); + } + + protected virtual Task AfterApplyDatabaseMigrations(ApplyDatabaseMigrationsEto eventData, bool schemaMigrated) + { + return Task.CompletedTask; + } + + public virtual async Task HandleEventAsync(TenantCreatedEto eventData) + { + var schemaMigrated = false; + try + { + schemaMigrated = await MigrateDatabaseSchemaAsync(eventData.Id); + await SeedAsync(eventData.Id); + + if (schemaMigrated) + { + await DistributedEventBus.PublishAsync( + new AppliedDatabaseMigrationsEto + { + DatabaseName = DatabaseName, + TenantId = eventData.Id + } + ); + } + } + catch (Exception ex) + { + await HandleErrorTenantCreatedAsync(eventData, ex); + } + + await AfterTenantCreated(eventData, schemaMigrated); + } + + protected virtual Task AfterTenantCreated(TenantCreatedEto eventData, bool schemaMigrated) + { + return Task.CompletedTask; + } + + public virtual async Task HandleEventAsync(TenantConnectionStringUpdatedEto eventData) + { + if (eventData.ConnectionStringName != DatabaseName && + eventData.ConnectionStringName != Volo.Abp.Data.ConnectionStrings.DefaultConnectionStringName || + eventData.NewValue.IsNullOrWhiteSpace()) + { + return; + } + + var schemaMigrated = false; + try + { + schemaMigrated = await MigrateDatabaseSchemaAsync(eventData.Id); + await SeedAsync(eventData.Id); + + if (schemaMigrated) + { + await DistributedEventBus.PublishAsync( + new AppliedDatabaseMigrationsEto + { + DatabaseName = DatabaseName, + TenantId = eventData.Id + } + ); + } + } + catch (Exception ex) + { + await HandleErrorTenantConnectionStringUpdatedAsync(eventData, ex); + } + + await AfterTenantConnectionStringUpdated(eventData, schemaMigrated); + } + + protected virtual Task AfterTenantConnectionStringUpdated(TenantConnectionStringUpdatedEto eventData, + bool schemaMigrated) + { + return Task.CompletedTask; + } + + protected virtual Task SeedAsync(Guid? tenantId) + { + return Task.CompletedTask; + } + + /// + /// Apply pending EF Core schema migrations to the database. + /// Returns true if any migration has applied. + /// + protected virtual async Task MigrateDatabaseSchemaAsync(Guid? tenantId) + { + var result = false; + + using (CurrentTenant.Change(tenantId)) + { + using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) + { + async Task MigrateDatabaseSchemaWithDbContextAsync() + { + var dbContext = await uow.ServiceProvider + .GetRequiredService>() + .GetDbContextAsync(); + + if ((await dbContext.Database.GetPendingMigrationsAsync()).Any()) + { + await dbContext.Database.MigrateAsync(); + return true; + } + + return false; + } + + if (tenantId == null) + { + //Migrating the host database + Logger.LogInformation($"Migrating database of host. Database Name = {DatabaseName}"); + result = await MigrateDatabaseSchemaWithDbContextAsync(); + } + else + { + var tenantConfiguration = await TenantStore.FindAsync(tenantId.Value); + if (!tenantConfiguration.ConnectionStrings.Default.IsNullOrWhiteSpace() || + !tenantConfiguration.ConnectionStrings.GetOrDefault(DatabaseName).IsNullOrWhiteSpace()) + { + //Migrating the tenant database (only if tenant has a separate database) + Logger.LogInformation( + $"Migrating separate database of tenant. Database Name = {DatabaseName}, TenantId = {tenantId}"); + result = await MigrateDatabaseSchemaWithDbContextAsync(); + } + } + + await uow.CompleteAsync(); + } + } + + return result; + } + + protected virtual async Task HandleErrorOnApplyDatabaseMigrationAsync( + ApplyDatabaseMigrationsEto eventData, + Exception exception) + { + var tryCount = IncrementEventTryCount(eventData); + if (tryCount <= MaxEventTryCount) + { + Logger.LogWarning( + $"Could not apply database migrations. Re-queueing the operation. TenantId = {eventData.TenantId}, Database Name = {eventData.DatabaseName}."); + Logger.LogException(exception, LogLevel.Warning); + + await Task.Delay(RandomHelper.GetRandom(MinValueToWaitOnFailure, MaxValueToWaitOnFailure)); + await DistributedEventBus.PublishAsync(eventData); + } + else + { + Logger.LogError( + $"Could not apply database migrations. Canceling the operation. TenantId = {eventData.TenantId}, DatabaseName = {eventData.DatabaseName}."); + Logger.LogException(exception); + } + } + + protected virtual async Task HandleErrorTenantCreatedAsync( + TenantCreatedEto eventData, + Exception exception) + { + var tryCount = IncrementEventTryCount(eventData); + if (tryCount <= MaxEventTryCount) + { + Logger.LogWarning( + $"Could not perform tenant created event. Re-queueing the operation. TenantId = {eventData.Id}, TenantName = {eventData.Name}."); + Logger.LogException(exception, LogLevel.Warning); + + await Task.Delay(RandomHelper.GetRandom(5000, 15000)); + await DistributedEventBus.PublishAsync(eventData); + } + else + { + Logger.LogError( + $"Could not perform tenant created event. Canceling the operation. TenantId = {eventData.Id}, TenantName = {eventData.Name}."); + Logger.LogException(exception); + } + } + + protected virtual async Task HandleErrorTenantConnectionStringUpdatedAsync( + TenantConnectionStringUpdatedEto eventData, + Exception exception) + { + var tryCount = IncrementEventTryCount(eventData); + if (tryCount <= MaxEventTryCount) + { + Logger.LogWarning( + $"Could not perform tenant connection string updated event. Re-queueing the operation. TenantId = {eventData.Id}, TenantName = {eventData.Name}."); + Logger.LogException(exception, LogLevel.Warning); + + await Task.Delay(RandomHelper.GetRandom(5000, 15000)); + await DistributedEventBus.PublishAsync(eventData); + } + else + { + Logger.LogError( + $"Could not perform tenant connection string updated event. Canceling the operation. TenantId = {eventData.Id}, TenantName = {eventData.Name}."); + Logger.LogException(exception); + } + } + + private static int GetEventTryCount(EtoBase eventData) + { + var tryCountAsString = eventData.Properties.GetOrDefault(TryCountPropertyName); + if (tryCountAsString.IsNullOrEmpty()) + { + return 0; + } + + return int.Parse(tryCountAsString); + } + + private static void SetEventTryCount(EtoBase eventData, int count) + { + eventData.Properties[TryCountPropertyName] = count.ToString(); + } + + private static int IncrementEventTryCount(EtoBase eventData) + { + var count = GetEventTryCount(eventData) + 1; + SetEventTryCount(eventData, count); + return count; + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/Migrations/EfCoreRuntimeDatabaseMigratorBase.cs b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/Migrations/EfCoreRuntimeDatabaseMigratorBase.cs new file mode 100644 index 0000000000..9b9719ff1b --- /dev/null +++ b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/Migrations/EfCoreRuntimeDatabaseMigratorBase.cs @@ -0,0 +1,145 @@ +using System; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Volo.Abp.Data; +using Volo.Abp.DependencyInjection; +using Volo.Abp.DistributedLocking; +using Volo.Abp.EventBus.Distributed; +using Volo.Abp.MultiTenancy; +using Volo.Abp.Uow; + +namespace Volo.Abp.EntityFrameworkCore.Migrations; + +public abstract class EfCoreRuntimeDatabaseMigratorBase : ITransientDependency + where TDbContext : DbContext, IEfCoreDbContext +{ + protected int MinValueToWaitOnFailure { get; set; } = 5000; + protected int MaxValueToWaitOnFailure { get; set; } = 15000; + + protected string DatabaseName { get; } + + /// + /// Enabling this might be inefficient if you have many tenants! + /// If disabled (default), tenant databases will be seeded only + /// if there is a schema migration applied to the host database. + /// If enabled, tenant databases will be seeded always on every service startup. + /// + protected bool AlwaysSeedTenantDatabases { get; set; } = false; + + protected IUnitOfWorkManager UnitOfWorkManager { get; } + protected IServiceProvider ServiceProvider { get; } + protected ICurrentTenant CurrentTenant { get; } + protected IAbpDistributedLock DistributedLock { get; } + protected IDistributedEventBus DistributedEventBus { get; } + protected ILogger> Logger { get; } + + protected EfCoreRuntimeDatabaseMigratorBase( + string databaseName, + IUnitOfWorkManager unitOfWorkManager, + IServiceProvider serviceProvider, + ICurrentTenant currentTenant, + IAbpDistributedLock abpDistributedLock, + IDistributedEventBus distributedEventBus, + ILoggerFactory loggerFactory) + { + DatabaseName = databaseName; + UnitOfWorkManager = unitOfWorkManager; + ServiceProvider = serviceProvider; + CurrentTenant = currentTenant; + DistributedLock = abpDistributedLock; + DistributedEventBus = distributedEventBus; + Logger = loggerFactory.CreateLogger>(); + } + + public virtual async Task CheckAndApplyDatabaseMigrationsAsync() + { + await TryAsync(LockAndApplyDatabaseMigrationsAsync); + } + + protected virtual async Task LockAndApplyDatabaseMigrationsAsync() + { + Logger.LogInformation($"Trying to acquire the distributed lock for database migration: {DatabaseName}."); + + var schemaMigrated = false; + + await using (var handle = await DistributedLock.TryAcquireAsync("DatabaseMigration_" + DatabaseName)) + { + if (handle is null) + { + Logger.LogInformation($"Distributed lock could not be acquired for database migration: {DatabaseName}. Operation cancelled."); + return; + } + + Logger.LogInformation($"Distributed lock is acquired for database migration: {DatabaseName}..."); + + using (CurrentTenant.Change(null)) + { + // Create database tables if needed + using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) + { + var dbContext = await ServiceProvider + .GetRequiredService>() + .GetDbContextAsync(); + + var pendingMigrations = await dbContext + .Database + .GetPendingMigrationsAsync(); + + if (pendingMigrations.Any()) + { + await dbContext.Database.MigrateAsync(); + schemaMigrated = true; + } + + await uow.CompleteAsync(); + } + } + + await SeedAsync(); + + if (schemaMigrated || AlwaysSeedTenantDatabases) + { + await DistributedEventBus.PublishAsync( + new AppliedDatabaseMigrationsEto + { + DatabaseName = DatabaseName, + TenantId = null + } + ); + } + } + + Logger.LogInformation($"Distributed lock has been released for database migration: {DatabaseName}..."); + } + + protected virtual Task SeedAsync() + { + return Task.CompletedTask; + } + + protected virtual async Task TryAsync(Func task, int maxTryCount = 3) + { + try + { + await task(); + } + catch (Exception ex) + { + maxTryCount--; + + if (maxTryCount <= 0) + { + throw; + } + + Logger.LogWarning($"{ex.GetType().Name} has been thrown. The operation will be tried {maxTryCount} times more. Exception:\n{ex.Message}. Stack Trace:\n{ex.StackTrace}"); + + await Task.Delay(RandomHelper.GetRandom(MinValueToWaitOnFailure, MaxValueToWaitOnFailure)); + + await TryAsync(task, maxTryCount); + } + } +} \ No newline at end of file