7 changed files with 519 additions and 246 deletions
@ -1,185 +1,7 @@ |
|||||
using System; |
using Volo.Abp.DependencyInjection; |
||||
using System.Collections.Generic; |
|
||||
using System.Linq; |
|
||||
using System.Threading.Tasks; |
|
||||
using Microsoft.EntityFrameworkCore; |
|
||||
using Microsoft.Extensions.DependencyInjection; |
|
||||
using Microsoft.Extensions.Logging; |
|
||||
using Microsoft.Extensions.Logging.Abstractions; |
|
||||
using Serilog; |
|
||||
using Volo.Abp; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DependencyInjection; |
|
||||
using Volo.Abp.Domain.Entities.Events.Distributed; |
|
||||
using Volo.Abp.EntityFrameworkCore; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations |
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations; |
||||
|
|
||||
|
public abstract class DatabaseMigrationEventHandlerBase : ITransientDependency |
||||
{ |
{ |
||||
public abstract class DatabaseMigrationEventHandlerBase<TDbContext> : |
} |
||||
ITransientDependency |
|
||||
where TDbContext : DbContext, IEfCoreDbContext |
|
||||
{ |
|
||||
protected const string TryCountPropertyName = "TryCount"; |
|
||||
protected const int MaxEventTryCount = 3; |
|
||||
|
|
||||
protected ICurrentTenant CurrentTenant { get; } |
|
||||
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
|
||||
protected ITenantStore TenantStore { get; } |
|
||||
protected IDistributedEventBus DistributedEventBus { get; } |
|
||||
protected ILogger<DatabaseMigrationEventHandlerBase<TDbContext>> Logger { get; set; } |
|
||||
protected string DatabaseName { get; } |
|
||||
|
|
||||
protected DatabaseMigrationEventHandlerBase( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
string databaseName) |
|
||||
{ |
|
||||
CurrentTenant = currentTenant; |
|
||||
UnitOfWorkManager = unitOfWorkManager; |
|
||||
TenantStore = tenantStore; |
|
||||
DatabaseName = databaseName; |
|
||||
DistributedEventBus = distributedEventBus; |
|
||||
|
|
||||
Logger = NullLogger<DatabaseMigrationEventHandlerBase<TDbContext>>.Instance; |
|
||||
} |
|
||||
|
|
||||
/// <summary>
|
|
||||
/// Apply pending EF Core schema migrations to the database.
|
|
||||
/// Returns true if any migration has applied.
|
|
||||
/// </summary>
|
|
||||
protected virtual async Task<bool> MigrateDatabaseSchemaAsync(Guid? tenantId) |
|
||||
{ |
|
||||
var result = false; |
|
||||
|
|
||||
using (CurrentTenant.Change(tenantId)) |
|
||||
{ |
|
||||
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) |
|
||||
{ |
|
||||
async Task<bool> MigrateDatabaseSchemaWithDbContextAsync() |
|
||||
{ |
|
||||
var dbContext = await uow.ServiceProvider |
|
||||
.GetRequiredService<IDbContextProvider<TDbContext>>() |
|
||||
.GetDbContextAsync(); |
|
||||
|
|
||||
if ((await dbContext.Database.GetPendingMigrationsAsync()).Any()) |
|
||||
{ |
|
||||
await dbContext.Database.MigrateAsync(); |
|
||||
return true; |
|
||||
} |
|
||||
|
|
||||
return false; |
|
||||
} |
|
||||
|
|
||||
if (tenantId == null) |
|
||||
{ |
|
||||
//Migrating the host database
|
|
||||
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)
|
|
||||
result = await MigrateDatabaseSchemaWithDbContextAsync(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
await uow.CompleteAsync(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
return result; |
|
||||
} |
|
||||
|
|
||||
protected virtual async Task HandleErrorOnApplyDatabaseMigrationAsync( |
|
||||
ApplyDatabaseMigrationsEto eventData, |
|
||||
Exception exception) |
|
||||
{ |
|
||||
var tryCount = IncrementEventTryCount(eventData); |
|
||||
if (tryCount <= MaxEventTryCount) |
|
||||
{ |
|
||||
Log.Warning($"Could not apply database migrations. Re-queueing the operation. TenantId = {eventData.TenantId}, Database Name = {eventData.DatabaseName}."); |
|
||||
Log.Error(exception.ToString()); |
|
||||
|
|
||||
await Task.Delay(RandomHelper.GetRandom(5000, 15000)); |
|
||||
Log.Warning("Re publishing the event!"); |
|
||||
await DistributedEventBus.PublishAsync(eventData); |
|
||||
} |
|
||||
else |
|
||||
{ |
|
||||
Log.Error($"Could not apply database migrations. Canceling the operation. TenantId = {eventData.TenantId}, DatabaseName = {eventData.DatabaseName}."); |
|
||||
Log.Error(exception.ToString()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
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; |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -0,0 +1,188 @@ |
|||||
|
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 Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Serilog; |
||||
|
using Volo.Abp; |
||||
|
using Volo.Abp.Data; |
||||
|
using Volo.Abp.Domain.Entities.Events.Distributed; |
||||
|
using Volo.Abp.EntityFrameworkCore; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
using Volo.Abp.MultiTenancy; |
||||
|
using Volo.Abp.Uow; |
||||
|
|
||||
|
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.EfCore; |
||||
|
|
||||
|
public abstract class DatabaseEfCoreMigrationEventHandler<TDbContext> : DatabaseMigrationEventHandlerBase |
||||
|
where TDbContext : DbContext, IEfCoreDbContext |
||||
|
{ |
||||
|
protected const string TryCountPropertyName = "TryCount"; |
||||
|
protected const int MaxEventTryCount = 3; |
||||
|
|
||||
|
protected ICurrentTenant CurrentTenant { get; } |
||||
|
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
||||
|
protected ITenantStore TenantStore { get; } |
||||
|
protected IDistributedEventBus DistributedEventBus { get; } |
||||
|
protected ILogger<DatabaseEfCoreMigrationEventHandler<TDbContext>> Logger { get; set; } |
||||
|
protected string DatabaseName { get; } |
||||
|
|
||||
|
protected DatabaseEfCoreMigrationEventHandler( |
||||
|
ICurrentTenant currentTenant, |
||||
|
IUnitOfWorkManager unitOfWorkManager, |
||||
|
ITenantStore tenantStore, |
||||
|
IDistributedEventBus distributedEventBus, |
||||
|
string databaseName) |
||||
|
{ |
||||
|
CurrentTenant = currentTenant; |
||||
|
UnitOfWorkManager = unitOfWorkManager; |
||||
|
TenantStore = tenantStore; |
||||
|
DatabaseName = databaseName; |
||||
|
DistributedEventBus = distributedEventBus; |
||||
|
|
||||
|
Logger = NullLogger<DatabaseEfCoreMigrationEventHandler<TDbContext>>.Instance; |
||||
|
} |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Apply pending EF Core schema migrations to the database.
|
||||
|
/// Returns true if any migration has applied.
|
||||
|
/// </summary>
|
||||
|
protected virtual async Task<bool> MigrateDatabaseSchemaAsync(Guid? tenantId) |
||||
|
{ |
||||
|
var result = false; |
||||
|
|
||||
|
using (CurrentTenant.Change(tenantId)) |
||||
|
{ |
||||
|
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) |
||||
|
{ |
||||
|
async Task<bool> MigrateDatabaseSchemaWithDbContextAsync() |
||||
|
{ |
||||
|
var dbContext = await uow.ServiceProvider |
||||
|
.GetRequiredService<IDbContextProvider<TDbContext>>() |
||||
|
.GetDbContextAsync(); |
||||
|
|
||||
|
if ((await dbContext.Database.GetPendingMigrationsAsync()).Any()) |
||||
|
{ |
||||
|
await dbContext.Database.MigrateAsync(); |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
if (tenantId == null) |
||||
|
{ |
||||
|
//Migrating the host database
|
||||
|
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)
|
||||
|
result = await MigrateDatabaseSchemaWithDbContextAsync(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
await uow.CompleteAsync(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
return result; |
||||
|
} |
||||
|
|
||||
|
protected virtual async Task HandleErrorOnApplyDatabaseMigrationAsync( |
||||
|
ApplyDatabaseMigrationsEto eventData, |
||||
|
Exception exception) |
||||
|
{ |
||||
|
var tryCount = IncrementEventTryCount(eventData); |
||||
|
if (tryCount <= MaxEventTryCount) |
||||
|
{ |
||||
|
Log.Warning( |
||||
|
$"Could not apply database migrations. Re-queueing the operation. TenantId = {eventData.TenantId}, Database Name = {eventData.DatabaseName}."); |
||||
|
Log.Error(exception.ToString()); |
||||
|
|
||||
|
await Task.Delay(RandomHelper.GetRandom(5000, 15000)); |
||||
|
Log.Warning("Re publishing the event!"); |
||||
|
await DistributedEventBus.PublishAsync(eventData); |
||||
|
} |
||||
|
else |
||||
|
{ |
||||
|
Log.Error( |
||||
|
$"Could not apply database migrations. Canceling the operation. TenantId = {eventData.TenantId}, DatabaseName = {eventData.DatabaseName}."); |
||||
|
Log.Error(exception.ToString()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
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; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,67 @@ |
|||||
|
using System; |
||||
|
using System.Linq; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.EntityFrameworkCore; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Data; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
using Volo.Abp.MultiTenancy; |
||||
|
using Volo.Abp.Uow; |
||||
|
|
||||
|
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.EfCore; |
||||
|
|
||||
|
public abstract class PendingEfCoreMigrationsChecker<TDbContext> : PendingMigrationsCheckerBase |
||||
|
where TDbContext : DbContext |
||||
|
{ |
||||
|
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
protected ICurrentTenant CurrentTenant { get; } |
||||
|
protected IDistributedEventBus DistributedEventBus { get; } |
||||
|
protected string DatabaseName { get; } |
||||
|
|
||||
|
protected PendingEfCoreMigrationsChecker( |
||||
|
IUnitOfWorkManager unitOfWorkManager, |
||||
|
IServiceProvider serviceProvider, |
||||
|
ICurrentTenant currentTenant, |
||||
|
IDistributedEventBus distributedEventBus, |
||||
|
string databaseName) |
||||
|
{ |
||||
|
UnitOfWorkManager = unitOfWorkManager; |
||||
|
ServiceProvider = serviceProvider; |
||||
|
CurrentTenant = currentTenant; |
||||
|
DistributedEventBus = distributedEventBus; |
||||
|
DatabaseName = databaseName; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<bool> CheckAsync() |
||||
|
{ |
||||
|
var isMigrationRequired = false; |
||||
|
|
||||
|
using (CurrentTenant.Change(null)) |
||||
|
{ |
||||
|
// Create database tables if needed
|
||||
|
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) |
||||
|
{ |
||||
|
var pendingMigrations = await ServiceProvider |
||||
|
.GetRequiredService<TDbContext>() |
||||
|
.Database |
||||
|
.GetPendingMigrationsAsync(); |
||||
|
|
||||
|
if (pendingMigrations.Any()) |
||||
|
{ |
||||
|
await DistributedEventBus.PublishAsync( |
||||
|
new ApplyDatabaseMigrationsEto |
||||
|
{ |
||||
|
DatabaseName = DatabaseName |
||||
|
} |
||||
|
); |
||||
|
isMigrationRequired = true; |
||||
|
} |
||||
|
|
||||
|
await uow.CompleteAsync(); |
||||
|
} |
||||
|
|
||||
|
return isMigrationRequired; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,204 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Linq; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using MongoDB.Driver; |
||||
|
using Volo.Abp; |
||||
|
using Volo.Abp.Data; |
||||
|
using Volo.Abp.Domain.Entities.Events.Distributed; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
using Volo.Abp.MongoDB; |
||||
|
using Volo.Abp.MultiTenancy; |
||||
|
using Volo.Abp.Uow; |
||||
|
|
||||
|
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.MongoDb; |
||||
|
|
||||
|
public abstract class DatabaseMongoDbMigrationEventHandler<TDbContext> : DatabaseMigrationEventHandlerBase |
||||
|
where TDbContext : AbpMongoDbContext, IAbpMongoDbContext |
||||
|
{ |
||||
|
protected const string TryCountPropertyName = "TryCount"; |
||||
|
protected const int MaxEventTryCount = 3; |
||||
|
|
||||
|
protected ICurrentTenant CurrentTenant { get; } |
||||
|
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
||||
|
protected ITenantStore TenantStore { get; } |
||||
|
protected IDistributedEventBus DistributedEventBus { get; } |
||||
|
protected ILogger<DatabaseMongoDbMigrationEventHandler<TDbContext>> Logger { get; set; } |
||||
|
|
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
protected string DatabaseName { get; } |
||||
|
|
||||
|
protected DatabaseMongoDbMigrationEventHandler( |
||||
|
ICurrentTenant currentTenant, |
||||
|
IUnitOfWorkManager unitOfWorkManager, |
||||
|
ITenantStore tenantStore, |
||||
|
IDistributedEventBus distributedEventBus, |
||||
|
string databaseName, IServiceProvider serviceProvider) |
||||
|
{ |
||||
|
CurrentTenant = currentTenant; |
||||
|
UnitOfWorkManager = unitOfWorkManager; |
||||
|
TenantStore = tenantStore; |
||||
|
DatabaseName = databaseName; |
||||
|
ServiceProvider = serviceProvider; |
||||
|
DistributedEventBus = distributedEventBus; |
||||
|
|
||||
|
Logger = NullLogger<DatabaseMongoDbMigrationEventHandler<TDbContext>>.Instance; |
||||
|
} |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Apply pending EF Core schema migrations to the database.
|
||||
|
/// Returns true if any migration has applied.
|
||||
|
/// </summary>
|
||||
|
protected virtual async Task<bool> MigrateDatabaseSchemaAsync(Guid? tenantId) |
||||
|
{ |
||||
|
var result = false; |
||||
|
|
||||
|
using (CurrentTenant.Change(tenantId)) |
||||
|
{ |
||||
|
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) |
||||
|
{ |
||||
|
async Task<bool> MigrateDatabaseSchemaWithDbContextAsync() |
||||
|
{ |
||||
|
var dbContexts = ServiceProvider.GetServices<IAbpMongoDbContext>(); |
||||
|
var connectionStringResolver = ServiceProvider.GetRequiredService<IConnectionStringResolver>(); |
||||
|
|
||||
|
foreach (var dbContext in dbContexts) |
||||
|
{ |
||||
|
var connectionString = |
||||
|
await connectionStringResolver.ResolveAsync( |
||||
|
ConnectionStringNameAttribute.GetConnStringName(dbContext.GetType())); |
||||
|
if (connectionString.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
continue; |
||||
|
} |
||||
|
|
||||
|
var mongoUrl = new MongoUrl(connectionString); |
||||
|
var databaseName = mongoUrl.DatabaseName; |
||||
|
var client = new MongoClient(mongoUrl); |
||||
|
|
||||
|
if (databaseName.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
databaseName = ConnectionStringNameAttribute.GetConnStringName(dbContext.GetType()); |
||||
|
} |
||||
|
|
||||
|
(dbContext as AbpMongoDbContext)?.InitializeCollections(client.GetDatabase(databaseName)); |
||||
|
} |
||||
|
|
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
if (tenantId == null) |
||||
|
{ |
||||
|
//Migrating the host database
|
||||
|
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)
|
||||
|
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(5000, 15000)); |
||||
|
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; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,52 @@ |
|||||
|
using System; |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp.Data; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
using Volo.Abp.MongoDB; |
||||
|
using Volo.Abp.MultiTenancy; |
||||
|
using Volo.Abp.Uow; |
||||
|
|
||||
|
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations; |
||||
|
|
||||
|
public class PendingMongoDbMigrationsChecker<TDbContext> : PendingMigrationsCheckerBase |
||||
|
where TDbContext : AbpMongoDbContext |
||||
|
{ |
||||
|
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
protected ICurrentTenant CurrentTenant { get; } |
||||
|
protected IDistributedEventBus DistributedEventBus { get; } |
||||
|
protected string DatabaseName { get; } |
||||
|
|
||||
|
protected PendingMongoDbMigrationsChecker( |
||||
|
IUnitOfWorkManager unitOfWorkManager, |
||||
|
IServiceProvider serviceProvider, |
||||
|
ICurrentTenant currentTenant, |
||||
|
IDistributedEventBus distributedEventBus, |
||||
|
string databaseName) |
||||
|
{ |
||||
|
UnitOfWorkManager = unitOfWorkManager; |
||||
|
ServiceProvider = serviceProvider; |
||||
|
CurrentTenant = currentTenant; |
||||
|
DistributedEventBus = distributedEventBus; |
||||
|
DatabaseName = databaseName; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task CheckAsync() |
||||
|
{ |
||||
|
using (CurrentTenant.Change(null)) |
||||
|
{ |
||||
|
// Create database tables if needed
|
||||
|
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) |
||||
|
{ |
||||
|
await DistributedEventBus.PublishAsync( |
||||
|
new ApplyDatabaseMigrationsEto |
||||
|
{ |
||||
|
DatabaseName = DatabaseName |
||||
|
} |
||||
|
); |
||||
|
|
||||
|
await uow.CompleteAsync(); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -1,69 +1,8 @@ |
|||||
using System; |
using Volo.Abp.DependencyInjection; |
||||
using System.Linq; |
|
||||
using System.Threading.Tasks; |
|
||||
using Microsoft.EntityFrameworkCore; |
|
||||
using Microsoft.Extensions.DependencyInjection; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DependencyInjection; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations |
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations |
||||
{ |
{ |
||||
public abstract class PendingMigrationsCheckerBase<TDbContext> : ITransientDependency |
public abstract class PendingMigrationsCheckerBase : ITransientDependency |
||||
where TDbContext : DbContext |
|
||||
{ |
{ |
||||
protected IUnitOfWorkManager UnitOfWorkManager { get; } |
|
||||
protected IServiceProvider ServiceProvider { get; } |
|
||||
protected ICurrentTenant CurrentTenant { get; } |
|
||||
protected IDistributedEventBus DistributedEventBus { get; } |
|
||||
protected string DatabaseName { get; } |
|
||||
|
|
||||
protected PendingMigrationsCheckerBase( |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
IServiceProvider serviceProvider, |
|
||||
ICurrentTenant currentTenant, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
string databaseName) |
|
||||
{ |
|
||||
UnitOfWorkManager = unitOfWorkManager; |
|
||||
ServiceProvider = serviceProvider; |
|
||||
CurrentTenant = currentTenant; |
|
||||
DistributedEventBus = distributedEventBus; |
|
||||
DatabaseName = databaseName; |
|
||||
} |
|
||||
|
|
||||
public virtual async Task<bool> CheckAsync() |
|
||||
{ |
|
||||
var isMigrationRequired = false; |
|
||||
|
|
||||
using (CurrentTenant.Change(null)) |
|
||||
{ |
|
||||
// Create database tables if needed
|
|
||||
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: false)) |
|
||||
{ |
|
||||
var pendingMigrations = await ServiceProvider |
|
||||
.GetRequiredService<TDbContext>() |
|
||||
.Database |
|
||||
.GetPendingMigrationsAsync(); |
|
||||
|
|
||||
if (pendingMigrations.Any()) |
|
||||
{ |
|
||||
await DistributedEventBus.PublishAsync( |
|
||||
new ApplyDatabaseMigrationsEto |
|
||||
{ |
|
||||
DatabaseName = DatabaseName |
|
||||
} |
|
||||
); |
|
||||
isMigrationRequired = true; |
|
||||
} |
|
||||
|
|
||||
await uow.CompleteAsync(); |
|
||||
} |
|
||||
|
|
||||
return isMigrationRequired; |
|
||||
} |
|
||||
} |
|
||||
} |
} |
||||
} |
} |
||||
Loading…
Reference in new issue