diff --git a/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizer.cs b/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizer.cs new file mode 100644 index 0000000000..7d87d07cff --- /dev/null +++ b/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizer.cs @@ -0,0 +1,157 @@ +using System; +using System.ComponentModel; +using System.Globalization; +using System.Threading.Tasks; +using JetBrains.Annotations; +using Volo.Abp.Auditing; +using Volo.Abp.Domain.Repositories; +using Volo.Abp.EventBus.Distributed; +using Volo.Abp.ObjectMapping; +using Volo.Abp.Uow; + +namespace Volo.Abp.Domain.Entities.Events.Distributed; + +public abstract class ExternalEntitySynchronizer : + ExternalEntitySynchronizer + where TEntity : class, IEntity, IHasRemoteModificationTime + where TExternalEntityEto : EntityEto, IHasModificationTime +{ + private readonly IRepository _repository; + + protected ExternalEntitySynchronizer(IObjectMapper objectMapper, IRepository repository) : + base(objectMapper, repository) + { + _repository = repository; + } + + protected override Task FindLocalEntityAsync(TExternalEntityEto eto) + { + return _repository.FindAsync(GetExternalEntityId(eto)); + } + + protected virtual TKey GetExternalEntityId(TExternalEntityEto eto) + { + var keyType = typeof(TKey); + var keyValue = Check.NotNullOrEmpty(eto.KeysAsString, nameof(eto.KeysAsString)); + + if (keyType == typeof(Guid)) + { + return (TKey)TypeDescriptor.GetConverter(keyType).ConvertFromInvariantString(keyValue); + } + + return (TKey)Convert.ChangeType(keyValue, keyType, CultureInfo.InvariantCulture); + } +} + +public abstract class ExternalEntitySynchronizer : + IDistributedEventHandler>, + IDistributedEventHandler>, + IDistributedEventHandler>, + IUnitOfWorkEnabled + where TEntity : class, IEntity, IHasRemoteModificationTime + where TExternalEntityEto : EntityEto, IHasModificationTime +{ + protected IObjectMapper ObjectMapper { get; } + private readonly IRepository _repository; + + protected virtual bool IgnoreEntityCreatedEvent { get; set; } + protected virtual bool IgnoreEntityUpdatedEvent { get; set; } + protected virtual bool IgnoreEntityDeletedEvent { get; set; } + + public ExternalEntitySynchronizer( + IObjectMapper objectMapper, + IRepository repository) + { + ObjectMapper = objectMapper; + _repository = repository; + } + + public virtual async Task HandleEventAsync(EntityCreatedEto eventData) + { + if (IgnoreEntityCreatedEvent) + { + return; + } + + await CreateOrUpdateEntityAsync(eventData.Entity); + } + + public virtual async Task HandleEventAsync(EntityUpdatedEto eventData) + { + if (IgnoreEntityUpdatedEvent) + { + return; + } + + await CreateOrUpdateEntityAsync(eventData.Entity); + } + + public virtual async Task HandleEventAsync(EntityDeletedEto eventData) + { + if (IgnoreEntityDeletedEvent) + { + return; + } + + await TryDeleteEntityAsync(eventData.Entity); + } + + protected virtual async Task CreateOrUpdateEntityAsync(TExternalEntityEto eto) + { + var localEntity = await FindLocalEntityAsync(eto); + + if (!await IsEtoNewerAsync(eto, localEntity)) + { + return; + } + + if (localEntity == null) + { + localEntity = await MapToEntityAsync(eto); + ObjectHelper.TrySetProperty(localEntity, x => x.RemoteLastModificationTime, () => eto.LastModificationTime); + + await _repository.InsertAsync(localEntity, true); + } + else + { + await MapToEntityAsync(eto, localEntity); + ObjectHelper.TrySetProperty(localEntity, x => x.RemoteLastModificationTime, () => eto.LastModificationTime); + + await _repository.UpdateAsync(localEntity, true); + } + } + + protected virtual Task MapToEntityAsync(TExternalEntityEto eto) + { + return Task.FromResult(ObjectMapper.Map(eto)); + } + + protected virtual Task MapToEntityAsync(TExternalEntityEto eto, TEntity localEntity) + { + ObjectMapper.Map(eto, localEntity); + return Task.CompletedTask; + } + + protected virtual async Task TryDeleteEntityAsync(TExternalEntityEto eto) + { + var localEntity = await FindLocalEntityAsync(eto); + + if (localEntity == null) + { + return; + } + + await _repository.DeleteAsync(localEntity, true); + } + + [ItemCanBeNull] + protected abstract Task FindLocalEntityAsync(TExternalEntityEto eto); + + protected virtual Task IsEtoNewerAsync(TExternalEntityEto eto, [CanBeNull] TEntity localEntity) + { + return Task.FromResult( + localEntity?.RemoteLastModificationTime == null || + eto.LastModificationTime > localEntity.RemoteLastModificationTime + ); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/Distributed/IHasRemoteModificationTime.cs b/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/Distributed/IHasRemoteModificationTime.cs new file mode 100644 index 0000000000..2f1cb1fc69 --- /dev/null +++ b/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/Distributed/IHasRemoteModificationTime.cs @@ -0,0 +1,11 @@ +using System; + +namespace Volo.Abp.Domain.Entities.Events.Distributed; + +public interface IHasRemoteModificationTime +{ + /// + /// The last modified time for the synchronized remote entity. + /// + DateTime? RemoteLastModificationTime { get; } +} \ No newline at end of file diff --git a/framework/test/Volo.Abp.Ddd.Tests/Volo.Abp.Ddd.Tests.csproj b/framework/test/Volo.Abp.Ddd.Tests/Volo.Abp.Ddd.Tests.csproj index b6a478866f..ea331dcac3 100644 --- a/framework/test/Volo.Abp.Ddd.Tests/Volo.Abp.Ddd.Tests.csproj +++ b/framework/test/Volo.Abp.Ddd.Tests/Volo.Abp.Ddd.Tests.csproj @@ -8,7 +8,10 @@ + + + diff --git a/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/Book.cs b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/Book.cs new file mode 100644 index 0000000000..91f4b40d8c --- /dev/null +++ b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/Book.cs @@ -0,0 +1,19 @@ +using System; + +namespace Volo.Abp.Domain.Entities.Events.Distributed.ExternalEntitySynchronizers; + +public class Book : Entity, IHasRemoteModificationTime +{ + public virtual DateTime? RemoteLastModificationTime { get; protected set; } + + public virtual int Sold { get; set; } + + protected Book() + { + } + + public Book(Guid id, int sold) : base(id) + { + Sold = sold; + } +} \ No newline at end of file diff --git a/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/BookEntityJsonConverter.cs b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/BookEntityJsonConverter.cs new file mode 100644 index 0000000000..a49e5f7048 --- /dev/null +++ b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/BookEntityJsonConverter.cs @@ -0,0 +1,42 @@ +using System; +using System.Text.Json; +using System.Text.Json.Serialization; + +namespace Volo.Abp.Domain.Entities.Events.Distributed.ExternalEntitySynchronizers; + +public class BookEntityJsonConverter : JsonConverter +{ + private JsonSerializerOptions _writeJsonSerializerOptions; + + public override Book Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options) + { + var jsonDocument = JsonDocument.ParseValue(ref reader); + + if (jsonDocument.RootElement.ValueKind != JsonValueKind.Object) + { + throw new JsonException("RootElement's ValueKind is not Object!"); + } + + var entity = (Book)jsonDocument.RootElement.Deserialize(typeToConvert); + + if (entity == null) + { + throw new JsonException("RootElement's ValueKind is not Object!"); + } + + ObjectHelper.TrySetProperty(entity, x => x.RemoteLastModificationTime, () => + { + var property = jsonDocument.RootElement.GetProperty("RemoteLastModificationTime"); + + return property.ValueKind == JsonValueKind.Null ? null : property.GetDateTime(); + }); + + return entity; + } + + public override void Write(Utf8JsonWriter writer, Book value, JsonSerializerOptions options) + { + _writeJsonSerializerOptions ??= JsonSerializerOptionsHelper.Create(options, this); + JsonSerializer.Serialize(writer, value, _writeJsonSerializerOptions); + } +} \ No newline at end of file diff --git a/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/BookSynchronizer.cs b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/BookSynchronizer.cs new file mode 100644 index 0000000000..492e7d81f0 --- /dev/null +++ b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/BookSynchronizer.cs @@ -0,0 +1,13 @@ +using System; +using Volo.Abp.DependencyInjection; +using Volo.Abp.Domain.Repositories; +using Volo.Abp.ObjectMapping; + +namespace Volo.Abp.Domain.Entities.Events.Distributed.ExternalEntitySynchronizers; + +public class BookSynchronizer : ExternalEntitySynchronizer, ITransientDependency +{ + public BookSynchronizer(IObjectMapper objectMapper, IRepository repository) : base(objectMapper, repository) + { + } +} \ No newline at end of file diff --git a/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/ExternalEntitySynchronizer_Tests.cs b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/ExternalEntitySynchronizer_Tests.cs new file mode 100644 index 0000000000..73b7f2c1c7 --- /dev/null +++ b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/ExternalEntitySynchronizer_Tests.cs @@ -0,0 +1,181 @@ +using System; +using System.Collections.Generic; +using System.Threading.Tasks; +using AutoMapper; +using Microsoft.Extensions.DependencyInjection; +using Shouldly; +using Volo.Abp.Autofac; +using Volo.Abp.AutoMapper; +using Volo.Abp.Data; +using Volo.Abp.Domain.Repositories; +using Volo.Abp.Domain.Repositories.MemoryDb; +using Volo.Abp.MemoryDb; +using Volo.Abp.Modularity; +using Volo.Abp.Testing; +using Volo.Abp.Uow; +using Xunit; + +namespace Volo.Abp.Domain.Entities.Events.Distributed.ExternalEntitySynchronizers; + +public class ExternalEntitySynchronizer_Tests : AbpIntegratedTest +{ + [Fact] + public async Task Should_Handle_Entity_Created_Event() + { + var bookId = Guid.NewGuid(); + + var uowManager = GetRequiredService(); + using var uow = uowManager.Begin(); + + var bookSynchronizer = GetRequiredService(); + var repository = GetRequiredService>(); + + (await repository.FindAsync(bookId)).ShouldBeNull(); + + var remoteBookEto = new RemoteBookEto { + KeysAsString = bookId.ToString(), LastModificationTime = DateTime.Now, Sold = 1 + }; + + await bookSynchronizer.HandleEventAsync(new EntityCreatedEto(remoteBookEto)); + + var book = await repository.FindAsync(bookId); + book.ShouldNotBeNull(); + book.RemoteLastModificationTime.ShouldBe(remoteBookEto.LastModificationTime); + book.Sold.ShouldBe(1); + } + + [Fact] + public async Task Should_Handle_Entity_Update_Event() + { + var bookId = Guid.NewGuid(); + + var uowManager = GetRequiredService(); + using var uow = uowManager.Begin(); + + var bookSynchronizer = GetRequiredService(); + var repository = GetRequiredService>(); + + (await repository.FindAsync(bookId)).ShouldBeNull(); + + var remoteBookEto = new RemoteBookEto { + KeysAsString = bookId.ToString(), LastModificationTime = DateTime.Now, Sold = 1 + }; + + await bookSynchronizer.HandleEventAsync(new EntityUpdatedEto(remoteBookEto)); + + var book = await repository.FindAsync(bookId); + book.ShouldNotBeNull(); + book.RemoteLastModificationTime.ShouldBe(remoteBookEto.LastModificationTime); + book.Sold.ShouldBe(1); + + remoteBookEto.LastModificationTime = DateTime.Now; + remoteBookEto.Sold = 2; + + await bookSynchronizer.HandleEventAsync(new EntityUpdatedEto(remoteBookEto)); + + book = await repository.FindAsync(bookId); + book.ShouldNotBeNull(); + book.RemoteLastModificationTime.ShouldBe(remoteBookEto.LastModificationTime); + book.Sold.ShouldBe(2); + + // Should skip synchronizing older remote entities. + var originalLastModificationTime = remoteBookEto.LastModificationTime; + remoteBookEto.LastModificationTime = remoteBookEto.LastModificationTime.Value.AddTicks(-1); + remoteBookEto.Sold = 3; + + await bookSynchronizer.HandleEventAsync(new EntityUpdatedEto(remoteBookEto)); + + book = await repository.FindAsync(bookId); + book.ShouldNotBeNull(); + book.RemoteLastModificationTime.ShouldBe(originalLastModificationTime); + book.Sold.ShouldBe(2); + } + + [Fact] + public async Task Should_Handle_Entity_Deleted_Event() + { + var bookId = Guid.NewGuid(); + + var uowManager = GetRequiredService(); + using var uow = uowManager.Begin(); + + var bookSynchronizer = GetRequiredService(); + var repository = GetRequiredService>(); + + await repository.InsertAsync(new Book(bookId, 1), true); + + var book = await repository.FindAsync(bookId); + book.ShouldNotBeNull(); + book.Id.ShouldBe(bookId); + book.RemoteLastModificationTime.ShouldBeNull(); + + var remoteBookEto = new RemoteBookEto { + KeysAsString = bookId.ToString(), LastModificationTime = DateTime.Now, Sold = 1 + }; + + await bookSynchronizer.HandleEventAsync(new EntityDeletedEto(remoteBookEto)); + + (await repository.FindAsync(bookId)).ShouldBeNull(); + + await bookSynchronizer.HandleEventAsync(new EntityDeletedEto(remoteBookEto)); + + (await repository.FindAsync(bookId)).ShouldBeNull(); + } + + protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) + { + options.UseAutofac(); + } + + [DependsOn( + typeof(AbpAutofacModule), + typeof(AbpMemoryDbModule), + typeof(AbpDddDomainModule), + typeof(AbpAutoMapperModule) + )] + public class TestModule : AbpModule + { + public override void ConfigureServices(ServiceConfigurationContext context) + { + var connStr = Guid.NewGuid().ToString(); + + Configure(options => + { + options.ConnectionStrings.Default = connStr; + }); + + context.Services.AddMemoryDbContext(options => + { + options.AddDefaultRepositories(includeAllEntities: true); + }); + + Configure(options => + { + options.JsonSerializerOptions.Converters.Add(new BookEntityJsonConverter()); + }); + + context.Services.AddAutoMapperObjectMapper(); + Configure(options => + { + options.AddMaps(validate: true); + }); + } + } + + public class MyMemoryDbContext : MemoryDbContext + { + public override IReadOnlyList GetEntityTypes() + { + return new List { typeof(Book) }; + } + } + + public class MyAutoMapperProfile : Profile + { + public MyAutoMapperProfile() + { + CreateMap(MemberList.None) + .ForMember(x => x.Id, options => options.MapFrom(x => Guid.Parse(x.KeysAsString))); + } + } +} \ No newline at end of file diff --git a/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/RemoteBookEto.cs b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/RemoteBookEto.cs new file mode 100644 index 0000000000..4741c8d3de --- /dev/null +++ b/framework/test/Volo.Abp.Ddd.Tests/Volo/Abp/Domain/Entities/Events/Distributed/ExternalEntitySynchronizers/RemoteBookEto.cs @@ -0,0 +1,11 @@ +using System; +using Volo.Abp.Auditing; + +namespace Volo.Abp.Domain.Entities.Events.Distributed.ExternalEntitySynchronizers; + +public class RemoteBookEto : EntityEto, IHasModificationTime +{ + public DateTime? LastModificationTime { get; set; } + + public int Sold { get; set; } +} \ No newline at end of file