29 changed files with 240 additions and 793 deletions
@ -1,94 +0,0 @@ |
|||||
using EShopOnAbp.AdministrationService.EntityFrameworkCore; |
|
||||
using EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.EfCore; |
|
||||
using Serilog; |
|
||||
using System; |
|
||||
using System.Linq; |
|
||||
using System.Threading.Tasks; |
|
||||
using Volo.Abp.Authorization.Permissions; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.PermissionManagement; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.AdministrationService.DbMigrations; |
|
||||
|
|
||||
public class AdministrationServiceDatabaseMigrationEventHandler |
|
||||
: DatabaseEfCoreMigrationEventHandler<AdministrationServiceDbContext>, |
|
||||
IDistributedEventHandler<ApplyDatabaseMigrationsEto> |
|
||||
{ |
|
||||
private readonly IPermissionDefinitionManager _permissionDefinitionManager; |
|
||||
private readonly IPermissionDataSeeder _permissionDataSeeder; |
|
||||
|
|
||||
public AdministrationServiceDatabaseMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IPermissionDefinitionManager permissionDefinitionManager, |
|
||||
IPermissionDataSeeder permissionDataSeeder, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
IAbpDistributedLock distributedLockProvider |
|
||||
) : base( |
|
||||
currentTenant, |
|
||||
unitOfWorkManager, |
|
||||
tenantStore, |
|
||||
distributedEventBus, |
|
||||
AdministrationServiceDbProperties.ConnectionStringName, |
|
||||
distributedLockProvider |
|
||||
) |
|
||||
{ |
|
||||
_permissionDefinitionManager = permissionDefinitionManager; |
|
||||
_permissionDataSeeder = permissionDataSeeder; |
|
||||
} |
|
||||
|
|
||||
public async Task HandleEventAsync(ApplyDatabaseMigrationsEto eventData) |
|
||||
{ |
|
||||
if (eventData.DatabaseName != DatabaseName) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
try |
|
||||
{ |
|
||||
await using (var handle = await DistributedLockProvider.TryAcquireAsync(DatabaseName)) |
|
||||
{ |
|
||||
Log.Information("AdministrationService acquired lock for db migration and seeding..."); |
|
||||
|
|
||||
if (handle != null) |
|
||||
{ |
|
||||
await MigrateDatabaseSchemaAsync(); |
|
||||
await SeedDataAsync(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
catch (Exception ex) |
|
||||
{ |
|
||||
await HandleErrorOnApplyDatabaseMigrationAsync(eventData, ex); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private async Task SeedDataAsync() |
|
||||
{ |
|
||||
using (var uow = UnitOfWorkManager.Begin(requiresNew: true, isTransactional: true)) |
|
||||
{ |
|
||||
var multiTenancySide = MultiTenancySides.Host; |
|
||||
|
|
||||
var permissionNames = _permissionDefinitionManager |
|
||||
.GetPermissions() |
|
||||
.Where(p => p.MultiTenancySide.HasFlag(multiTenancySide)) |
|
||||
.Where(p => !p.Providers.Any() || |
|
||||
p.Providers.Contains(RolePermissionValueProvider.ProviderName)) |
|
||||
.Select(p => p.Name) |
|
||||
.ToArray(); |
|
||||
|
|
||||
await _permissionDataSeeder.SeedAsync( |
|
||||
RolePermissionValueProvider.ProviderName, |
|
||||
"admin", |
|
||||
permissionNames |
|
||||
); |
|
||||
|
|
||||
await uow.CompleteAsync(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,66 +0,0 @@ |
|||||
using EShopOnAbp.CatalogService.MongoDB; |
|
||||
using EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.MongoDb; |
|
||||
using System; |
|
||||
using System.Threading.Tasks; |
|
||||
using Serilog; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.CatalogService.DbMigrations; |
|
||||
|
|
||||
public class CatalogServiceDatabaseMigrationEventHandler |
|
||||
: DatabaseMongoDbMigrationEventHandler<CatalogServiceMongoDbContext>, |
|
||||
IDistributedEventHandler<ApplyDatabaseMigrationsEto> |
|
||||
{ |
|
||||
public CatalogServiceDatabaseMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
IServiceProvider serviceProvider, |
|
||||
IAbpDistributedLock distributedLockProvider |
|
||||
) : base( |
|
||||
currentTenant, |
|
||||
unitOfWorkManager, |
|
||||
tenantStore, |
|
||||
distributedEventBus, |
|
||||
CatalogServiceDbProperties.ConnectionStringName, |
|
||||
serviceProvider, |
|
||||
distributedLockProvider) |
|
||||
{ |
|
||||
} |
|
||||
|
|
||||
public async Task HandleEventAsync(ApplyDatabaseMigrationsEto eventData) |
|
||||
{ |
|
||||
if (eventData.DatabaseName != DatabaseName) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
if (eventData.TenantId != null) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
try |
|
||||
{ |
|
||||
Log.Information("CatalogService has acquired lock for db migration..."); |
|
||||
|
|
||||
await using (var handle = await DistributedLockProvider.TryAcquireAsync(DatabaseName)) |
|
||||
{ |
|
||||
if (handle != null) |
|
||||
{ |
|
||||
Log.Information("CatalogService is migrating database..."); |
|
||||
await MigrateDatabaseSchemaAsync(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
catch (Exception ex) |
|
||||
{ |
|
||||
await HandleErrorOnApplyDatabaseMigrationAsync(eventData, ex); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,9 +0,0 @@ |
|||||
using Volo.Abp.Domain.Entities.Events.Distributed; |
|
||||
using Volo.Abp.EventBus; |
|
||||
|
|
||||
namespace EShopOnAbp.IdentityService.DbMigrations; |
|
||||
|
|
||||
[EventName("abp.identity.apply_database_seeds")] |
|
||||
public class ApplyDatabaseSeedsEto : EtoBase |
|
||||
{ |
|
||||
} |
|
||||
@ -1,21 +0,0 @@ |
|||||
using System.Threading.Tasks; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DependencyInjection; |
|
||||
using Volo.Abp.EventBus; |
|
||||
|
|
||||
namespace EShopOnAbp.IdentityService.DbMigrations; |
|
||||
|
|
||||
public class DataSeederEventHandler : ILocalEventHandler<ApplyDatabaseSeedsEto>, ITransientDependency |
|
||||
{ |
|
||||
protected IDataSeeder DataSeeder { get; } |
|
||||
|
|
||||
public DataSeederEventHandler(IDataSeeder dataSeeder) |
|
||||
{ |
|
||||
DataSeeder = dataSeeder; |
|
||||
} |
|
||||
|
|
||||
public async Task HandleEventAsync(ApplyDatabaseSeedsEto eventData) |
|
||||
{ |
|
||||
await DataSeeder.SeedAsync(); |
|
||||
} |
|
||||
} |
|
||||
@ -1,90 +0,0 @@ |
|||||
using EShopOnAbp.IdentityService.EntityFrameworkCore; |
|
||||
using EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.EfCore; |
|
||||
using Serilog; |
|
||||
using System; |
|
||||
using System.Threading.Tasks; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.EventBus.Local; |
|
||||
using Volo.Abp.Identity; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.IdentityService.DbMigrations; |
|
||||
|
|
||||
public class IdentityServiceDatabaseMigrationEventHandler |
|
||||
: DatabaseEfCoreMigrationEventHandler<IdentityServiceDbContext>, |
|
||||
IDistributedEventHandler<ApplyDatabaseMigrationsEto> |
|
||||
{ |
|
||||
private readonly IIdentityDataSeeder _identityDataSeeder; |
|
||||
private readonly IdentityServerDataSeeder _identityServerDataSeeder; |
|
||||
private readonly ILocalEventBus _localEventBus; |
|
||||
|
|
||||
public IdentityServiceDatabaseMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IIdentityDataSeeder identityDataSeeder, |
|
||||
IdentityServerDataSeeder identityServerDataSeeder, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
ILocalEventBus localEventBus, |
|
||||
IAbpDistributedLock distributedLockProvider |
|
||||
) : base( |
|
||||
currentTenant, |
|
||||
unitOfWorkManager, |
|
||||
tenantStore, |
|
||||
distributedEventBus, |
|
||||
IdentityServiceDbProperties.ConnectionStringName, |
|
||||
distributedLockProvider) |
|
||||
{ |
|
||||
_identityDataSeeder = identityDataSeeder; |
|
||||
_identityServerDataSeeder = identityServerDataSeeder; |
|
||||
_localEventBus = localEventBus; |
|
||||
} |
|
||||
|
|
||||
public async Task HandleEventAsync(ApplyDatabaseMigrationsEto eventData) |
|
||||
{ |
|
||||
if (eventData.DatabaseName != DatabaseName) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
try |
|
||||
{ |
|
||||
await using (var handle = await DistributedLockProvider.TryAcquireAsync(DatabaseName)) |
|
||||
{ |
|
||||
Log.Information("IdentityService has acquired lock for db migration..."); |
|
||||
|
|
||||
if (handle != null) |
|
||||
{ |
|
||||
Log.Information("IdentityService is migrating database..."); |
|
||||
await MigrateDatabaseSchemaAsync(); |
|
||||
Log.Information("IdentityService is seeding data..."); |
|
||||
await SeedDataAsync( |
|
||||
adminEmail: IdentityServiceDbProperties.DefaultAdminEmailAddress, |
|
||||
adminPassword: IdentityServiceDbProperties.DefaultAdminPassword |
|
||||
); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
await _localEventBus.PublishAsync(new ApplyDatabaseSeedsEto()); |
|
||||
} |
|
||||
catch (Exception ex) |
|
||||
{ |
|
||||
await HandleErrorOnApplyDatabaseMigrationAsync(eventData, ex); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private async Task SeedDataAsync(string adminEmail, string adminPassword) |
|
||||
{ |
|
||||
Log.Information($"Seeding IdentityServer data..."); |
|
||||
await _identityServerDataSeeder.SeedAsync(); |
|
||||
|
|
||||
Log.Information($"Seeding user data..."); |
|
||||
await _identityDataSeeder.SeedAsync( |
|
||||
adminEmail, |
|
||||
adminPassword |
|
||||
); |
|
||||
} |
|
||||
} |
|
||||
@ -1,70 +0,0 @@ |
|||||
using EShopOnAbp.OrderingService.EntityFrameworkCore; |
|
||||
using EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.EfCore; |
|
||||
using System; |
|
||||
using System.Threading.Tasks; |
|
||||
using Serilog; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.OrderingService.DbMigrations; |
|
||||
|
|
||||
public class OrderingServiceDatabaseMigrationEventHandler |
|
||||
: DatabaseEfCoreMigrationEventHandler<OrderingServiceDbContext>, |
|
||||
IDistributedEventHandler<ApplyDatabaseMigrationsEto> |
|
||||
{ |
|
||||
private readonly IDataSeeder _dataSeeder; |
|
||||
|
|
||||
public OrderingServiceDatabaseMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
IDataSeeder dataSeeder, |
|
||||
IAbpDistributedLock distributedLockProvider) |
|
||||
: base( |
|
||||
currentTenant, |
|
||||
unitOfWorkManager, |
|
||||
tenantStore, |
|
||||
distributedEventBus, |
|
||||
OrderingServiceDbProperties.ConnectionStringName, |
|
||||
distributedLockProvider) |
|
||||
{ |
|
||||
_dataSeeder = dataSeeder; |
|
||||
} |
|
||||
|
|
||||
public async Task HandleEventAsync(ApplyDatabaseMigrationsEto eventData) |
|
||||
{ |
|
||||
if (eventData.DatabaseName != DatabaseName) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
if (eventData.TenantId != null) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
try |
|
||||
{ |
|
||||
await using (var handle = await DistributedLockProvider.TryAcquireAsync(DatabaseName)) |
|
||||
{ |
|
||||
Log.Information("OrderingService has acquired lock for db migration..."); |
|
||||
|
|
||||
if (handle != null) |
|
||||
{ |
|
||||
Log.Information("OrderingService is migrating database..."); |
|
||||
await MigrateDatabaseSchemaAsync(); |
|
||||
Log.Information("OrderingService is seeding data..."); |
|
||||
await _dataSeeder.SeedAsync(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
catch (Exception ex) |
|
||||
{ |
|
||||
await HandleErrorOnApplyDatabaseMigrationAsync(eventData, ex); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,64 +0,0 @@ |
|||||
using EShopOnAbp.PaymentService.EntityFrameworkCore; |
|
||||
using EShopOnAbp.Shared.Hosting.Microservices.DbMigrations.EfCore; |
|
||||
using System; |
|
||||
using System.Threading.Tasks; |
|
||||
using Serilog; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
using Volo.Abp.EventBus.Distributed; |
|
||||
using Volo.Abp.MultiTenancy; |
|
||||
using Volo.Abp.Uow; |
|
||||
|
|
||||
namespace EShopOnAbp.PaymentService.DbMigrations; |
|
||||
|
|
||||
public class PaymentServiceDatabaseMigrationEventHandler |
|
||||
: DatabaseEfCoreMigrationEventHandler<PaymentServiceDbContext>, |
|
||||
IDistributedEventHandler<ApplyDatabaseMigrationsEto> |
|
||||
{ |
|
||||
public PaymentServiceDatabaseMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
IAbpDistributedLock distributedLockProvider) |
|
||||
: base( |
|
||||
currentTenant, |
|
||||
unitOfWorkManager, |
|
||||
tenantStore, |
|
||||
distributedEventBus, |
|
||||
PaymentServiceDbProperties.ConnectionStringName, |
|
||||
distributedLockProvider) |
|
||||
{ |
|
||||
} |
|
||||
|
|
||||
public async Task HandleEventAsync(ApplyDatabaseMigrationsEto eventData) |
|
||||
{ |
|
||||
if (eventData.DatabaseName != DatabaseName) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
if (eventData.TenantId != null) |
|
||||
{ |
|
||||
return; |
|
||||
} |
|
||||
|
|
||||
try |
|
||||
{ |
|
||||
Log.Information("PaymentService has acquired lock for db migration..."); |
|
||||
|
|
||||
await using (var handle = await DistributedLockProvider.TryAcquireAsync(DatabaseName)) |
|
||||
{ |
|
||||
if (handle != null) |
|
||||
{ |
|
||||
Log.Information("PaymentService is migrating database..."); |
|
||||
await MigrateDatabaseSchemaAsync(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
catch (Exception ex) |
|
||||
{ |
|
||||
await HandleErrorOnApplyDatabaseMigrationAsync(eventData, ex); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,7 +0,0 @@ |
|||||
using Volo.Abp.DependencyInjection; |
|
||||
|
|
||||
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations; |
|
||||
|
|
||||
public abstract class DatabaseMigrationEventHandlerBase : ITransientDependency |
|
||||
{ |
|
||||
} |
|
||||
@ -1,133 +0,0 @@ |
|||||
using Microsoft.EntityFrameworkCore; |
|
||||
using Microsoft.Extensions.DependencyInjection; |
|
||||
using Microsoft.Extensions.Logging; |
|
||||
using Microsoft.Extensions.Logging.Abstractions; |
|
||||
using Serilog; |
|
||||
using System; |
|
||||
using System.Collections.Generic; |
|
||||
using System.Linq; |
|
||||
using System.Threading.Tasks; |
|
||||
using Volo.Abp; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
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 IAbpDistributedLock DistributedLockProvider { get; } |
|
||||
|
|
||||
protected DatabaseEfCoreMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
string databaseName, |
|
||||
IAbpDistributedLock distributedLockProvider |
|
||||
) |
|
||||
{ |
|
||||
CurrentTenant = currentTenant; |
|
||||
UnitOfWorkManager = unitOfWorkManager; |
|
||||
TenantStore = tenantStore; |
|
||||
DatabaseName = databaseName; |
|
||||
DistributedEventBus = distributedEventBus; |
|
||||
DistributedLockProvider = distributedLockProvider; |
|
||||
|
|
||||
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() |
|
||||
{ |
|
||||
var result = false; |
|
||||
|
|
||||
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; |
|
||||
} |
|
||||
|
|
||||
//Migrating the host database
|
|
||||
Log.Information($"There is no tenant. Migrating {DatabaseName}..."); |
|
||||
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.Warning( |
|
||||
$"Could not apply database migrations. Canceling the operation. TenantId = {eventData.TenantId}, DatabaseName = {eventData.DatabaseName}."); |
|
||||
Log.Error(exception.ToString()); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
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; |
|
||||
} |
|
||||
} |
|
||||
@ -1,149 +0,0 @@ |
|||||
using Microsoft.Extensions.DependencyInjection; |
|
||||
using Microsoft.Extensions.Logging; |
|
||||
using Microsoft.Extensions.Logging.Abstractions; |
|
||||
using MongoDB.Driver; |
|
||||
using System; |
|
||||
using System.Collections.Generic; |
|
||||
using System.Threading.Tasks; |
|
||||
using Volo.Abp; |
|
||||
using Volo.Abp.Data; |
|
||||
using Volo.Abp.DistributedLocking; |
|
||||
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 IAbpDistributedLock DistributedLockProvider { get; } |
|
||||
|
|
||||
|
|
||||
protected DatabaseMongoDbMigrationEventHandler( |
|
||||
ICurrentTenant currentTenant, |
|
||||
IUnitOfWorkManager unitOfWorkManager, |
|
||||
ITenantStore tenantStore, |
|
||||
IDistributedEventBus distributedEventBus, |
|
||||
string databaseName, |
|
||||
IServiceProvider serviceProvider, |
|
||||
IAbpDistributedLock distributedLockProvider |
|
||||
) |
|
||||
{ |
|
||||
CurrentTenant = currentTenant; |
|
||||
UnitOfWorkManager = unitOfWorkManager; |
|
||||
TenantStore = tenantStore; |
|
||||
DatabaseName = databaseName; |
|
||||
ServiceProvider = serviceProvider; |
|
||||
DistributedEventBus = distributedEventBus; |
|
||||
DistributedLockProvider = distributedLockProvider; |
|
||||
|
|
||||
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() |
|
||||
{ |
|
||||
var result = false; |
|
||||
|
|
||||
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; |
|
||||
} |
|
||||
|
|
||||
//Migrating the host 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); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
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; |
|
||||
} |
|
||||
} |
|
||||
@ -1,7 +1,33 @@ |
|||||
using Volo.Abp.DependencyInjection; |
using Serilog; |
||||
|
using System; |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations; |
namespace EShopOnAbp.Shared.Hosting.Microservices.DbMigrations; |
||||
|
|
||||
public abstract class PendingMigrationsCheckerBase : ITransientDependency |
public abstract class PendingMigrationsCheckerBase : ITransientDependency |
||||
{ |
{ |
||||
|
public async Task TryAsync(Func<Task> task, int retryCount = 3) |
||||
|
{ |
||||
|
try |
||||
|
{ |
||||
|
await task(); |
||||
|
} |
||||
|
catch (Exception ex) |
||||
|
{ |
||||
|
retryCount--; |
||||
|
|
||||
|
if (retryCount <= 0) |
||||
|
{ |
||||
|
throw; |
||||
|
} |
||||
|
|
||||
|
Log.Warning($"{ex.GetType().Name} has been thrown. The operation will be tried {retryCount} times more. Exception:\n{ex.Message}"); |
||||
|
|
||||
|
await Task.Delay(RandomHelper.GetRandom(5000, 15000)); |
||||
|
|
||||
|
await TryAsync(task, retryCount); |
||||
|
} |
||||
|
} |
||||
} |
} |
||||
Loading…
Reference in new issue