mirror of https://github.com/abpframework/abp.git
8 changed files with 437 additions and 0 deletions
@ -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<TEntity, TKey, TExternalEntityEto> : |
||||
|
ExternalEntitySynchronizer<TEntity, TExternalEntityEto> |
||||
|
where TEntity : class, IEntity<TKey>, IHasRemoteModificationTime |
||||
|
where TExternalEntityEto : EntityEto, IHasModificationTime |
||||
|
{ |
||||
|
private readonly IRepository<TEntity, TKey> _repository; |
||||
|
|
||||
|
protected ExternalEntitySynchronizer(IObjectMapper objectMapper, IRepository<TEntity, TKey> repository) : |
||||
|
base(objectMapper, repository) |
||||
|
{ |
||||
|
_repository = repository; |
||||
|
} |
||||
|
|
||||
|
protected override Task<TEntity> 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<TEntity, TExternalEntityEto> : |
||||
|
IDistributedEventHandler<EntityCreatedEto<TExternalEntityEto>>, |
||||
|
IDistributedEventHandler<EntityUpdatedEto<TExternalEntityEto>>, |
||||
|
IDistributedEventHandler<EntityDeletedEto<TExternalEntityEto>>, |
||||
|
IUnitOfWorkEnabled |
||||
|
where TEntity : class, IEntity, IHasRemoteModificationTime |
||||
|
where TExternalEntityEto : EntityEto, IHasModificationTime |
||||
|
{ |
||||
|
protected IObjectMapper ObjectMapper { get; } |
||||
|
private readonly IRepository<TEntity> _repository; |
||||
|
|
||||
|
protected virtual bool IgnoreEntityCreatedEvent { get; set; } |
||||
|
protected virtual bool IgnoreEntityUpdatedEvent { get; set; } |
||||
|
protected virtual bool IgnoreEntityDeletedEvent { get; set; } |
||||
|
|
||||
|
public ExternalEntitySynchronizer( |
||||
|
IObjectMapper objectMapper, |
||||
|
IRepository<TEntity> repository) |
||||
|
{ |
||||
|
ObjectMapper = objectMapper; |
||||
|
_repository = repository; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task HandleEventAsync(EntityCreatedEto<TExternalEntityEto> eventData) |
||||
|
{ |
||||
|
if (IgnoreEntityCreatedEvent) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
await CreateOrUpdateEntityAsync(eventData.Entity); |
||||
|
} |
||||
|
|
||||
|
public virtual async Task HandleEventAsync(EntityUpdatedEto<TExternalEntityEto> eventData) |
||||
|
{ |
||||
|
if (IgnoreEntityUpdatedEvent) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
await CreateOrUpdateEntityAsync(eventData.Entity); |
||||
|
} |
||||
|
|
||||
|
public virtual async Task HandleEventAsync(EntityDeletedEto<TExternalEntityEto> 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<TEntity> MapToEntityAsync(TExternalEntityEto eto) |
||||
|
{ |
||||
|
return Task.FromResult(ObjectMapper.Map<TExternalEntityEto, TEntity>(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<TEntity> FindLocalEntityAsync(TExternalEntityEto eto); |
||||
|
|
||||
|
protected virtual Task<bool> IsEtoNewerAsync(TExternalEntityEto eto, [CanBeNull] TEntity localEntity) |
||||
|
{ |
||||
|
return Task.FromResult( |
||||
|
localEntity?.RemoteLastModificationTime == null || |
||||
|
eto.LastModificationTime > localEntity.RemoteLastModificationTime |
||||
|
); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.Domain.Entities.Events.Distributed; |
||||
|
|
||||
|
public interface IHasRemoteModificationTime |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// The last modified time for the synchronized remote entity.
|
||||
|
/// </summary>
|
||||
|
DateTime? RemoteLastModificationTime { get; } |
||||
|
} |
||||
@ -0,0 +1,19 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.Domain.Entities.Events.Distributed.ExternalEntitySynchronizers; |
||||
|
|
||||
|
public class Book : Entity<Guid>, 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; |
||||
|
} |
||||
|
} |
||||
@ -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<Book> |
||||
|
{ |
||||
|
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); |
||||
|
} |
||||
|
} |
||||
@ -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<Book, Guid, RemoteBookEto>, ITransientDependency |
||||
|
{ |
||||
|
public BookSynchronizer(IObjectMapper objectMapper, IRepository<Book, Guid> repository) : base(objectMapper, repository) |
||||
|
{ |
||||
|
} |
||||
|
} |
||||
@ -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<ExternalEntitySynchronizer_Tests.TestModule> |
||||
|
{ |
||||
|
[Fact] |
||||
|
public async Task Should_Handle_Entity_Created_Event() |
||||
|
{ |
||||
|
var bookId = Guid.NewGuid(); |
||||
|
|
||||
|
var uowManager = GetRequiredService<IUnitOfWorkManager>(); |
||||
|
using var uow = uowManager.Begin(); |
||||
|
|
||||
|
var bookSynchronizer = GetRequiredService<BookSynchronizer>(); |
||||
|
var repository = GetRequiredService<IRepository<Book, Guid>>(); |
||||
|
|
||||
|
(await repository.FindAsync(bookId)).ShouldBeNull(); |
||||
|
|
||||
|
var remoteBookEto = new RemoteBookEto { |
||||
|
KeysAsString = bookId.ToString(), LastModificationTime = DateTime.Now, Sold = 1 |
||||
|
}; |
||||
|
|
||||
|
await bookSynchronizer.HandleEventAsync(new EntityCreatedEto<RemoteBookEto>(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<IUnitOfWorkManager>(); |
||||
|
using var uow = uowManager.Begin(); |
||||
|
|
||||
|
var bookSynchronizer = GetRequiredService<BookSynchronizer>(); |
||||
|
var repository = GetRequiredService<IRepository<Book, Guid>>(); |
||||
|
|
||||
|
(await repository.FindAsync(bookId)).ShouldBeNull(); |
||||
|
|
||||
|
var remoteBookEto = new RemoteBookEto { |
||||
|
KeysAsString = bookId.ToString(), LastModificationTime = DateTime.Now, Sold = 1 |
||||
|
}; |
||||
|
|
||||
|
await bookSynchronizer.HandleEventAsync(new EntityUpdatedEto<RemoteBookEto>(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>(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>(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<IUnitOfWorkManager>(); |
||||
|
using var uow = uowManager.Begin(); |
||||
|
|
||||
|
var bookSynchronizer = GetRequiredService<BookSynchronizer>(); |
||||
|
var repository = GetRequiredService<IRepository<Book, Guid>>(); |
||||
|
|
||||
|
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>(remoteBookEto)); |
||||
|
|
||||
|
(await repository.FindAsync(bookId)).ShouldBeNull(); |
||||
|
|
||||
|
await bookSynchronizer.HandleEventAsync(new EntityDeletedEto<RemoteBookEto>(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<AbpDbConnectionOptions>(options => |
||||
|
{ |
||||
|
options.ConnectionStrings.Default = connStr; |
||||
|
}); |
||||
|
|
||||
|
context.Services.AddMemoryDbContext<MyMemoryDbContext>(options => |
||||
|
{ |
||||
|
options.AddDefaultRepositories(includeAllEntities: true); |
||||
|
}); |
||||
|
|
||||
|
Configure<Utf8JsonMemoryDbSerializerOptions>(options => |
||||
|
{ |
||||
|
options.JsonSerializerOptions.Converters.Add(new BookEntityJsonConverter()); |
||||
|
}); |
||||
|
|
||||
|
context.Services.AddAutoMapperObjectMapper<TestModule>(); |
||||
|
Configure<AbpAutoMapperOptions>(options => |
||||
|
{ |
||||
|
options.AddMaps<TestModule>(validate: true); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public class MyMemoryDbContext : MemoryDbContext |
||||
|
{ |
||||
|
public override IReadOnlyList<Type> GetEntityTypes() |
||||
|
{ |
||||
|
return new List<Type> { typeof(Book) }; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public class MyAutoMapperProfile : Profile |
||||
|
{ |
||||
|
public MyAutoMapperProfile() |
||||
|
{ |
||||
|
CreateMap<RemoteBookEto, Book>(MemberList.None) |
||||
|
.ForMember(x => x.Id, options => options.MapFrom(x => Guid.Parse(x.KeysAsString))); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -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; } |
||||
|
} |
||||
Loading…
Reference in new issue