Browse Source

Merge pull request #9909 from abpframework/dist-events

Publish events while unit of work is being completed, in the same unit of work
pull/9943/head
maliming 5 years ago
committed by GitHub
parent
commit
50cf7cd729
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 25
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/BasicAggregateRoot.cs
  2. 15
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/DomainEventRecord.cs
  3. 5
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/DomainEventEntry.cs
  4. 260
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityChangeEventHelper.cs
  5. 32
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityChangeReport.cs
  6. 1
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityChangingEventData.cs
  7. 1
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityCreatingEventData.cs
  8. 1
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityDeletingEventData.cs
  9. 22
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityEventReport.cs
  10. 1
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityUpdatingEventData.cs
  11. 16
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/IEntityChangeEventHelper.cs
  12. 41
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/NullEntityChangeEventHelper.cs
  13. 4
      framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/IGeneratesDomainEvents.cs
  14. 176
      framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/AbpDbContext.cs
  15. 11
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs
  16. 12
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  17. 11
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs
  18. 8
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs
  19. 4
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs
  20. 29
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs
  21. 6
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventBus.cs
  22. 11
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs
  23. 4
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs
  24. 40
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/UnitOfWorkEventPublisher.cs
  25. 59
      framework/src/Volo.Abp.MemoryDb/Volo/Abp/Domain/Repositories/MemoryDb/MemoryDbRepository.cs
  26. 65
      framework/src/Volo.Abp.MongoDB/Volo/Abp/Domain/Repositories/MongoDB/MongoDbRepository.cs
  27. 13
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/AmbientUnitOfWork.cs
  28. 14
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/ChildUnitOfWork.cs
  29. 14
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/EventOrderGenerator.cs
  30. 2
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/IAmbientUnitOfWork.cs
  31. 10
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/IUnitOfWork.cs
  32. 12
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/IUnitOfWorkEventPublisher.cs
  33. 19
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/NullUnitOfWorkEventPublisher.cs
  34. 89
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWork.cs
  35. 29
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWorkEventRecord.cs
  36. 15
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWorkManager.cs
  37. 10
      framework/test/Volo.Abp.AspNetCore.Mvc.Tests/Volo/Abp/AspNetCore/Mvc/Uow/TestUnitOfWork.cs
  38. 10
      framework/test/Volo.Abp.MemoryDb.Tests/Volo/Abp/MemoryDb/DomainEvents/DomainEvents_Tests.cs
  39. 119
      framework/test/Volo.Abp.TestApp/Volo/Abp/TestApp/Testing/DomainEvents_Tests.cs
  40. 16
      framework/test/Volo.Abp.TestApp/Volo/Abp/TestApp/Testing/EntityChangeEvents_Tests.cs
  41. 1
      framework/test/Volo.Abp.TestApp/Volo/Abp/TestApp/Testing/TestAppTestBase.cs
  42. 2
      modules/setting-management/src/Volo.Abp.SettingManagement.Domain/Volo/Abp/SettingManagement/SettingCacheItemInvalidator.cs
  43. 29
      test/DistEvents/DistDemoApp/DemoService.cs
  44. 34
      test/DistEvents/DistDemoApp/DistDemoApp.csproj
  45. 37
      test/DistEvents/DistDemoApp/DistDemoAppModule.cs
  46. 61
      test/DistEvents/DistDemoApp/Migrations/20210825110134_Initial.Designer.cs
  47. 33
      test/DistEvents/DistDemoApp/Migrations/20210825110134_Initial.cs
  48. 95
      test/DistEvents/DistDemoApp/Migrations/20210825112717_Added_Summary_Table.Designer.cs
  49. 34
      test/DistEvents/DistDemoApp/Migrations/20210825112717_Added_Summary_Table.cs
  50. 93
      test/DistEvents/DistDemoApp/Migrations/TodoDbContextModelSnapshot.cs
  51. 39
      test/DistEvents/DistDemoApp/MyProjectNameHostedService.cs
  52. 57
      test/DistEvents/DistDemoApp/Program.cs
  53. 28
      test/DistEvents/DistDemoApp/TodoDbContext.cs
  54. 29
      test/DistEvents/DistDemoApp/TodoDbContextFactory.cs
  55. 65
      test/DistEvents/DistDemoApp/TodoEventHandler.cs
  56. 15
      test/DistEvents/DistDemoApp/TodoItem.cs
  57. 12
      test/DistEvents/DistDemoApp/TodoItemEto.cs
  58. 24
      test/DistEvents/DistDemoApp/TodoItemObjectMapper.cs
  59. 41
      test/DistEvents/DistDemoApp/TodoSummary.cs
  60. 5
      test/DistEvents/DistDemoApp/appsettings.json
  61. 16
      test/DistEvents/DistEventsDemo.sln

25
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/BasicAggregateRoot.cs

@ -1,6 +1,7 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Collections.ObjectModel; using System.Collections.ObjectModel;
using Volo.Abp.Uow;
namespace Volo.Abp.Domain.Entities namespace Volo.Abp.Domain.Entities
{ {
@ -9,15 +10,15 @@ namespace Volo.Abp.Domain.Entities
IAggregateRoot, IAggregateRoot,
IGeneratesDomainEvents IGeneratesDomainEvents
{ {
private readonly ICollection<object> _distributedEvents = new Collection<object>(); private readonly ICollection<DomainEventRecord> _distributedEvents = new Collection<DomainEventRecord>();
private readonly ICollection<object> _localEvents = new Collection<object>(); private readonly ICollection<DomainEventRecord> _localEvents = new Collection<DomainEventRecord>();
public virtual IEnumerable<object> GetLocalEvents() public virtual IEnumerable<DomainEventRecord> GetLocalEvents()
{ {
return _localEvents; return _localEvents;
} }
public virtual IEnumerable<object> GetDistributedEvents() public virtual IEnumerable<DomainEventRecord> GetDistributedEvents()
{ {
return _distributedEvents; return _distributedEvents;
} }
@ -34,12 +35,12 @@ namespace Volo.Abp.Domain.Entities
protected virtual void AddLocalEvent(object eventData) protected virtual void AddLocalEvent(object eventData)
{ {
_localEvents.Add(eventData); _localEvents.Add(new DomainEventRecord(eventData, EventOrderGenerator.GetNext()));
} }
protected virtual void AddDistributedEvent(object eventData) protected virtual void AddDistributedEvent(object eventData)
{ {
_distributedEvents.Add(eventData); _distributedEvents.Add(new DomainEventRecord(eventData, EventOrderGenerator.GetNext()));
} }
} }
@ -48,8 +49,8 @@ namespace Volo.Abp.Domain.Entities
IAggregateRoot<TKey>, IAggregateRoot<TKey>,
IGeneratesDomainEvents IGeneratesDomainEvents
{ {
private readonly ICollection<object> _distributedEvents = new Collection<object>(); private readonly ICollection<DomainEventRecord> _distributedEvents = new Collection<DomainEventRecord>();
private readonly ICollection<object> _localEvents = new Collection<object>(); private readonly ICollection<DomainEventRecord> _localEvents = new Collection<DomainEventRecord>();
protected BasicAggregateRoot() protected BasicAggregateRoot()
{ {
@ -62,12 +63,12 @@ namespace Volo.Abp.Domain.Entities
} }
public virtual IEnumerable<object> GetLocalEvents() public virtual IEnumerable<DomainEventRecord> GetLocalEvents()
{ {
return _localEvents; return _localEvents;
} }
public virtual IEnumerable<object> GetDistributedEvents() public virtual IEnumerable<DomainEventRecord> GetDistributedEvents()
{ {
return _distributedEvents; return _distributedEvents;
} }
@ -84,12 +85,12 @@ namespace Volo.Abp.Domain.Entities
protected virtual void AddLocalEvent(object eventData) protected virtual void AddLocalEvent(object eventData)
{ {
_localEvents.Add(eventData); _localEvents.Add(new DomainEventRecord(eventData, EventOrderGenerator.GetNext()));
} }
protected virtual void AddDistributedEvent(object eventData) protected virtual void AddDistributedEvent(object eventData)
{ {
_distributedEvents.Add(eventData); _distributedEvents.Add(new DomainEventRecord(eventData, EventOrderGenerator.GetNext()));
} }
} }
} }

15
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/DomainEventRecord.cs

@ -0,0 +1,15 @@
namespace Volo.Abp.Domain.Entities
{
public class DomainEventRecord
{
public object EventData { get; }
public long EventOrder { get; }
public DomainEventRecord(object eventData, long eventOrder)
{
EventData = eventData;
EventOrder = eventOrder;
}
}
}

5
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/DomainEventEntry.cs

@ -8,11 +8,14 @@ namespace Volo.Abp.Domain.Entities.Events
public object SourceEntity { get; } public object SourceEntity { get; }
public object EventData { get; } public object EventData { get; }
public long EventOrder { get; }
public DomainEventEntry(object sourceEntity, object eventData) public DomainEventEntry(object sourceEntity, object eventData, long eventOrder)
{ {
SourceEntity = sourceEntity; SourceEntity = sourceEntity;
EventData = eventData; EventData = eventData;
EventOrder = eventOrder;
} }
} }
} }

260
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityChangeEventHelper.cs

@ -1,11 +1,8 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
using Volo.Abp.Auditing;
using Volo.Abp.DependencyInjection; using Volo.Abp.DependencyInjection;
using Volo.Abp.Domain.Entities.Events.Distributed; using Volo.Abp.Domain.Entities.Events.Distributed;
using Volo.Abp.DynamicProxy; using Volo.Abp.DynamicProxy;
@ -21,6 +18,8 @@ namespace Volo.Abp.Domain.Entities.Events
/// </summary> /// </summary>
public class EntityChangeEventHelper : IEntityChangeEventHelper, ITransientDependency public class EntityChangeEventHelper : IEntityChangeEventHelper, ITransientDependency
{ {
private const string UnitOfWorkEventRecordEntityPropName = "_Abp_Entity";
public ILogger<EntityChangeEventHelper> Logger { get; set; } public ILogger<EntityChangeEventHelper> Logger { get; set; }
public ILocalEventBus LocalEventBus { get; set; } public ILocalEventBus LocalEventBus { get; set; }
public IDistributedEventBus DistributedEventBus { get; set; } public IDistributedEventBus DistributedEventBus { get; set; }
@ -43,37 +42,25 @@ namespace Volo.Abp.Domain.Entities.Events
Logger = NullLogger<EntityChangeEventHelper>.Instance; Logger = NullLogger<EntityChangeEventHelper>.Instance;
} }
public async Task TriggerEventsAsync(EntityChangeReport changeReport) public virtual void PublishEntityCreatingEvent(object entity)
{ {
await TriggerEventsInternalAsync(changeReport); TriggerEventWithEntity(
if (changeReport.IsEmpty() || UnitOfWorkManager.Current == null)
{
return;
}
await UnitOfWorkManager.Current.SaveChangesAsync();
}
public virtual async Task TriggerEntityCreatingEventAsync(object entity)
{
await TriggerEventWithEntity(
LocalEventBus, LocalEventBus,
#pragma warning disable 618
typeof(EntityCreatingEventData<>), typeof(EntityCreatingEventData<>),
#pragma warning restore 618
entity, entity,
entity, entity
true
); );
} }
public virtual async Task TriggerEntityCreatedEventOnUowCompletedAsync(object entity) public virtual void PublishEntityCreatedEvent(object entity)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
LocalEventBus, LocalEventBus,
typeof(EntityCreatedEventData<>), typeof(EntityCreatedEventData<>),
entity, entity,
entity, entity
false
); );
if (ShouldPublishDistributedEventForEntity(entity)) if (ShouldPublishDistributedEventForEntity(entity))
@ -81,12 +68,11 @@ namespace Volo.Abp.Domain.Entities.Events
var eto = EntityToEtoMapper.Map(entity); var eto = EntityToEtoMapper.Map(entity);
if (eto != null) if (eto != null)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
DistributedEventBus, DistributedEventBus,
typeof(EntityCreatedEto<>), typeof(EntityCreatedEto<>),
eto, eto,
entity, entity
false
); );
} }
} }
@ -103,25 +89,25 @@ namespace Volo.Abp.Domain.Entities.Events
); );
} }
public virtual async Task TriggerEntityUpdatingEventAsync(object entity) public virtual void PublishEntityUpdatingEvent(object entity)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
LocalEventBus, LocalEventBus,
#pragma warning disable 618
typeof(EntityUpdatingEventData<>), typeof(EntityUpdatingEventData<>),
#pragma warning restore 618
entity, entity,
entity, entity
true
); );
} }
public virtual async Task TriggerEntityUpdatedEventOnUowCompletedAsync(object entity) public virtual void PublishEntityUpdatedEvent(object entity)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
LocalEventBus, LocalEventBus,
typeof(EntityUpdatedEventData<>), typeof(EntityUpdatedEventData<>),
entity, entity,
entity, entity
false
); );
if (ShouldPublishDistributedEventForEntity(entity)) if (ShouldPublishDistributedEventForEntity(entity))
@ -129,36 +115,35 @@ namespace Volo.Abp.Domain.Entities.Events
var eto = EntityToEtoMapper.Map(entity); var eto = EntityToEtoMapper.Map(entity);
if (eto != null) if (eto != null)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
DistributedEventBus, DistributedEventBus,
typeof(EntityUpdatedEto<>), typeof(EntityUpdatedEto<>),
eto, eto,
entity, entity
false
); );
} }
} }
} }
public virtual async Task TriggerEntityDeletingEventAsync(object entity) public virtual void PublishEntityDeletingEvent(object entity)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
LocalEventBus, LocalEventBus,
#pragma warning disable 618
typeof(EntityDeletingEventData<>), typeof(EntityDeletingEventData<>),
#pragma warning restore 618
entity, entity,
entity, entity
true
); );
} }
public virtual async Task TriggerEntityDeletedEventOnUowCompletedAsync(object entity) public virtual void PublishEntityDeletedEvent(object entity)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
LocalEventBus, LocalEventBus,
typeof(EntityDeletedEventData<>), typeof(EntityDeletedEventData<>),
entity, entity,
entity, entity
false
); );
if (ShouldPublishDistributedEventForEntity(entity)) if (ShouldPublishDistributedEventForEntity(entity))
@ -166,183 +151,78 @@ namespace Volo.Abp.Domain.Entities.Events
var eto = EntityToEtoMapper.Map(entity); var eto = EntityToEtoMapper.Map(entity);
if (eto != null) if (eto != null)
{ {
await TriggerEventWithEntity( TriggerEventWithEntity(
DistributedEventBus, DistributedEventBus,
typeof(EntityDeletedEto<>), typeof(EntityDeletedEto<>),
eto, eto,
entity, entity
false
); );
} }
} }
} }
protected virtual async Task TriggerEventsInternalAsync(EntityChangeReport changeReport) protected virtual void TriggerEventWithEntity(
{
await TriggerEntityChangeEvents(changeReport.ChangedEntities);
await TriggerLocalEvents(changeReport.DomainEvents);
await TriggerDistributedEvents(changeReport.DistributedEvents);
}
protected virtual async Task TriggerEntityChangeEvents(List<EntityChangeEntry> changedEntities)
{
foreach (var changedEntity in changedEntities)
{
switch (changedEntity.ChangeType)
{
case EntityChangeType.Created:
await TriggerEntityCreatingEventAsync(changedEntity.Entity);
await TriggerEntityCreatedEventOnUowCompletedAsync(changedEntity.Entity);
break;
case EntityChangeType.Updated:
await TriggerEntityUpdatingEventAsync(changedEntity.Entity);
await TriggerEntityUpdatedEventOnUowCompletedAsync(changedEntity.Entity);
break;
case EntityChangeType.Deleted:
await TriggerEntityDeletingEventAsync(changedEntity.Entity);
await TriggerEntityDeletedEventOnUowCompletedAsync(changedEntity.Entity);
break;
default:
throw new AbpException("Unknown EntityChangeType: " + changedEntity.ChangeType);
}
}
}
protected virtual async Task TriggerLocalEvents(List<DomainEventEntry> localEvents)
{
foreach (var localEvent in localEvents)
{
await LocalEventBus.PublishAsync(localEvent.EventData.GetType(), localEvent.EventData);
}
}
protected virtual async Task TriggerDistributedEvents(List<DomainEventEntry> distributedEvents)
{
foreach (var distributedEvent in distributedEvents)
{
await DistributedEventBus.PublishAsync(distributedEvent.EventData.GetType(),
distributedEvent.EventData);
}
}
protected virtual async Task TriggerEventWithEntity(
IEventBus eventPublisher, IEventBus eventPublisher,
Type genericEventType, Type genericEventType,
object entityOrEto, object entityOrEto,
object originalEntity, object originalEntity)
bool triggerInCurrentUnitOfWork)
{ {
var entityType = ProxyHelper.UnProxy(entityOrEto).GetType(); var entityType = ProxyHelper.UnProxy(entityOrEto).GetType();
var eventType = genericEventType.MakeGenericType(entityType); var eventType = genericEventType.MakeGenericType(entityType);
var eventData = Activator.CreateInstance(eventType, entityOrEto);
var currentUow = UnitOfWorkManager.Current; var currentUow = UnitOfWorkManager.Current;
if (triggerInCurrentUnitOfWork || currentUow == null) if (currentUow == null)
{ {
await eventPublisher.PublishAsync( Logger.LogWarning("UnitOfWorkManager.Current is null! Can not publish the event.");
eventType,
Activator.CreateInstance(eventType, entityOrEto)
);
return; return;
} }
var eventList = GetEventList(currentUow); var eventRecord = new UnitOfWorkEventRecord(eventType, eventData, EventOrderGenerator.GetNext())
var isFirstEvent = !eventList.Any(); {
Properties =
eventList.AddUniqueEvent(eventPublisher, eventType, entityOrEto, originalEntity); {
{ UnitOfWorkEventRecordEntityPropName, originalEntity },
/* Register to OnCompleted if this is the first item. }
* Other items will already be in the list once the UOW completes. };
/* We are trying to eliminate same events for the same entity.
* In this way, for example, we don't trigger update event for an entity multiple times
* even if it is updated multiple times in the current UOW.
*/ */
if (isFirstEvent)
if (eventPublisher == DistributedEventBus)
{ {
currentUow.OnCompleted( currentUow.AddOrReplaceDistributedEvent(
async () => eventRecord,
{ otherRecord => IsSameEntityEventRecord(eventRecord, otherRecord)
foreach (var eventEntry in eventList)
{
try
{
await eventEntry.EventBus.PublishAsync(
eventEntry.EventType,
Activator.CreateInstance(eventEntry.EventType, eventEntry.EntityOrEto)
);
}
catch (Exception ex)
{
Logger.LogError(
$"Caught an exception while publishing the event '{eventType.FullName}' for the entity '{entityOrEto}'");
Logger.LogException(ex);
}
}
}
); );
} }
} else
private EntityChangeEventList GetEventList(IUnitOfWork currentUow)
{
return (EntityChangeEventList) currentUow.Items.GetOrAdd(
"AbpEntityChangeEventList",
() => new EntityChangeEventList()
);
}
private class EntityChangeEventList : List<EntityChangeEventEntry>
{
public void AddUniqueEvent(IEventBus eventBus, Type eventType, object entityOrEto, object originalEntity)
{ {
var newEntry = new EntityChangeEventEntry(eventBus, eventType, entityOrEto, originalEntity); currentUow.AddOrReplaceLocalEvent(
eventRecord,
//Latest "same" event overrides the previous events. otherRecord => IsSameEntityEventRecord(eventRecord, otherRecord)
for (var i = 0; i < Count; i++) );
{
if (this[i].IsSameEvent(newEntry))
{
this[i] = newEntry;
return;
}
}
//If this is a "new" event, add to the end
Add(newEntry);
} }
} }
private class EntityChangeEventEntry public bool IsSameEntityEventRecord(UnitOfWorkEventRecord record1, UnitOfWorkEventRecord record2)
{ {
public IEventBus EventBus { get; } if (record1.EventType != record2.EventType)
public Type EventType { get; }
public object EntityOrEto { get; }
public object OriginalEntity { get; }
public EntityChangeEventEntry(IEventBus eventBus, Type eventType, object entityOrEto, object originalEntity)
{ {
EventType = eventType; return false;
EntityOrEto = entityOrEto;
OriginalEntity = originalEntity;
EventBus = eventBus;
} }
public bool IsSameEvent(EntityChangeEventEntry otherEntry) var record1OriginalEntity = record1.Properties.GetOrDefault(UnitOfWorkEventRecordEntityPropName) as IEntity;
{ var record2OriginalEntity = record2.Properties.GetOrDefault(UnitOfWorkEventRecordEntityPropName) as IEntity;
if (EventBus != otherEntry.EventBus || EventType != otherEntry.EventType)
{
return false;
}
var originalEntityRef = OriginalEntity as IEntity;
var otherOriginalEntityRef = otherEntry.OriginalEntity as IEntity;
if (originalEntityRef == null || otherOriginalEntityRef == null)
{
return false;
}
return EntityHelper.EntityEquals(originalEntityRef, otherOriginalEntityRef); if (record1OriginalEntity == null || record2OriginalEntity == null)
{
return false;
} }
return EntityHelper.EntityEquals(record1OriginalEntity, record2OriginalEntity);
} }
} }
} }

32
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityChangeReport.cs

@ -1,32 +0,0 @@
using System.Collections.Generic;
namespace Volo.Abp.Domain.Entities.Events
{
public class EntityChangeReport
{
public List<EntityChangeEntry> ChangedEntities { get; }
public List<DomainEventEntry> DomainEvents { get; }
public List<DomainEventEntry> DistributedEvents { get; }
public EntityChangeReport()
{
ChangedEntities = new List<EntityChangeEntry>();
DomainEvents = new List<DomainEventEntry>();
DistributedEvents = new List<DomainEventEntry>();
}
public bool IsEmpty()
{
return ChangedEntities.Count <= 0 &&
DomainEvents.Count <= 0 &&
DistributedEvents.Count <= 0;
}
public override string ToString()
{
return $"[EntityChangeReport] ChangedEntities: {ChangedEntities.Count}, DomainEvents: {DomainEvents.Count}, DistributedEvents: {DistributedEvents.Count}";
}
}
}

1
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityChangingEventData.cs

@ -8,6 +8,7 @@ namespace Volo.Abp.Domain.Entities.Events
/// </summary> /// </summary>
/// <typeparam name="TEntity">Entity type</typeparam> /// <typeparam name="TEntity">Entity type</typeparam>
[Serializable] [Serializable]
[Obsolete("This event is no longer needed and identical to EntityChangedEventData. Please use EntityChangedEventData instead.")]
public class EntityChangingEventData<TEntity> : EntityEventData<TEntity> public class EntityChangingEventData<TEntity> : EntityEventData<TEntity>
{ {
/// <summary> /// <summary>

1
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityCreatingEventData.cs

@ -7,6 +7,7 @@ namespace Volo.Abp.Domain.Entities.Events
/// </summary> /// </summary>
/// <typeparam name="TEntity">Entity type</typeparam> /// <typeparam name="TEntity">Entity type</typeparam>
[Serializable] [Serializable]
[Obsolete("This event is no longer needed and identical to EntityCreatedEventData. Please use EntityCreatedEventData instead.")]
public class EntityCreatingEventData<TEntity> : EntityChangingEventData<TEntity> public class EntityCreatingEventData<TEntity> : EntityChangingEventData<TEntity>
{ {
/// <summary> /// <summary>

1
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityDeletingEventData.cs

@ -7,6 +7,7 @@ namespace Volo.Abp.Domain.Entities.Events
/// </summary> /// </summary>
/// <typeparam name="TEntity">Entity type</typeparam> /// <typeparam name="TEntity">Entity type</typeparam>
[Serializable] [Serializable]
[Obsolete("This event is no longer needed and identical to EntityDeleteEventData. Please use EntityDeleteEventData instead.")]
public class EntityDeletingEventData<TEntity> : EntityChangingEventData<TEntity> public class EntityDeletingEventData<TEntity> : EntityChangingEventData<TEntity>
{ {
/// <summary> /// <summary>

22
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityEventReport.cs

@ -0,0 +1,22 @@
using System.Collections.Generic;
namespace Volo.Abp.Domain.Entities.Events
{
public class EntityEventReport
{
public List<DomainEventEntry> DomainEvents { get; }
public List<DomainEventEntry> DistributedEvents { get; }
public EntityEventReport()
{
DomainEvents = new List<DomainEventEntry>();
DistributedEvents = new List<DomainEventEntry>();
}
public override string ToString()
{
return $"[{nameof(EntityEventReport)}] DomainEvents: {DomainEvents.Count}, DistributedEvents: {DistributedEvents.Count}";
}
}
}

1
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/EntityUpdatingEventData.cs

@ -7,6 +7,7 @@ namespace Volo.Abp.Domain.Entities.Events
/// </summary> /// </summary>
/// <typeparam name="TEntity">Entity type</typeparam> /// <typeparam name="TEntity">Entity type</typeparam>
[Serializable] [Serializable]
[Obsolete("This event is no longer needed and identical to EntityUpdatedEventData. Please use EntityUpdatedEventData instead.")]
public class EntityUpdatingEventData<TEntity> : EntityChangingEventData<TEntity> public class EntityUpdatingEventData<TEntity> : EntityChangingEventData<TEntity>
{ {
/// <summary> /// <summary>

16
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/IEntityChangeEventHelper.cs

@ -1,5 +1,3 @@
using System.Threading.Tasks;
namespace Volo.Abp.Domain.Entities.Events namespace Volo.Abp.Domain.Entities.Events
{ {
/// <summary> /// <summary>
@ -7,15 +5,13 @@ namespace Volo.Abp.Domain.Entities.Events
/// </summary> /// </summary>
public interface IEntityChangeEventHelper public interface IEntityChangeEventHelper
{ {
Task TriggerEventsAsync(EntityChangeReport changeReport); void PublishEntityCreatingEvent(object entity);
void PublishEntityCreatedEvent(object entity);
Task TriggerEntityCreatingEventAsync(object entity);
Task TriggerEntityCreatedEventOnUowCompletedAsync(object entity);
Task TriggerEntityUpdatingEventAsync(object entity); void PublishEntityUpdatingEvent(object entity);
Task TriggerEntityUpdatedEventOnUowCompletedAsync(object entity); void PublishEntityUpdatedEvent(object entity);
Task TriggerEntityDeletingEventAsync(object entity); void PublishEntityDeletingEvent(object entity);
Task TriggerEntityDeletedEventOnUowCompletedAsync(object entity); void PublishEntityDeletedEvent(object entity);
} }
} }

41
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/Events/NullEntityChangeEventHelper.cs

@ -1,5 +1,3 @@
using System.Threading.Tasks;
namespace Volo.Abp.Domain.Entities.Events namespace Volo.Abp.Domain.Entities.Events
{ {
/// <summary> /// <summary>
@ -14,57 +12,30 @@ namespace Volo.Abp.Domain.Entities.Events
private NullEntityChangeEventHelper() private NullEntityChangeEventHelper()
{ {
}
public Task TriggerEntityCreatingEventAsync(object entity)
{
return Task.CompletedTask;
}
public Task TriggerEntityCreatedEventAsync(object entity)
{
return Task.CompletedTask;
}
public Task TriggerEntityCreatedEventOnUowCompletedAsync(object entity)
{
return Task.CompletedTask;
}
public Task TriggerEntityUpdatingEventAsync(object entity)
{
return Task.CompletedTask;
} }
public Task TriggerEntityUpdatedEventAsync(object entity) public void PublishEntityCreatingEvent(object entity)
{ {
return Task.CompletedTask;
} }
public Task TriggerEntityUpdatedEventOnUowCompletedAsync(object entity) public void PublishEntityCreatedEvent(object entity)
{ {
return Task.CompletedTask;
} }
public Task TriggerEntityDeletingEventAsync(object entity) public void PublishEntityUpdatingEvent(object entity)
{ {
return Task.CompletedTask;
} }
public Task TriggerEntityDeletedEventAsync(object entity) public void PublishEntityUpdatedEvent(object entity)
{ {
return Task.CompletedTask;
} }
public Task TriggerEntityDeletedEventOnUowCompletedAsync(object entity) public void PublishEntityDeletingEvent(object entity)
{ {
return Task.CompletedTask;
} }
public Task TriggerEventsAsync(EntityChangeReport changeReport) public void PublishEntityDeletedEvent(object entity)
{ {
return Task.CompletedTask;
} }
} }
} }

4
framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/Entities/IGeneratesDomainEvents.cs

@ -6,9 +6,9 @@ namespace Volo.Abp.Domain.Entities
public interface IGeneratesDomainEvents public interface IGeneratesDomainEvents
{ {
IEnumerable<object> GetLocalEvents(); IEnumerable<DomainEventRecord> GetLocalEvents();
IEnumerable<object> GetDistributedEvents(); IEnumerable<DomainEventRecord> GetDistributedEvents();
void ClearLocalEvents(); void ClearLocalEvents();

176
framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/AbpDbContext.cs

@ -22,6 +22,8 @@ using Volo.Abp.Domain.Repositories;
using Volo.Abp.EntityFrameworkCore.EntityHistory; using Volo.Abp.EntityFrameworkCore.EntityHistory;
using Volo.Abp.EntityFrameworkCore.Modeling; using Volo.Abp.EntityFrameworkCore.Modeling;
using Volo.Abp.EntityFrameworkCore.ValueConverters; using Volo.Abp.EntityFrameworkCore.ValueConverters;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Local;
using Volo.Abp.Guids; using Volo.Abp.Guids;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.ObjectExtending; using Volo.Abp.ObjectExtending;
@ -59,6 +61,10 @@ namespace Volo.Abp.EntityFrameworkCore
public IUnitOfWorkManager UnitOfWorkManager => LazyServiceProvider.LazyGetRequiredService<IUnitOfWorkManager>(); public IUnitOfWorkManager UnitOfWorkManager => LazyServiceProvider.LazyGetRequiredService<IUnitOfWorkManager>();
public IClock Clock => LazyServiceProvider.LazyGetRequiredService<IClock>(); public IClock Clock => LazyServiceProvider.LazyGetRequiredService<IClock>();
public IDistributedEventBus DistributedEventBus => LazyServiceProvider.LazyGetRequiredService<IDistributedEventBus>();
public ILocalEventBus LocalEventBus => LazyServiceProvider.LazyGetRequiredService<ILocalEventBus>();
public ILogger<AbpDbContext<TDbContext>> Logger => LazyServiceProvider.LazyGetService<ILogger<AbpDbContext<TDbContext>>>(NullLogger<AbpDbContext<TDbContext>>.Instance); public ILogger<AbpDbContext<TDbContext>> Logger => LazyServiceProvider.LazyGetService<ILogger<AbpDbContext<TDbContext>>>(NullLogger<AbpDbContext<TDbContext>>.Instance);
@ -150,20 +156,21 @@ namespace Volo.Abp.EntityFrameworkCore
try try
{ {
var auditLog = AuditingManager?.Current?.Log; var auditLog = AuditingManager?.Current?.Log;
List<EntityChangeInfo> entityChangeList = null; List<EntityChangeInfo> entityChangeList = null;
if (auditLog != null) if (auditLog != null)
{ {
entityChangeList = EntityHistoryHelper.CreateChangeList(ChangeTracker.Entries().ToList()); entityChangeList = EntityHistoryHelper.CreateChangeList(ChangeTracker.Entries().ToList());
} }
var changeReport = ApplyAbpConcepts(); ApplyAbpConcepts();
var eventReport = CreateEventReport();
var result = await base.SaveChangesAsync(acceptAllChangesOnSuccess, cancellationToken); var result = await base.SaveChangesAsync(acceptAllChangesOnSuccess, cancellationToken);
await EntityChangeEventHelper.TriggerEventsAsync(changeReport); PublishEntityEvents(eventReport);
if (auditLog != null) if (entityChangeList != null)
{ {
EntityHistoryHelper.UpdateChangeList(entityChangeList); EntityHistoryHelper.UpdateChangeList(entityChangeList);
auditLog.EntityChanges.AddRange(entityChangeList); auditLog.EntityChanges.AddRange(entityChangeList);
@ -182,6 +189,23 @@ namespace Volo.Abp.EntityFrameworkCore
} }
} }
private void PublishEntityEvents(EntityEventReport changeReport)
{
foreach (var localEvent in changeReport.DomainEvents)
{
UnitOfWorkManager.Current?.AddOrReplaceLocalEvent(
new UnitOfWorkEventRecord(localEvent.EventData.GetType(), localEvent.EventData, localEvent.EventOrder)
);
}
foreach (var distributedEvent in changeReport.DistributedEvents)
{
UnitOfWorkManager.Current?.AddOrReplaceDistributedEvent(
new UnitOfWorkEventRecord(distributedEvent.EventData.GetType(), distributedEvent.EventData, distributedEvent.EventOrder)
);
}
}
/// <summary> /// <summary>
/// This method will call the DbContext <see cref="SaveChangesAsync(bool, CancellationToken)"/> method directly of EF Core, which doesn't apply concepts of abp. /// This method will call the DbContext <see cref="SaveChangesAsync(bool, CancellationToken)"/> method directly of EF Core, which doesn't apply concepts of abp.
/// </summary> /// </summary>
@ -202,11 +226,18 @@ namespace Volo.Abp.EntityFrameworkCore
ChangeTracker.CascadeDeleteTiming = CascadeTiming.OnSaveChanges; ChangeTracker.CascadeDeleteTiming = CascadeTiming.OnSaveChanges;
ChangeTracker.Tracked += ChangeTracker_Tracked; ChangeTracker.Tracked += ChangeTracker_Tracked;
ChangeTracker.StateChanged += ChangeTracker_StateChanged;
} }
protected virtual void ChangeTracker_Tracked(object sender, EntityTrackedEventArgs e) protected virtual void ChangeTracker_Tracked(object sender, EntityTrackedEventArgs e)
{ {
FillExtraPropertiesForTrackedEntities(e); FillExtraPropertiesForTrackedEntities(e);
PublishEventsForTrackedEntity(e.Entry);
}
protected virtual void ChangeTracker_StateChanged(object sender, EntityStateChangedEventArgs e)
{
PublishEventsForTrackedEntity(e.Entry);
} }
protected virtual void FillExtraPropertiesForTrackedEntities(EntityTrackedEventArgs e) protected virtual void FillExtraPropertiesForTrackedEntities(EntityTrackedEventArgs e)
@ -254,37 +285,107 @@ namespace Volo.Abp.EntityFrameworkCore
} }
} }
} }
protected virtual EntityChangeReport ApplyAbpConcepts() private void PublishEventsForTrackedEntity(EntityEntry entry)
{ {
var changeReport = new EntityChangeReport(); switch (entry.State)
{
case EntityState.Added:
EntityChangeEventHelper.PublishEntityCreatingEvent(entry.Entity);
EntityChangeEventHelper.PublishEntityCreatedEvent(entry.Entity);
break;
case EntityState.Modified:
if (entry.Properties.Any(x => x.IsModified && x.Metadata.ValueGenerated == ValueGenerated.Never))
{
if (entry.Entity is ISoftDelete && entry.Entity.As<ISoftDelete>().IsDeleted)
{
EntityChangeEventHelper.PublishEntityDeletingEvent(entry.Entity);
EntityChangeEventHelper.PublishEntityDeletedEvent(entry.Entity);
}
else
{
EntityChangeEventHelper.PublishEntityUpdatingEvent(entry.Entity);
EntityChangeEventHelper.PublishEntityUpdatedEvent(entry.Entity);
}
}
break;
case EntityState.Deleted:
EntityChangeEventHelper.PublishEntityDeletingEvent(entry.Entity);
EntityChangeEventHelper.PublishEntityDeletedEvent(entry.Entity);
break;
}
}
protected virtual void ApplyAbpConcepts()
{
foreach (var entry in ChangeTracker.Entries().ToList()) foreach (var entry in ChangeTracker.Entries().ToList())
{ {
ApplyAbpConcepts(entry, changeReport); ApplyAbpConcepts(entry);
} }
return changeReport;
} }
protected virtual EntityEventReport CreateEventReport()
{
var eventReport = new EntityEventReport();
foreach (var entry in ChangeTracker.Entries().ToList())
{
var generatesDomainEventsEntity = entry.Entity as IGeneratesDomainEvents;
if (generatesDomainEventsEntity == null)
{
continue;
}
protected virtual void ApplyAbpConcepts(EntityEntry entry, EntityChangeReport changeReport) var localEvents = generatesDomainEventsEntity.GetLocalEvents()?.ToArray();
if (localEvents != null && localEvents.Any())
{
eventReport.DomainEvents.AddRange(
localEvents.Select(
eventRecord => new DomainEventEntry(
entry.Entity,
eventRecord.EventData,
eventRecord.EventOrder
)
)
);
generatesDomainEventsEntity.ClearLocalEvents();
}
var distributedEvents = generatesDomainEventsEntity.GetDistributedEvents()?.ToArray();
if (distributedEvents != null && distributedEvents.Any())
{
eventReport.DistributedEvents.AddRange(
distributedEvents.Select(
eventRecord => new DomainEventEntry(
entry.Entity,
eventRecord.EventData,
eventRecord.EventOrder)
)
);
generatesDomainEventsEntity.ClearDistributedEvents();
}
}
return eventReport;
}
protected virtual void ApplyAbpConcepts(EntityEntry entry)
{ {
switch (entry.State) switch (entry.State)
{ {
case EntityState.Added: case EntityState.Added:
ApplyAbpConceptsForAddedEntity(entry, changeReport); ApplyAbpConceptsForAddedEntity(entry);
break; break;
case EntityState.Modified: case EntityState.Modified:
ApplyAbpConceptsForModifiedEntity(entry, changeReport); ApplyAbpConceptsForModifiedEntity(entry);
break; break;
case EntityState.Deleted: case EntityState.Deleted:
ApplyAbpConceptsForDeletedEntity(entry, changeReport); ApplyAbpConceptsForDeletedEntity(entry);
break; break;
} }
HandleExtraPropertiesOnSave(entry); HandleExtraPropertiesOnSave(entry);
AddDomainEvents(changeReport, entry.Entity);
} }
protected virtual void HandleExtraPropertiesOnSave(EntityEntry entry) protected virtual void HandleExtraPropertiesOnSave(EntityEntry entry)
@ -361,15 +462,14 @@ namespace Volo.Abp.EntityFrameworkCore
} }
} }
protected virtual void ApplyAbpConceptsForAddedEntity(EntityEntry entry, EntityChangeReport changeReport) protected virtual void ApplyAbpConceptsForAddedEntity(EntityEntry entry)
{ {
CheckAndSetId(entry); CheckAndSetId(entry);
SetConcurrencyStampIfNull(entry); SetConcurrencyStampIfNull(entry);
SetCreationAuditProperties(entry); SetCreationAuditProperties(entry);
changeReport.ChangedEntities.Add(new EntityChangeEntry(entry.Entity, EntityChangeType.Created));
} }
protected virtual void ApplyAbpConceptsForModifiedEntity(EntityEntry entry, EntityChangeReport changeReport) protected virtual void ApplyAbpConceptsForModifiedEntity(EntityEntry entry)
{ {
if (entry.State == EntityState.Modified && entry.Properties.Any(x => x.IsModified && x.Metadata.ValueGenerated == ValueGenerated.Never)) if (entry.State == EntityState.Modified && entry.Properties.Any(x => x.IsModified && x.Metadata.ValueGenerated == ValueGenerated.Never))
{ {
@ -379,24 +479,17 @@ namespace Volo.Abp.EntityFrameworkCore
if (entry.Entity is ISoftDelete && entry.Entity.As<ISoftDelete>().IsDeleted) if (entry.Entity is ISoftDelete && entry.Entity.As<ISoftDelete>().IsDeleted)
{ {
SetDeletionAuditProperties(entry); SetDeletionAuditProperties(entry);
changeReport.ChangedEntities.Add(new EntityChangeEntry(entry.Entity, EntityChangeType.Deleted));
}
else
{
changeReport.ChangedEntities.Add(new EntityChangeEntry(entry.Entity, EntityChangeType.Updated));
} }
} }
} }
protected virtual void ApplyAbpConceptsForDeletedEntity(EntityEntry entry, EntityChangeReport changeReport) protected virtual void ApplyAbpConceptsForDeletedEntity(EntityEntry entry)
{ {
if (TryCancelDeletionForSoftDelete(entry)) if (TryCancelDeletionForSoftDelete(entry))
{ {
UpdateConcurrencyStamp(entry); UpdateConcurrencyStamp(entry);
SetDeletionAuditProperties(entry); SetDeletionAuditProperties(entry);
} }
changeReport.ChangedEntities.Add(new EntityChangeEntry(entry.Entity, EntityChangeType.Deleted));
} }
protected virtual bool IsHardDeleted(EntityEntry entry) protected virtual bool IsHardDeleted(EntityEntry entry)
@ -410,29 +503,6 @@ namespace Volo.Abp.EntityFrameworkCore
return hardDeletedEntities.Contains(entry.Entity); return hardDeletedEntities.Contains(entry.Entity);
} }
protected virtual void AddDomainEvents(EntityChangeReport changeReport, object entityAsObj)
{
var generatesDomainEventsEntity = entityAsObj as IGeneratesDomainEvents;
if (generatesDomainEventsEntity == null)
{
return;
}
var localEvents = generatesDomainEventsEntity.GetLocalEvents()?.ToArray();
if (localEvents != null && localEvents.Any())
{
changeReport.DomainEvents.AddRange(localEvents.Select(eventData => new DomainEventEntry(entityAsObj, eventData)));
generatesDomainEventsEntity.ClearLocalEvents();
}
var distributedEvents = generatesDomainEventsEntity.GetDistributedEvents()?.ToArray();
if (distributedEvents != null && distributedEvents.Any())
{
changeReport.DistributedEvents.AddRange(distributedEvents.Select(eventData => new DomainEventEntry(entityAsObj, eventData)));
generatesDomainEventsEntity.ClearDistributedEvents();
}
}
protected virtual void UpdateConcurrencyStamp(EntityEntry entry) protected virtual void UpdateConcurrencyStamp(EntityEntry entry)
{ {
var entity = entry.Entity as IHasConcurrencyStamp; var entity = entry.Entity as IHasConcurrencyStamp;

11
framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs

@ -12,6 +12,7 @@ using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Kafka; using Volo.Abp.Kafka;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.Threading; using Volo.Abp.Threading;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus.Kafka namespace Volo.Abp.EventBus.Kafka
{ {
@ -33,6 +34,7 @@ namespace Volo.Abp.EventBus.Kafka
public KafkaDistributedEventBus( public KafkaDistributedEventBus(
IServiceScopeFactory serviceScopeFactory, IServiceScopeFactory serviceScopeFactory,
ICurrentTenant currentTenant, ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager,
IOptions<AbpKafkaEventBusOptions> abpKafkaEventBusOptions, IOptions<AbpKafkaEventBusOptions> abpKafkaEventBusOptions,
IKafkaMessageConsumerFactory messageConsumerFactory, IKafkaMessageConsumerFactory messageConsumerFactory,
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions, IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions,
@ -40,7 +42,7 @@ namespace Volo.Abp.EventBus.Kafka
IProducerPool producerPool, IProducerPool producerPool,
IEventErrorHandler errorHandler, IEventErrorHandler errorHandler,
IOptions<AbpEventBusOptions> abpEventBusOptions) IOptions<AbpEventBusOptions> abpEventBusOptions)
: base(serviceScopeFactory, currentTenant, errorHandler) : base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler)
{ {
AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value;
AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value;
@ -165,11 +167,16 @@ namespace Volo.Abp.EventBus.Kafka
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear()); GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear());
} }
public override async Task PublishAsync(Type eventType, object eventData) protected override async Task PublishToEventBusAsync(Type eventType, object eventData)
{ {
await PublishAsync(eventType, eventData, new Headers {{"messageId", Serializer.Serialize(Guid.NewGuid())}}, null); await PublishAsync(eventType, eventData, new Headers {{"messageId", Serializer.Serialize(Guid.NewGuid())}}, null);
} }
protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord)
{
unitOfWork.AddOrReplaceDistributedEvent(eventRecord);
}
public virtual async Task PublishAsync(Type eventType, object eventData, Headers headers, Dictionary<string, object> headersArguments) public virtual async Task PublishAsync(Type eventType, object eventData, Headers headers, Dictionary<string, object> headersArguments)
{ {
await PublishAsync(AbpKafkaEventBusOptions.TopicName, eventType, eventData, headers, headersArguments); await PublishAsync(AbpKafkaEventBusOptions.TopicName, eventType, eventData, headers, headersArguments);

12
framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs

@ -13,6 +13,7 @@ using Volo.Abp.EventBus.Distributed;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.RabbitMQ; using Volo.Abp.RabbitMQ;
using Volo.Abp.Threading; using Volo.Abp.Threading;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus.RabbitMq namespace Volo.Abp.EventBus.RabbitMq
{ {
@ -44,9 +45,10 @@ namespace Volo.Abp.EventBus.RabbitMq
IOptions<AbpDistributedEventBusOptions> distributedEventBusOptions, IOptions<AbpDistributedEventBusOptions> distributedEventBusOptions,
IRabbitMqMessageConsumerFactory messageConsumerFactory, IRabbitMqMessageConsumerFactory messageConsumerFactory,
ICurrentTenant currentTenant, ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager,
IEventErrorHandler errorHandler, IEventErrorHandler errorHandler,
IOptions<AbpEventBusOptions> abpEventBusOptions) IOptions<AbpEventBusOptions> abpEventBusOptions)
: base(serviceScopeFactory, currentTenant, errorHandler) : base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler)
{ {
ConnectionPool = connectionPool; ConnectionPool = connectionPool;
Serializer = serializer; Serializer = serializer;
@ -189,14 +191,18 @@ namespace Volo.Abp.EventBus.RabbitMq
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear()); GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear());
} }
public override async Task PublishAsync(Type eventType, object eventData) protected override async Task PublishToEventBusAsync(Type eventType, object eventData)
{ {
await PublishAsync(eventType, eventData, null); await PublishAsync(eventType, eventData, null);
} }
public Task PublishAsync(Type eventType, object eventData, IBasicProperties properties, Dictionary<string, object> headersArguments = null) protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord)
{ {
unitOfWork.AddOrReplaceDistributedEvent(eventRecord);
}
public Task PublishAsync(Type eventType, object eventData, IBasicProperties properties, Dictionary<string, object> headersArguments = null)
{
var eventName = EventNameAttribute.GetNameOrDefault(eventType); var eventName = EventNameAttribute.GetNameOrDefault(eventType);
var body = Serializer.Serialize(eventData); var body = Serializer.Serialize(eventData);

11
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs

@ -10,6 +10,7 @@ using Volo.Abp.DependencyInjection;
using Volo.Abp.EventBus.Distributed; using Volo.Abp.EventBus.Distributed;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.Threading; using Volo.Abp.Threading;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus.Rebus namespace Volo.Abp.EventBus.Rebus
{ {
@ -28,11 +29,12 @@ namespace Volo.Abp.EventBus.Rebus
public RebusDistributedEventBus( public RebusDistributedEventBus(
IServiceScopeFactory serviceScopeFactory, IServiceScopeFactory serviceScopeFactory,
ICurrentTenant currentTenant, ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager,
IBus rebus, IBus rebus,
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions, IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions,
IOptions<AbpRebusEventBusOptions> abpEventBusRebusOptions, IOptions<AbpRebusEventBusOptions> abpEventBusRebusOptions,
IEventErrorHandler errorHandler) : IEventErrorHandler errorHandler) :
base(serviceScopeFactory, currentTenant, errorHandler) base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler)
{ {
Rebus = rebus; Rebus = rebus;
AbpRebusEventBusOptions = abpEventBusRebusOptions.Value; AbpRebusEventBusOptions = abpEventBusRebusOptions.Value;
@ -125,11 +127,16 @@ namespace Volo.Abp.EventBus.Rebus
return Subscribe(typeof(TEvent), handler); return Subscribe(typeof(TEvent), handler);
} }
public override async Task PublishAsync(Type eventType, object eventData) protected override async Task PublishToEventBusAsync(Type eventType, object eventData)
{ {
await AbpRebusEventBusOptions.Publish(Rebus, eventType, eventData); await AbpRebusEventBusOptions.Publish(Rebus, eventType, eventData);
} }
protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord)
{
unitOfWork.AddOrReplaceDistributedEvent(eventRecord);
}
private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType) private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType)
{ {
return HandlerFactories.GetOrAdd( return HandlerFactories.GetOrAdd(

8
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/LocalDistributedEventBus.cs

@ -122,15 +122,15 @@ namespace Volo.Abp.EventBus.Distributed
_localEventBus.UnsubscribeAll(eventType); _localEventBus.UnsubscribeAll(eventType);
} }
public Task PublishAsync<TEvent>(TEvent eventData) public Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true)
where TEvent : class where TEvent : class
{ {
return _localEventBus.PublishAsync(eventData); return _localEventBus.PublishAsync(eventData, onUnitOfWorkComplete);
} }
public Task PublishAsync(Type eventType, object eventData) public Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true)
{ {
return _localEventBus.PublishAsync(eventType, eventData); return _localEventBus.PublishAsync(eventType, eventData, onUnitOfWorkComplete);
} }
} }
} }

4
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/NullDistributedEventBus.cs

@ -77,12 +77,12 @@ namespace Volo.Abp.EventBus.Distributed
} }
public Task PublishAsync<TEvent>(TEvent eventData) where TEvent : class public Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true) where TEvent : class
{ {
return Task.CompletedTask; return Task.CompletedTask;
} }
public Task PublishAsync(Type eventType, object eventData) public Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true)
{ {
return Task.CompletedTask; return Task.CompletedTask;
} }

29
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs

@ -10,6 +10,7 @@ using Volo.Abp.Collections;
using Volo.Abp.EventBus.Distributed; using Volo.Abp.EventBus.Distributed;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.Reflection; using Volo.Abp.Reflection;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus namespace Volo.Abp.EventBus
{ {
@ -19,15 +20,19 @@ namespace Volo.Abp.EventBus
protected ICurrentTenant CurrentTenant { get; } protected ICurrentTenant CurrentTenant { get; }
protected IUnitOfWorkManager UnitOfWorkManager { get; }
protected IEventErrorHandler ErrorHandler { get; } protected IEventErrorHandler ErrorHandler { get; }
protected EventBusBase( protected EventBusBase(
IServiceScopeFactory serviceScopeFactory, IServiceScopeFactory serviceScopeFactory,
ICurrentTenant currentTenant, ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager,
IEventErrorHandler errorHandler) IEventErrorHandler errorHandler)
{ {
ServiceScopeFactory = serviceScopeFactory; ServiceScopeFactory = serviceScopeFactory;
CurrentTenant = currentTenant; CurrentTenant = currentTenant;
UnitOfWorkManager = unitOfWorkManager;
ErrorHandler = errorHandler; ErrorHandler = errorHandler;
} }
@ -87,13 +92,29 @@ namespace Volo.Abp.EventBus
public abstract void UnsubscribeAll(Type eventType); public abstract void UnsubscribeAll(Type eventType);
/// <inheritdoc/> /// <inheritdoc/>
public virtual Task PublishAsync<TEvent>(TEvent eventData) where TEvent : class public Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true) where TEvent : class
{ {
return PublishAsync(typeof(TEvent), eventData); return PublishAsync(typeof(TEvent), eventData, onUnitOfWorkComplete);
} }
/// <inheritdoc/> /// <inheritdoc/>
public abstract Task PublishAsync(Type eventType, object eventData); public async Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true)
{
if (onUnitOfWorkComplete && UnitOfWorkManager.Current != null)
{
AddToUnitOfWork(
UnitOfWorkManager.Current,
new UnitOfWorkEventRecord(eventType, eventData, EventOrderGenerator.GetNext())
);
return;
}
await PublishToEventBusAsync(eventType, eventData);
}
protected abstract Task PublishToEventBusAsync(Type eventType, object eventData);
protected abstract void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord);
public virtual async Task TriggerHandlersAsync(Type eventType, object eventData, Action<EventExecutionErrorContext> onErrorAction = null) public virtual async Task TriggerHandlersAsync(Type eventType, object eventData, Action<EventExecutionErrorContext> onErrorAction = null)
{ {
@ -133,7 +154,7 @@ namespace Volo.Abp.EventBus
var baseEventType = eventType.GetGenericTypeDefinition().MakeGenericType(baseArg); var baseEventType = eventType.GetGenericTypeDefinition().MakeGenericType(baseArg);
var constructorArgs = ((IEventDataWithInheritableGenericArgument)eventData).GetConstructorArgs(); var constructorArgs = ((IEventDataWithInheritableGenericArgument)eventData).GetConstructorArgs();
var baseEventData = Activator.CreateInstance(baseEventType, constructorArgs); var baseEventData = Activator.CreateInstance(baseEventType, constructorArgs);
await PublishAsync(baseEventType, baseEventData); await PublishToEventBusAsync(baseEventType, baseEventData);
} }
} }
} }

6
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventBus.cs

@ -10,8 +10,9 @@ namespace Volo.Abp.EventBus
/// </summary> /// </summary>
/// <typeparam name="TEvent">Event type</typeparam> /// <typeparam name="TEvent">Event type</typeparam>
/// <param name="eventData">Related data for the event</param> /// <param name="eventData">Related data for the event</param>
/// <param name="onUnitOfWorkComplete">True, to publish the event at the end of the current unit of work, if available</param>
/// <returns>The task to handle async operation</returns> /// <returns>The task to handle async operation</returns>
Task PublishAsync<TEvent>(TEvent eventData) Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true)
where TEvent : class; where TEvent : class;
/// <summary> /// <summary>
@ -19,8 +20,9 @@ namespace Volo.Abp.EventBus
/// </summary> /// </summary>
/// <param name="eventType">Event type</param> /// <param name="eventType">Event type</param>
/// <param name="eventData">Related data for the event</param> /// <param name="eventData">Related data for the event</param>
/// <param name="onUnitOfWorkComplete">True, to publish the event at the end of the current unit of work, if available</param>
/// <returns>The task to handle async operation</returns> /// <returns>The task to handle async operation</returns>
Task PublishAsync(Type eventType, object eventData); Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true);
/// <summary> /// <summary>
/// Registers to an event. /// Registers to an event.

11
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs

@ -12,6 +12,7 @@ using Volo.Abp.DependencyInjection;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.Threading; using Volo.Abp.Threading;
using Volo.Abp.Json; using Volo.Abp.Json;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus.Local namespace Volo.Abp.EventBus.Local
{ {
@ -34,8 +35,9 @@ namespace Volo.Abp.EventBus.Local
IOptions<AbpLocalEventBusOptions> options, IOptions<AbpLocalEventBusOptions> options,
IServiceScopeFactory serviceScopeFactory, IServiceScopeFactory serviceScopeFactory,
ICurrentTenant currentTenant, ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager,
IEventErrorHandler errorHandler) IEventErrorHandler errorHandler)
: base(serviceScopeFactory, currentTenant, errorHandler) : base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler)
{ {
Options = options.Value; Options = options.Value;
Logger = NullLogger<LocalEventBus>.Instance; Logger = NullLogger<LocalEventBus>.Instance;
@ -120,11 +122,16 @@ namespace Volo.Abp.EventBus.Local
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear()); GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear());
} }
public override async Task PublishAsync(Type eventType, object eventData) protected override async Task PublishToEventBusAsync(Type eventType, object eventData)
{ {
await PublishAsync(new LocalEventMessage(Guid.NewGuid(), eventData, eventType)); await PublishAsync(new LocalEventMessage(Guid.NewGuid(), eventData, eventType));
} }
protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord)
{
unitOfWork.AddOrReplaceLocalEvent(eventRecord);
}
public virtual async Task PublishAsync(LocalEventMessage localEventMessage) public virtual async Task PublishAsync(LocalEventMessage localEventMessage)
{ {
await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData, errorContext => await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData, errorContext =>

4
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/NullLocalEventBus.cs

@ -77,12 +77,12 @@ namespace Volo.Abp.EventBus.Local
} }
public Task PublishAsync<TEvent>(TEvent eventData) where TEvent : class public Task PublishAsync<TEvent>(TEvent eventData, bool onUnitOfWorkComplete = true) where TEvent : class
{ {
return Task.CompletedTask; return Task.CompletedTask;
} }
public Task PublishAsync(Type eventType, object eventData) public Task PublishAsync(Type eventType, object eventData, bool onUnitOfWorkComplete = true)
{ {
return Task.CompletedTask; return Task.CompletedTask;
} }

40
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/UnitOfWorkEventPublisher.cs

@ -0,0 +1,40 @@
using System.Collections.Generic;
using System.Threading.Tasks;
using Volo.Abp.DependencyInjection;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Local;
using Volo.Abp.Uow;
namespace Volo.Abp.EventBus
{
[Dependency(ReplaceServices = true)]
public class UnitOfWorkEventPublisher : IUnitOfWorkEventPublisher, ITransientDependency
{
private readonly ILocalEventBus _localEventBus;
private readonly IDistributedEventBus _distributedEventBus;
public UnitOfWorkEventPublisher(
ILocalEventBus localEventBus,
IDistributedEventBus distributedEventBus)
{
_localEventBus = localEventBus;
_distributedEventBus = distributedEventBus;
}
public async Task PublishLocalEventsAsync(IEnumerable<UnitOfWorkEventRecord> localEvents)
{
foreach (var localEvent in localEvents)
{
await _localEventBus.PublishAsync(localEvent.EventType, localEvent.EventData, onUnitOfWorkComplete: false);
}
}
public async Task PublishDistributedEventsAsync(IEnumerable<UnitOfWorkEventRecord> distributedEvents)
{
foreach (var distributedEvent in distributedEvents)
{
await _distributedEventBus.PublishAsync(distributedEvent.EventType, distributedEvent.EventData, onUnitOfWorkComplete: false);
}
}
}
}

59
framework/src/Volo.Abp.MemoryDb/Volo/Abp/Domain/Repositories/MemoryDb/MemoryDbRepository.cs

@ -12,6 +12,7 @@ using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Local; using Volo.Abp.EventBus.Local;
using Volo.Abp.Guids; using Volo.Abp.Guids;
using Volo.Abp.MemoryDb; using Volo.Abp.MemoryDb;
using Volo.Abp.Uow;
namespace Volo.Abp.Domain.Repositories.MemoryDb namespace Volo.Abp.Domain.Repositories.MemoryDb
{ {
@ -65,7 +66,7 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
return ApplyDataFilters((await GetCollectionAsync()).AsQueryable()); return ApplyDataFilters((await GetCollectionAsync()).AsQueryable());
} }
protected virtual async Task TriggerDomainEventsAsync(object entity) protected virtual void TriggerDomainEvents(object entity)
{ {
var generatesDomainEventsEntity = entity as IGeneratesDomainEvents; var generatesDomainEventsEntity = entity as IGeneratesDomainEvents;
if (generatesDomainEventsEntity == null) if (generatesDomainEventsEntity == null)
@ -78,7 +79,13 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
{ {
foreach (var localEvent in localEvents) foreach (var localEvent in localEvents)
{ {
await LocalEventBus.PublishAsync(localEvent.GetType(), localEvent); UnitOfWorkManager.Current?.AddOrReplaceLocalEvent(
new UnitOfWorkEventRecord(
localEvent.EventData.GetType(),
localEvent.EventData,
localEvent.EventOrder
)
);
} }
generatesDomainEventsEntity.ClearLocalEvents(); generatesDomainEventsEntity.ClearLocalEvents();
@ -89,7 +96,13 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
{ {
foreach (var distributedEvent in distributedEvents) foreach (var distributedEvent in distributedEvents)
{ {
await DistributedEventBus.PublishAsync(distributedEvent.GetType(), distributedEvent); UnitOfWorkManager.Current?.AddOrReplaceDistributedEvent(
new UnitOfWorkEventRecord(
distributedEvent.EventData.GetType(),
distributedEvent.EventData,
distributedEvent.EventOrder
)
);
} }
generatesDomainEventsEntity.ClearDistributedEvents(); generatesDomainEventsEntity.ClearDistributedEvents();
@ -143,37 +156,37 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
AuditPropertySetter.SetDeletionProperties(entity); AuditPropertySetter.SetDeletionProperties(entity);
} }
protected virtual async Task TriggerEntityCreateEvents(TEntity entity) protected virtual void TriggerEntityCreateEvents(TEntity entity)
{ {
await EntityChangeEventHelper.TriggerEntityCreatedEventOnUowCompletedAsync(entity); EntityChangeEventHelper.PublishEntityCreatingEvent(entity);
await EntityChangeEventHelper.TriggerEntityCreatingEventAsync(entity); EntityChangeEventHelper.PublishEntityCreatedEvent(entity);
} }
protected virtual async Task TriggerEntityUpdateEventsAsync(TEntity entity) protected virtual void TriggerEntityUpdateEvents(TEntity entity)
{ {
await EntityChangeEventHelper.TriggerEntityUpdatedEventOnUowCompletedAsync(entity); EntityChangeEventHelper.PublishEntityUpdatingEvent(entity);
await EntityChangeEventHelper.TriggerEntityUpdatingEventAsync(entity); EntityChangeEventHelper.PublishEntityUpdatedEvent(entity);
} }
protected virtual async Task TriggerEntityDeleteEventsAsync(TEntity entity) protected virtual void TriggerEntityDeleteEvents(TEntity entity)
{ {
await EntityChangeEventHelper.TriggerEntityDeletedEventOnUowCompletedAsync(entity); EntityChangeEventHelper.PublishEntityDeletingEvent(entity);
await EntityChangeEventHelper.TriggerEntityDeletingEventAsync(entity); EntityChangeEventHelper.PublishEntityDeletedEvent(entity);
} }
protected virtual async Task ApplyAbpConceptsForAddedEntityAsync(TEntity entity) protected virtual void ApplyAbpConceptsForAddedEntity(TEntity entity)
{ {
CheckAndSetId(entity); CheckAndSetId(entity);
SetCreationAuditProperties(entity); SetCreationAuditProperties(entity);
await TriggerEntityCreateEvents(entity); TriggerEntityCreateEvents(entity);
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
} }
protected virtual async Task ApplyAbpConceptsForDeletedEntityAsync(TEntity entity) protected virtual void ApplyAbpConceptsForDeletedEntity(TEntity entity)
{ {
SetDeletionAuditProperties(entity); SetDeletionAuditProperties(entity);
await TriggerEntityDeleteEventsAsync(entity); TriggerEntityDeleteEvents(entity);
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
} }
public override async Task<TEntity> FindAsync( public override async Task<TEntity> FindAsync(
@ -199,7 +212,7 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
bool autoSave = false, bool autoSave = false,
CancellationToken cancellationToken = default) CancellationToken cancellationToken = default)
{ {
await ApplyAbpConceptsForAddedEntityAsync(entity); ApplyAbpConceptsForAddedEntity(entity);
(await GetCollectionAsync()).Add(entity); (await GetCollectionAsync()).Add(entity);
@ -216,14 +229,14 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
if (entity is ISoftDelete softDeleteEntity && softDeleteEntity.IsDeleted) if (entity is ISoftDelete softDeleteEntity && softDeleteEntity.IsDeleted)
{ {
SetDeletionAuditProperties(entity); SetDeletionAuditProperties(entity);
await TriggerEntityDeleteEventsAsync(entity); TriggerEntityDeleteEvents(entity);
} }
else else
{ {
await TriggerEntityUpdateEventsAsync(entity); TriggerEntityUpdateEvents(entity);
} }
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
(await GetCollectionAsync()).Update(entity); (await GetCollectionAsync()).Update(entity);
@ -235,7 +248,7 @@ namespace Volo.Abp.Domain.Repositories.MemoryDb
bool autoSave = false, bool autoSave = false,
CancellationToken cancellationToken = default) CancellationToken cancellationToken = default)
{ {
await ApplyAbpConceptsForDeletedEntityAsync(entity); ApplyAbpConceptsForDeletedEntity(entity);
if (entity is ISoftDelete softDeleteEntity && !IsHardDeleted(entity)) if (entity is ISoftDelete softDeleteEntity && !IsHardDeleted(entity))
{ {

65
framework/src/Volo.Abp.MongoDB/Volo/Abp/Domain/Repositories/MongoDB/MongoDbRepository.cs

@ -17,6 +17,7 @@ using Volo.Abp.EventBus.Local;
using Volo.Abp.Guids; using Volo.Abp.Guids;
using Volo.Abp.MongoDB; using Volo.Abp.MongoDB;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.Uow;
namespace Volo.Abp.Domain.Repositories.MongoDB namespace Volo.Abp.Domain.Repositories.MongoDB
{ {
@ -182,14 +183,14 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
if (entity is ISoftDelete softDeleteEntity && softDeleteEntity.IsDeleted) if (entity is ISoftDelete softDeleteEntity && softDeleteEntity.IsDeleted)
{ {
SetDeletionAuditProperties(entity); SetDeletionAuditProperties(entity);
await TriggerEntityDeleteEventsAsync(entity); TriggerEntityDeleteEvents(entity);
} }
else else
{ {
await TriggerEntityUpdateEventsAsync(entity); TriggerEntityUpdateEvents(entity);
} }
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
var oldConcurrencyStamp = SetNewConcurrencyStamp(entity); var oldConcurrencyStamp = SetNewConcurrencyStamp(entity);
ReplaceOneResult result; ReplaceOneResult result;
@ -235,14 +236,14 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
if (isSoftDeleteEntity) if (isSoftDeleteEntity)
{ {
SetDeletionAuditProperties(entity); SetDeletionAuditProperties(entity);
await TriggerEntityDeleteEventsAsync(entity); TriggerEntityDeleteEvents(entity);
} }
else else
{ {
await TriggerEntityUpdateEventsAsync(entity); TriggerEntityUpdateEvents(entity);
} }
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
SetNewConcurrencyStamp(entity); SetNewConcurrencyStamp(entity);
} }
@ -295,7 +296,7 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
if (typeof(ISoftDelete).IsAssignableFrom(typeof(TEntity)) && !IsHardDeleted(entity)) if (typeof(ISoftDelete).IsAssignableFrom(typeof(TEntity)) && !IsHardDeleted(entity))
{ {
((ISoftDelete)entity).IsDeleted = true; ((ISoftDelete)entity).IsDeleted = true;
await ApplyAbpConceptsForDeletedEntityAsync(entity); ApplyAbpConceptsForDeletedEntity(entity);
ReplaceOneResult result; ReplaceOneResult result;
@ -324,7 +325,7 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
} }
else else
{ {
await ApplyAbpConceptsForDeletedEntityAsync(entity); ApplyAbpConceptsForDeletedEntity(entity);
DeleteResult result; DeleteResult result;
@ -374,7 +375,7 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
hardDeletedEntities.Add(entity); hardDeletedEntities.Add(entity);
} }
await ApplyAbpConceptsForDeletedEntityAsync(entity); ApplyAbpConceptsForDeletedEntity(entity);
} }
var dbContext = await GetDbContextAsync(cancellationToken); var dbContext = await GetDbContextAsync(cancellationToken);
@ -573,33 +574,33 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
{ {
CheckAndSetId(entity); CheckAndSetId(entity);
SetCreationAuditProperties(entity); SetCreationAuditProperties(entity);
await TriggerEntityCreateEvents(entity); TriggerEntityCreateEvents(entity);
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
} }
private async Task TriggerEntityCreateEvents(TEntity entity) private void TriggerEntityCreateEvents(TEntity entity)
{ {
await EntityChangeEventHelper.TriggerEntityCreatedEventOnUowCompletedAsync(entity); EntityChangeEventHelper.PublishEntityCreatingEvent(entity);
await EntityChangeEventHelper.TriggerEntityCreatingEventAsync(entity); EntityChangeEventHelper.PublishEntityCreatedEvent(entity);
} }
protected virtual async Task TriggerEntityUpdateEventsAsync(TEntity entity) protected virtual void TriggerEntityUpdateEvents(TEntity entity)
{ {
await EntityChangeEventHelper.TriggerEntityUpdatedEventOnUowCompletedAsync(entity); EntityChangeEventHelper.PublishEntityUpdatingEvent(entity);
await EntityChangeEventHelper.TriggerEntityUpdatingEventAsync(entity); EntityChangeEventHelper.PublishEntityUpdatedEvent(entity);
} }
protected virtual async Task ApplyAbpConceptsForDeletedEntityAsync(TEntity entity) protected virtual void ApplyAbpConceptsForDeletedEntity(TEntity entity)
{ {
SetDeletionAuditProperties(entity); SetDeletionAuditProperties(entity);
await TriggerEntityDeleteEventsAsync(entity); TriggerEntityDeleteEvents(entity);
await TriggerDomainEventsAsync(entity); TriggerDomainEvents(entity);
} }
protected virtual async Task TriggerEntityDeleteEventsAsync(TEntity entity) protected virtual void TriggerEntityDeleteEvents(TEntity entity)
{ {
await EntityChangeEventHelper.TriggerEntityDeletedEventOnUowCompletedAsync(entity); EntityChangeEventHelper.PublishEntityDeletingEvent(entity);
await EntityChangeEventHelper.TriggerEntityDeletingEventAsync(entity); EntityChangeEventHelper.PublishEntityDeletedEvent(entity);
} }
protected virtual void CheckAndSetId(TEntity entity) protected virtual void CheckAndSetId(TEntity entity)
@ -639,7 +640,7 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
AuditPropertySetter.SetDeletionProperties(entity); AuditPropertySetter.SetDeletionProperties(entity);
} }
protected virtual async Task TriggerDomainEventsAsync(object entity) protected virtual void TriggerDomainEvents(object entity)
{ {
var generatesDomainEventsEntity = entity as IGeneratesDomainEvents; var generatesDomainEventsEntity = entity as IGeneratesDomainEvents;
if (generatesDomainEventsEntity == null) if (generatesDomainEventsEntity == null)
@ -652,7 +653,13 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
{ {
foreach (var localEvent in localEvents) foreach (var localEvent in localEvents)
{ {
await LocalEventBus.PublishAsync(localEvent.GetType(), localEvent); UnitOfWorkManager.Current?.AddOrReplaceLocalEvent(
new UnitOfWorkEventRecord(
localEvent.EventData.GetType(),
localEvent.EventData,
localEvent.EventOrder
)
);
} }
generatesDomainEventsEntity.ClearLocalEvents(); generatesDomainEventsEntity.ClearLocalEvents();
@ -663,7 +670,13 @@ namespace Volo.Abp.Domain.Repositories.MongoDB
{ {
foreach (var distributedEvent in distributedEvents) foreach (var distributedEvent in distributedEvents)
{ {
await DistributedEventBus.PublishAsync(distributedEvent.GetType(), distributedEvent); UnitOfWorkManager.Current?.AddOrReplaceDistributedEvent(
new UnitOfWorkEventRecord(
distributedEvent.EventData.GetType(),
distributedEvent.EventData,
distributedEvent.EventOrder
)
);
} }
generatesDomainEventsEntity.ClearDistributedEvents(); generatesDomainEventsEntity.ClearDistributedEvents();

13
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/AmbientUnitOfWork.cs

@ -19,5 +19,18 @@ namespace Volo.Abp.Uow
{ {
_currentUow.Value = unitOfWork; _currentUow.Value = unitOfWork;
} }
public IUnitOfWork GetCurrentByChecking()
{
var uow = UnitOfWork;
//Skip reserved unit of work
while (uow != null && (uow.IsReserved || uow.IsDisposed || uow.IsCompleted))
{
uow = uow.Outer;
}
return uow;
}
} }
} }

14
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/ChildUnitOfWork.cs

@ -76,6 +76,20 @@ namespace Volo.Abp.Uow
_parent.OnCompleted(handler); _parent.OnCompleted(handler);
} }
public void AddOrReplaceLocalEvent(
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null)
{
_parent.AddOrReplaceLocalEvent(eventRecord, replacementSelector);
}
public void AddOrReplaceDistributedEvent(
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null)
{
_parent.AddOrReplaceDistributedEvent(eventRecord, replacementSelector);
}
public IDatabaseApi FindDatabaseApi(string key) public IDatabaseApi FindDatabaseApi(string key)
{ {
return _parent.FindDatabaseApi(key); return _parent.FindDatabaseApi(key);

14
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/EventOrderGenerator.cs

@ -0,0 +1,14 @@
using System.Threading;
namespace Volo.Abp.Uow
{
public static class EventOrderGenerator
{
private static long _lastOrder;
public static long GetNext()
{
return Interlocked.Increment(ref _lastOrder);
}
}
}

2
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/IAmbientUnitOfWork.cs

@ -2,6 +2,6 @@
{ {
public interface IAmbientUnitOfWork : IUnitOfWorkAccessor public interface IAmbientUnitOfWork : IUnitOfWorkAccessor
{ {
IUnitOfWork GetCurrentByChecking();
} }
} }

10
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/IUnitOfWork.cs

@ -42,5 +42,15 @@ namespace Volo.Abp.Uow
Task RollbackAsync(CancellationToken cancellationToken = default); Task RollbackAsync(CancellationToken cancellationToken = default);
void OnCompleted(Func<Task> handler); void OnCompleted(Func<Task> handler);
void AddOrReplaceLocalEvent(
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null
);
void AddOrReplaceDistributedEvent(
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null
);
} }
} }

12
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/IUnitOfWorkEventPublisher.cs

@ -0,0 +1,12 @@
using System.Collections.Generic;
using System.Threading.Tasks;
namespace Volo.Abp.Uow
{
public interface IUnitOfWorkEventPublisher
{
Task PublishLocalEventsAsync(IEnumerable<UnitOfWorkEventRecord> localEvents);
Task PublishDistributedEventsAsync(IEnumerable<UnitOfWorkEventRecord> distributedEvents);
}
}

19
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/NullUnitOfWorkEventPublisher.cs

@ -0,0 +1,19 @@
using System.Collections.Generic;
using System.Threading.Tasks;
using Volo.Abp.DependencyInjection;
namespace Volo.Abp.Uow
{
public class NullUnitOfWorkEventPublisher : IUnitOfWorkEventPublisher, ISingletonDependency
{
public Task PublishLocalEventsAsync(IEnumerable<UnitOfWorkEventRecord> localEvents)
{
return Task.CompletedTask;
}
public Task PublishDistributedEventsAsync(IEnumerable<UnitOfWorkEventRecord> distributedEvents)
{
return Task.CompletedTask;
}
}
}

89
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWork.cs

@ -1,6 +1,7 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Collections.Immutable; using System.Collections.Immutable;
using System.Linq;
using System.Threading; using System.Threading;
using System.Threading.Tasks; using System.Threading.Tasks;
using JetBrains.Annotations; using JetBrains.Annotations;
@ -33,11 +34,14 @@ namespace Volo.Abp.Uow
public string ReservationName { get; set; } public string ReservationName { get; set; }
protected List<Func<Task>> CompletedHandlers { get; } = new List<Func<Task>>(); protected List<Func<Task>> CompletedHandlers { get; } = new List<Func<Task>>();
protected List<UnitOfWorkEventRecord> DistributedEvents { get; } = new List<UnitOfWorkEventRecord>();
protected List<UnitOfWorkEventRecord> LocalEvents { get; } = new List<UnitOfWorkEventRecord>();
public event EventHandler<UnitOfWorkFailedEventArgs> Failed; public event EventHandler<UnitOfWorkFailedEventArgs> Failed;
public event EventHandler<UnitOfWorkEventArgs> Disposed; public event EventHandler<UnitOfWorkEventArgs> Disposed;
public IServiceProvider ServiceProvider { get; } public IServiceProvider ServiceProvider { get; }
protected IUnitOfWorkEventPublisher UnitOfWorkEventPublisher { get; }
[NotNull] [NotNull]
public Dictionary<string, object> Items { get; } public Dictionary<string, object> Items { get; }
@ -50,9 +54,13 @@ namespace Volo.Abp.Uow
private bool _isCompleting; private bool _isCompleting;
private bool _isRolledback; private bool _isRolledback;
public UnitOfWork(IServiceProvider serviceProvider, IOptions<AbpUnitOfWorkDefaultOptions> options) public UnitOfWork(
IServiceProvider serviceProvider,
IUnitOfWorkEventPublisher unitOfWorkEventPublisher,
IOptions<AbpUnitOfWorkDefaultOptions> options)
{ {
ServiceProvider = serviceProvider; ServiceProvider = serviceProvider;
UnitOfWorkEventPublisher = unitOfWorkEventPublisher;
_defaultOptions = options.Value; _defaultOptions = options.Value;
_databaseApis = new Dictionary<string, IDatabaseApi>(); _databaseApis = new Dictionary<string, IDatabaseApi>();
@ -103,12 +111,12 @@ namespace Volo.Abp.Uow
} }
} }
public IReadOnlyList<IDatabaseApi> GetAllActiveDatabaseApis() public virtual IReadOnlyList<IDatabaseApi> GetAllActiveDatabaseApis()
{ {
return _databaseApis.Values.ToImmutableList(); return _databaseApis.Values.ToImmutableList();
} }
public IReadOnlyList<ITransactionApi> GetAllActiveTransactionApis() public virtual IReadOnlyList<ITransactionApi> GetAllActiveTransactionApis()
{ {
return _transactionApis.Values.ToImmutableList(); return _transactionApis.Values.ToImmutableList();
} }
@ -126,6 +134,30 @@ namespace Volo.Abp.Uow
{ {
_isCompleting = true; _isCompleting = true;
await SaveChangesAsync(cancellationToken); await SaveChangesAsync(cancellationToken);
while (LocalEvents.Any() || DistributedEvents.Any())
{
if (LocalEvents.Any())
{
var localEventsToBePublished = LocalEvents.OrderBy(e => e.EventOrder).ToArray();
LocalEvents.Clear();
await UnitOfWorkEventPublisher.PublishLocalEventsAsync(
localEventsToBePublished
);
}
if (DistributedEvents.Any())
{
var distributedEventsToBePublished = DistributedEvents.OrderBy(e => e.EventOrder).ToArray();
DistributedEvents.Clear();
await UnitOfWorkEventPublisher.PublishDistributedEventsAsync(
distributedEventsToBePublished
);
}
await SaveChangesAsync(cancellationToken);
}
await CommitTransactionsAsync(); await CommitTransactionsAsync();
IsCompleted = true; IsCompleted = true;
await OnCompletedAsync(); await OnCompletedAsync();
@ -149,12 +181,12 @@ namespace Volo.Abp.Uow
await RollbackAllAsync(cancellationToken); await RollbackAllAsync(cancellationToken);
} }
public IDatabaseApi FindDatabaseApi(string key) public virtual IDatabaseApi FindDatabaseApi(string key)
{ {
return _databaseApis.GetOrDefault(key); return _databaseApis.GetOrDefault(key);
} }
public void AddDatabaseApi(string key, IDatabaseApi api) public virtual void AddDatabaseApi(string key, IDatabaseApi api)
{ {
Check.NotNull(key, nameof(key)); Check.NotNull(key, nameof(key));
Check.NotNull(api, nameof(api)); Check.NotNull(api, nameof(api));
@ -167,7 +199,7 @@ namespace Volo.Abp.Uow
_databaseApis.Add(key, api); _databaseApis.Add(key, api);
} }
public IDatabaseApi GetOrAddDatabaseApi(string key, Func<IDatabaseApi> factory) public virtual IDatabaseApi GetOrAddDatabaseApi(string key, Func<IDatabaseApi> factory)
{ {
Check.NotNull(key, nameof(key)); Check.NotNull(key, nameof(key));
Check.NotNull(factory, nameof(factory)); Check.NotNull(factory, nameof(factory));
@ -175,14 +207,14 @@ namespace Volo.Abp.Uow
return _databaseApis.GetOrAdd(key, factory); return _databaseApis.GetOrAdd(key, factory);
} }
public ITransactionApi FindTransactionApi(string key) public virtual ITransactionApi FindTransactionApi(string key)
{ {
Check.NotNull(key, nameof(key)); Check.NotNull(key, nameof(key));
return _transactionApis.GetOrDefault(key); return _transactionApis.GetOrDefault(key);
} }
public void AddTransactionApi(string key, ITransactionApi api) public virtual void AddTransactionApi(string key, ITransactionApi api)
{ {
Check.NotNull(key, nameof(key)); Check.NotNull(key, nameof(key));
Check.NotNull(api, nameof(api)); Check.NotNull(api, nameof(api));
@ -195,7 +227,7 @@ namespace Volo.Abp.Uow
_transactionApis.Add(key, api); _transactionApis.Add(key, api);
} }
public ITransactionApi GetOrAddTransactionApi(string key, Func<ITransactionApi> factory) public virtual ITransactionApi GetOrAddTransactionApi(string key, Func<ITransactionApi> factory)
{ {
Check.NotNull(key, nameof(key)); Check.NotNull(key, nameof(key));
Check.NotNull(factory, nameof(factory)); Check.NotNull(factory, nameof(factory));
@ -203,11 +235,48 @@ namespace Volo.Abp.Uow
return _transactionApis.GetOrAdd(key, factory); return _transactionApis.GetOrAdd(key, factory);
} }
public void OnCompleted(Func<Task> handler) public virtual void OnCompleted(Func<Task> handler)
{ {
CompletedHandlers.Add(handler); CompletedHandlers.Add(handler);
} }
public virtual void AddOrReplaceLocalEvent(
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null)
{
AddOrReplaceEvent(LocalEvents, eventRecord, replacementSelector);
}
public virtual void AddOrReplaceDistributedEvent(
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null)
{
AddOrReplaceEvent(DistributedEvents, eventRecord, replacementSelector);
}
public virtual void AddOrReplaceEvent(
List<UnitOfWorkEventRecord> eventRecords,
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord> replacementSelector = null)
{
if (replacementSelector == null)
{
eventRecords.Add(eventRecord);
}
else
{
var foundIndex = eventRecords.FindIndex(replacementSelector);
if (foundIndex < 0)
{
eventRecords.Add(eventRecord);
}
else
{
eventRecords[foundIndex] = eventRecord;
}
}
}
protected virtual async Task OnCompletedAsync() protected virtual async Task OnCompletedAsync()
{ {
foreach (var handler in CompletedHandlers) foreach (var handler in CompletedHandlers)

29
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWorkEventRecord.cs

@ -0,0 +1,29 @@
using System;
using System.Collections.Generic;
namespace Volo.Abp.Uow
{
public class UnitOfWorkEventRecord
{
public object EventData { get; }
public Type EventType { get; }
public long EventOrder { get; }
/// <summary>
/// Extra properties can be used if needed.
/// </summary>
public Dictionary<string, object> Properties { get; } = new Dictionary<string, object>();
public UnitOfWorkEventRecord(
Type eventType,
object eventData,
long eventOrder)
{
EventType = eventType;
EventData = eventData;
EventOrder = eventOrder;
}
}
}

15
framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWorkManager.cs

@ -10,7 +10,7 @@ namespace Volo.Abp.Uow
[Obsolete("This will be removed in next versions.")] [Obsolete("This will be removed in next versions.")]
public static AsyncLocal<bool> DisableObsoleteDbContextCreationWarning { get; } = new AsyncLocal<bool>(); public static AsyncLocal<bool> DisableObsoleteDbContextCreationWarning { get; } = new AsyncLocal<bool>();
public IUnitOfWork Current => GetCurrentUnitOfWork(); public IUnitOfWork Current => _ambientUnitOfWork.GetCurrentByChecking();
private readonly IServiceScopeFactory _serviceScopeFactory; private readonly IServiceScopeFactory _serviceScopeFactory;
private readonly IAmbientUnitOfWork _ambientUnitOfWork; private readonly IAmbientUnitOfWork _ambientUnitOfWork;
@ -86,19 +86,6 @@ namespace Volo.Abp.Uow
return true; return true;
} }
private IUnitOfWork GetCurrentUnitOfWork()
{
var uow = _ambientUnitOfWork.UnitOfWork;
//Skip reserved unit of work
while (uow != null && (uow.IsReserved || uow.IsDisposed || uow.IsCompleted))
{
uow = uow.Outer;
}
return uow;
}
private IUnitOfWork CreateNewUnitOfWork() private IUnitOfWork CreateNewUnitOfWork()
{ {
var scope = _serviceScopeFactory.CreateScope(); var scope = _serviceScopeFactory.CreateScope();

10
framework/test/Volo.Abp.AspNetCore.Mvc.Tests/Volo/Abp/AspNetCore/Mvc/Uow/TestUnitOfWork.cs

@ -13,8 +13,14 @@ namespace Volo.Abp.AspNetCore.Mvc.Uow
{ {
private readonly TestUnitOfWorkConfig _config; private readonly TestUnitOfWorkConfig _config;
public TestUnitOfWork(IServiceProvider serviceProvider, IOptions<AbpUnitOfWorkDefaultOptions> options, TestUnitOfWorkConfig config) public TestUnitOfWork(
: base(serviceProvider, options) IServiceProvider serviceProvider,
IUnitOfWorkEventPublisher unitOfWorkEventPublisher,
IOptions<AbpUnitOfWorkDefaultOptions> options, TestUnitOfWorkConfig config)
: base(
serviceProvider,
unitOfWorkEventPublisher,
options)
{ {
_config = config; _config = config;
} }

10
framework/test/Volo.Abp.MemoryDb.Tests/Volo/Abp/MemoryDb/DomainEvents/DomainEvents_Tests.cs

@ -1,9 +1,15 @@
using Volo.Abp.TestApp.Testing; using System.Threading.Tasks;
using Volo.Abp.TestApp.Testing;
using Xunit;
namespace Volo.Abp.MemoryDb.DomainEvents namespace Volo.Abp.MemoryDb.DomainEvents
{ {
public class DomainEvents_Tests : DomainEvents_Tests<AbpMemoryDbTestModule> public class DomainEvents_Tests : DomainEvents_Tests<AbpMemoryDbTestModule>
{ {
[Fact(Skip = "MemoryDB doesn't support transactions.")]
public override Task Should_Rollback_Uow_If_Event_Handler_Throws_Exception()
{
return base.Should_Rollback_Uow_If_Event_Handler_Throws_Exception();
}
} }
} }

119
framework/test/Volo.Abp.TestApp/Volo/Abp/TestApp/Testing/DomainEvents_Tests.cs

@ -2,11 +2,13 @@
using System.Linq; using System.Linq;
using System.Threading.Tasks; using System.Threading.Tasks;
using Shouldly; using Shouldly;
using Volo.Abp.Domain.Entities.Events;
using Volo.Abp.Domain.Repositories; using Volo.Abp.Domain.Repositories;
using Volo.Abp.EventBus.Distributed; using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Local; using Volo.Abp.EventBus.Local;
using Volo.Abp.Modularity; using Volo.Abp.Modularity;
using Volo.Abp.TestApp.Domain; using Volo.Abp.TestApp.Domain;
using Volo.Abp.Uow;
using Xunit; using Xunit;
namespace Volo.Abp.TestApp.Testing namespace Volo.Abp.TestApp.Testing
@ -25,6 +27,113 @@ namespace Volo.Abp.TestApp.Testing
DistributedEventBus = GetRequiredService<IDistributedEventBus>(); DistributedEventBus = GetRequiredService<IDistributedEventBus>();
} }
[Fact]
public virtual async Task Should_Publish_Events_In_Order()
{
bool testPersonCreateHandled = false;
bool douglesUpdateHandled = false;
bool douglesNameChangeHandled = false;
bool customEventHandled = false;
bool customEvent2Handled = false;
LocalEventBus.Subscribe<EntityCreatedEventData<Person>>(data =>
{
data.Entity.Name.ShouldBe("TestPerson1");
testPersonCreateHandled = true;
douglesUpdateHandled.ShouldBeFalse();
douglesNameChangeHandled.ShouldBeFalse();
customEventHandled.ShouldBeFalse();
customEvent2Handled.ShouldBeFalse();
return Task.CompletedTask;
});
LocalEventBus.Subscribe<MyCustomEventData>(data =>
{
data.Value.ShouldBe("42");
customEventHandled = true;
testPersonCreateHandled.ShouldBeTrue();
douglesUpdateHandled.ShouldBeFalse();
douglesNameChangeHandled.ShouldBeFalse();
customEvent2Handled.ShouldBeFalse();
return Task.CompletedTask;
});
LocalEventBus.Subscribe<PersonNameChangedEvent>(data =>
{
data.OldName.ShouldBe("Douglas");
data.Person.Name.ShouldBe("Douglas-Updated");
douglesNameChangeHandled = true;
testPersonCreateHandled.ShouldBeTrue();
customEventHandled.ShouldBeTrue();
douglesUpdateHandled.ShouldBeFalse();
customEvent2Handled.ShouldBeFalse();
return Task.CompletedTask;
});
LocalEventBus.Subscribe<EntityUpdatedEventData<Person>>(data =>
{
data.Entity.Name.ShouldBe("Douglas-Updated");
douglesUpdateHandled = true;
testPersonCreateHandled.ShouldBeTrue();
customEventHandled.ShouldBeTrue();
douglesNameChangeHandled.ShouldBeTrue();
customEvent2Handled.ShouldBeFalse();
return Task.CompletedTask;
});
LocalEventBus.Subscribe<MyCustomEventData2>(data =>
{
data.Value.ShouldBe("44");
customEvent2Handled = true;
testPersonCreateHandled.ShouldBeTrue();
customEventHandled.ShouldBeTrue();
douglesUpdateHandled.ShouldBeTrue();
douglesNameChangeHandled.ShouldBeTrue();
return Task.CompletedTask;
});
await WithUnitOfWorkAsync(new AbpUnitOfWorkOptions{IsTransactional = true}, async () =>
{
await PersonRepository.InsertAsync(
new Person(Guid.NewGuid(), "TestPerson1", 42)
);
await LocalEventBus.PublishAsync(new MyCustomEventData { Value = "42" });
var douglas = await PersonRepository.GetAsync(TestDataBuilder.UserDouglasId);
douglas.ChangeName("Douglas-Updated");
await PersonRepository.UpdateAsync(douglas);
await LocalEventBus.PublishAsync(new MyCustomEventData2 { Value = "44" });
});
}
[Fact]
public virtual async Task Should_Rollback_Uow_If_Event_Handler_Throws_Exception()
{
(await PersonRepository.FindAsync(x => x.Name == "TestPerson1")).ShouldBeNull();
LocalEventBus.Subscribe<EntityCreatedEventData<Person>>(data =>
{
data.Entity.Name.ShouldBe("TestPerson1");
throw new ApplicationException("Just to rollback the UOW");
});
var exception = await Assert.ThrowsAsync<ApplicationException>(async () =>
{
await WithUnitOfWorkAsync(new AbpUnitOfWorkOptions{IsTransactional = true}, async () =>
{
await PersonRepository.InsertAsync(
new Person(Guid.NewGuid(), "TestPerson1", 42)
);
});
});
exception.Message.ShouldBe("Just to rollback the UOW");
(await PersonRepository.FindAsync(x => x.Name == "TestPerson1")).ShouldBeNull();
}
[Fact] [Fact]
public async Task Should_Trigger_Domain_Events_For_Aggregate_Root() public async Task Should_Trigger_Domain_Events_For_Aggregate_Root()
{ {
@ -63,5 +172,15 @@ namespace Volo.Abp.TestApp.Testing
isLocalEventTriggered.ShouldBeTrue(); isLocalEventTriggered.ShouldBeTrue();
isDistributedEventTriggered.ShouldBeTrue(); isDistributedEventTriggered.ShouldBeTrue();
} }
private class MyCustomEventData
{
public string Value { get; set; }
}
private class MyCustomEventData2
{
public string Value { get; set; }
}
} }
} }

16
framework/test/Volo.Abp.TestApp/Volo/Abp/TestApp/Testing/EntityChangeEvents_Tests.cs

@ -38,7 +38,9 @@ namespace Volo.Abp.TestApp.Testing
using (var uow = GetRequiredService<IUnitOfWorkManager>().Begin()) using (var uow = GetRequiredService<IUnitOfWorkManager>().Begin())
{ {
#pragma warning disable 618
LocalEventBus.Subscribe<EntityCreatingEventData<Person>>(data => LocalEventBus.Subscribe<EntityCreatingEventData<Person>>(data =>
#pragma warning restore 618
{ {
creatingEventTriggered.ShouldBeFalse(); creatingEventTriggered.ShouldBeFalse();
createdEventTriggered.ShouldBeFalse(); createdEventTriggered.ShouldBeFalse();
@ -88,12 +90,16 @@ namespace Volo.Abp.TestApp.Testing
[Fact] [Fact]
public async Task Multiple_Update_Should_Result_With_Single_Updated_Event_In_The_Same_Uow() public async Task Multiple_Update_Should_Result_With_Single_Updated_Event_In_The_Same_Uow()
{ {
var personId = Guid.NewGuid(); var createEventCount = 0;
await PersonRepository.InsertAsync(new Person(personId, Guid.NewGuid().ToString("D"), 42));
var updateEventCount = 0; var updateEventCount = 0;
var updatedAge = 0; var updatedAge = 0;
DistributedEventBus.Subscribe<EntityCreatedEto<PersonEto>>(eto =>
{
createEventCount++;
return Task.CompletedTask;
});
DistributedEventBus.Subscribe<EntityUpdatedEto<PersonEto>>(eto => DistributedEventBus.Subscribe<EntityUpdatedEto<PersonEto>>(eto =>
{ {
updateEventCount++; updateEventCount++;
@ -101,6 +107,9 @@ namespace Volo.Abp.TestApp.Testing
return Task.CompletedTask; return Task.CompletedTask;
}); });
var personId = Guid.NewGuid();
await PersonRepository.InsertAsync(new Person(personId, Guid.NewGuid().ToString("D"), 42));
using (var uow = GetRequiredService<IUnitOfWorkManager>().Begin()) using (var uow = GetRequiredService<IUnitOfWorkManager>().Begin())
{ {
var person = await PersonRepository.GetAsync(personId); var person = await PersonRepository.GetAsync(personId);
@ -120,6 +129,7 @@ namespace Volo.Abp.TestApp.Testing
await uow.CompleteAsync(); await uow.CompleteAsync();
} }
createEventCount.ShouldBe(1);
updateEventCount.ShouldBe(1); updateEventCount.ShouldBe(1);
updatedAge.ShouldBe(45); updatedAge.ShouldBe(45);
} }

1
framework/test/Volo.Abp.TestApp/Volo/Abp/TestApp/Testing/TestAppTestBase.cs

@ -31,7 +31,6 @@ namespace Volo.Abp.TestApp.Testing
using (var uow = uowManager.Begin(options)) using (var uow = uowManager.Begin(options))
{ {
await action(); await action();
await uow.CompleteAsync(); await uow.CompleteAsync();
} }
} }

2
modules/setting-management/src/Volo.Abp.SettingManagement.Domain/Volo/Abp/SettingManagement/SettingCacheItemInvalidator.cs

@ -23,7 +23,7 @@ namespace Volo.Abp.SettingManagement
eventData.Entity.ProviderKey eventData.Entity.ProviderKey
); );
await Cache.RemoveAsync(cacheKey); await Cache.RemoveAsync(cacheKey, considerUow: true);
} }
protected virtual string CalculateCacheKey(string name, string providerName, string providerKey) protected virtual string CalculateCacheKey(string name, string providerName, string providerKey)

29
test/DistEvents/DistDemoApp/DemoService.cs

@ -0,0 +1,29 @@
using System;
using System.Threading.Tasks;
using Volo.Abp.DependencyInjection;
using Volo.Abp.Domain.Repositories;
namespace DistDemoApp
{
public class DemoService : ITransientDependency
{
private readonly IRepository<TodoItem, Guid> _todoItemRepository;
public DemoService(IRepository<TodoItem, Guid> todoItemRepository)
{
_todoItemRepository = todoItemRepository;
}
public async Task CreateTodoItemAsync()
{
var todoItem = await _todoItemRepository.InsertAsync(
new TodoItem
{
Text = "todo item " + DateTime.Now.Ticks
}
);
Console.WriteLine("Created a new todo item: " + todoItem);
}
}
}

34
test/DistEvents/DistDemoApp/DistDemoApp.csproj

@ -0,0 +1,34 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net5.0</TargetFramework>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="5.0.*" />
<PackageReference Include="Serilog.Extensions.Hosting" Version="3.1.0" />
<PackageReference Include="Serilog.Sinks.Async" Version="1.4.0" />
<PackageReference Include="Serilog.Sinks.Console" Version="3.1.1" />
<PackageReference Include="Serilog.Sinks.File" Version="4.1.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\framework\src\Volo.Abp.EntityFrameworkCore.SqlServer\Volo.Abp.EntityFrameworkCore.SqlServer.csproj" />
<ProjectReference Include="..\..\..\framework\src\Volo.Abp.Autofac\Volo.Abp.Autofac.csproj" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="Microsoft.EntityFrameworkCore.Tools" Version="5.0.*">
<IncludeAssets>runtime; build; native; contentfiles; analyzers</IncludeAssets>
<PrivateAssets>compile; contentFiles; build; buildMultitargeting; buildTransitive; analyzers; native</PrivateAssets>
</PackageReference>
</ItemGroup>
<ItemGroup>
<None Update="appsettings.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
</ItemGroup>
</Project>

37
test/DistEvents/DistDemoApp/DistDemoAppModule.cs

@ -0,0 +1,37 @@
using Microsoft.Extensions.DependencyInjection;
using Volo.Abp.Autofac;
using Volo.Abp.Domain.Entities.Events.Distributed;
using Volo.Abp.EntityFrameworkCore;
using Volo.Abp.EntityFrameworkCore.SqlServer;
using Volo.Abp.Modularity;
namespace DistDemoApp
{
[DependsOn(
typeof(AbpEntityFrameworkCoreSqlServerModule),
typeof(AbpAutofacModule)
)]
public class DistDemoAppModule : AbpModule
{
public override void ConfigureServices(ServiceConfigurationContext context)
{
context.Services.AddHostedService<MyProjectNameHostedService>();
context.Services.AddAbpDbContext<TodoDbContext>(options =>
{
options.AddDefaultRepositories();
});
Configure<AbpDbContextOptions>(options =>
{
options.UseSqlServer();
});
Configure<AbpDistributedEntityEventOptions>(options =>
{
options.EtoMappings.Add<TodoItem, TodoItemEto>();
options.AutoEventSelectors.Add<TodoItem>();
});
}
}
}

61
test/DistEvents/DistDemoApp/Migrations/20210825110134_Initial.Designer.cs

@ -0,0 +1,61 @@
// <auto-generated />
using System;
using DistDemoApp;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Metadata;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Volo.Abp.EntityFrameworkCore;
namespace DistDemoApp.Migrations
{
[DbContext(typeof(TodoDbContext))]
[Migration("20210825110134_Initial")]
partial class Initial
{
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasAnnotation("_Abp_DatabaseProvider", EfCoreDatabaseProvider.SqlServer)
.HasAnnotation("Relational:MaxIdentifierLength", 128)
.HasAnnotation("ProductVersion", "5.0.9")
.HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn);
modelBuilder.Entity("DistDemoApp.TodoItem", b =>
{
b.Property<Guid>("Id")
.HasColumnType("uniqueidentifier");
b.Property<string>("ConcurrencyStamp")
.IsConcurrencyToken()
.HasMaxLength(40)
.HasColumnType("nvarchar(40)")
.HasColumnName("ConcurrencyStamp");
b.Property<DateTime>("CreationTime")
.HasColumnType("datetime2")
.HasColumnName("CreationTime");
b.Property<Guid?>("CreatorId")
.HasColumnType("uniqueidentifier")
.HasColumnName("CreatorId");
b.Property<string>("ExtraProperties")
.HasColumnType("nvarchar(max)")
.HasColumnName("ExtraProperties");
b.Property<string>("Text")
.IsRequired()
.HasMaxLength(128)
.HasColumnType("nvarchar(128)");
b.HasKey("Id");
b.ToTable("TodoItems");
});
#pragma warning restore 612, 618
}
}
}

33
test/DistEvents/DistDemoApp/Migrations/20210825110134_Initial.cs

@ -0,0 +1,33 @@
using System;
using Microsoft.EntityFrameworkCore.Migrations;
namespace DistDemoApp.Migrations
{
public partial class Initial : Migration
{
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.CreateTable(
name: "TodoItems",
columns: table => new
{
Id = table.Column<Guid>(type: "uniqueidentifier", nullable: false),
Text = table.Column<string>(type: "nvarchar(128)", maxLength: 128, nullable: false),
ExtraProperties = table.Column<string>(type: "nvarchar(max)", nullable: true),
ConcurrencyStamp = table.Column<string>(type: "nvarchar(40)", maxLength: 40, nullable: true),
CreationTime = table.Column<DateTime>(type: "datetime2", nullable: false),
CreatorId = table.Column<Guid>(type: "uniqueidentifier", nullable: true)
},
constraints: table =>
{
table.PrimaryKey("PK_TodoItems", x => x.Id);
});
}
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropTable(
name: "TodoItems");
}
}
}

95
test/DistEvents/DistDemoApp/Migrations/20210825112717_Added_Summary_Table.Designer.cs

@ -0,0 +1,95 @@
// <auto-generated />
using System;
using DistDemoApp;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Metadata;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Volo.Abp.EntityFrameworkCore;
namespace DistDemoApp.Migrations
{
[DbContext(typeof(TodoDbContext))]
[Migration("20210825112717_Added_Summary_Table")]
partial class Added_Summary_Table
{
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasAnnotation("_Abp_DatabaseProvider", EfCoreDatabaseProvider.SqlServer)
.HasAnnotation("Relational:MaxIdentifierLength", 128)
.HasAnnotation("ProductVersion", "5.0.9")
.HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn);
modelBuilder.Entity("DistDemoApp.TodoItem", b =>
{
b.Property<Guid>("Id")
.HasColumnType("uniqueidentifier");
b.Property<string>("ConcurrencyStamp")
.IsConcurrencyToken()
.HasMaxLength(40)
.HasColumnType("nvarchar(40)")
.HasColumnName("ConcurrencyStamp");
b.Property<DateTime>("CreationTime")
.HasColumnType("datetime2")
.HasColumnName("CreationTime");
b.Property<Guid?>("CreatorId")
.HasColumnType("uniqueidentifier")
.HasColumnName("CreatorId");
b.Property<string>("ExtraProperties")
.HasColumnType("nvarchar(max)")
.HasColumnName("ExtraProperties");
b.Property<string>("Text")
.IsRequired()
.HasMaxLength(128)
.HasColumnType("nvarchar(128)");
b.HasKey("Id");
b.ToTable("TodoItems");
});
modelBuilder.Entity("DistDemoApp.TodoSummary", b =>
{
b.Property<int>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("int")
.HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn);
b.Property<string>("ConcurrencyStamp")
.IsConcurrencyToken()
.HasMaxLength(40)
.HasColumnType("nvarchar(40)")
.HasColumnName("ConcurrencyStamp");
b.Property<byte>("Day")
.HasColumnType("tinyint");
b.Property<string>("ExtraProperties")
.HasColumnType("nvarchar(max)")
.HasColumnName("ExtraProperties");
b.Property<byte>("Month")
.HasColumnType("tinyint");
b.Property<int>("TotalCount")
.HasColumnType("int");
b.Property<int>("Year")
.HasColumnType("int");
b.HasKey("Id");
b.ToTable("TodoSummaries");
});
#pragma warning restore 612, 618
}
}
}

34
test/DistEvents/DistDemoApp/Migrations/20210825112717_Added_Summary_Table.cs

@ -0,0 +1,34 @@
using Microsoft.EntityFrameworkCore.Migrations;
namespace DistDemoApp.Migrations
{
public partial class Added_Summary_Table : Migration
{
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.CreateTable(
name: "TodoSummaries",
columns: table => new
{
Id = table.Column<int>(type: "int", nullable: false)
.Annotation("SqlServer:Identity", "1, 1"),
Year = table.Column<int>(type: "int", nullable: false),
Month = table.Column<byte>(type: "tinyint", nullable: false),
Day = table.Column<byte>(type: "tinyint", nullable: false),
TotalCount = table.Column<int>(type: "int", nullable: false),
ExtraProperties = table.Column<string>(type: "nvarchar(max)", nullable: true),
ConcurrencyStamp = table.Column<string>(type: "nvarchar(40)", maxLength: 40, nullable: true)
},
constraints: table =>
{
table.PrimaryKey("PK_TodoSummaries", x => x.Id);
});
}
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropTable(
name: "TodoSummaries");
}
}
}

93
test/DistEvents/DistDemoApp/Migrations/TodoDbContextModelSnapshot.cs

@ -0,0 +1,93 @@
// <auto-generated />
using System;
using DistDemoApp;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Metadata;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Volo.Abp.EntityFrameworkCore;
namespace DistDemoApp.Migrations
{
[DbContext(typeof(TodoDbContext))]
partial class TodoDbContextModelSnapshot : ModelSnapshot
{
protected override void BuildModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasAnnotation("_Abp_DatabaseProvider", EfCoreDatabaseProvider.SqlServer)
.HasAnnotation("Relational:MaxIdentifierLength", 128)
.HasAnnotation("ProductVersion", "5.0.9")
.HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn);
modelBuilder.Entity("DistDemoApp.TodoItem", b =>
{
b.Property<Guid>("Id")
.HasColumnType("uniqueidentifier");
b.Property<string>("ConcurrencyStamp")
.IsConcurrencyToken()
.HasMaxLength(40)
.HasColumnType("nvarchar(40)")
.HasColumnName("ConcurrencyStamp");
b.Property<DateTime>("CreationTime")
.HasColumnType("datetime2")
.HasColumnName("CreationTime");
b.Property<Guid?>("CreatorId")
.HasColumnType("uniqueidentifier")
.HasColumnName("CreatorId");
b.Property<string>("ExtraProperties")
.HasColumnType("nvarchar(max)")
.HasColumnName("ExtraProperties");
b.Property<string>("Text")
.IsRequired()
.HasMaxLength(128)
.HasColumnType("nvarchar(128)");
b.HasKey("Id");
b.ToTable("TodoItems");
});
modelBuilder.Entity("DistDemoApp.TodoSummary", b =>
{
b.Property<int>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("int")
.HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn);
b.Property<string>("ConcurrencyStamp")
.IsConcurrencyToken()
.HasMaxLength(40)
.HasColumnType("nvarchar(40)")
.HasColumnName("ConcurrencyStamp");
b.Property<byte>("Day")
.HasColumnType("tinyint");
b.Property<string>("ExtraProperties")
.HasColumnType("nvarchar(max)")
.HasColumnName("ExtraProperties");
b.Property<byte>("Month")
.HasColumnType("tinyint");
b.Property<int>("TotalCount")
.HasColumnType("int");
b.Property<int>("Year")
.HasColumnType("int");
b.HasKey("Id");
b.ToTable("TodoSummaries");
});
#pragma warning restore 612, 618
}
}
}

39
test/DistEvents/DistDemoApp/MyProjectNameHostedService.cs

@ -0,0 +1,39 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Hosting;
using Volo.Abp;
namespace DistDemoApp
{
public class MyProjectNameHostedService : IHostedService
{
private readonly IAbpApplicationWithExternalServiceProvider _application;
private readonly IServiceProvider _serviceProvider;
private readonly DemoService _demoService;
public MyProjectNameHostedService(
IAbpApplicationWithExternalServiceProvider application,
IServiceProvider serviceProvider,
DemoService demoService)
{
_application = application;
_serviceProvider = serviceProvider;
_demoService = demoService;
}
public async Task StartAsync(CancellationToken cancellationToken)
{
_application.Initialize(_serviceProvider);
await _demoService.CreateTodoItemAsync();
}
public Task StopAsync(CancellationToken cancellationToken)
{
_application.Shutdown();
return Task.CompletedTask;
}
}
}

57
test/DistEvents/DistDemoApp/Program.cs

@ -0,0 +1,57 @@
using System;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Serilog;
using Serilog.Events;
namespace DistDemoApp
{
public class Program
{
public static async Task<int> Main(string[] args)
{
Log.Logger = new LoggerConfiguration()
#if DEBUG
.MinimumLevel.Debug()
#else
.MinimumLevel.Information()
#endif
.MinimumLevel.Override("Microsoft", LogEventLevel.Warning)
.Enrich.FromLogContext()
.WriteTo.Async(c => c.File("Logs/logs.txt"))
.WriteTo.Async(c => c.Console())
.CreateLogger();
try
{
Log.Information("Starting console host.");
await CreateHostBuilder(args).RunConsoleAsync();
return 0;
}
catch (Exception ex)
{
Log.Fatal(ex, "Host terminated unexpectedly!");
return 1;
}
finally
{
Log.CloseAndFlush();
}
}
internal static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.UseAutofac()
.UseSerilog()
.ConfigureAppConfiguration((context, config) =>
{
//setup your additional configuration sources
})
.ConfigureServices((hostContext, services) =>
{
services.AddApplication<DistDemoAppModule>();
});
}
}

28
test/DistEvents/DistDemoApp/TodoDbContext.cs

@ -0,0 +1,28 @@
using Microsoft.EntityFrameworkCore;
using Volo.Abp.Domain.Entities;
using Volo.Abp.EntityFrameworkCore;
namespace DistDemoApp
{
public class TodoDbContext : AbpDbContext<TodoDbContext>
{
public DbSet<TodoItem> TodoItems { get; set; }
public DbSet<TodoSummary> TodoSummaries { get; set; }
public TodoDbContext(DbContextOptions<TodoDbContext> options)
: base(options)
{
}
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
base.OnModelCreating(modelBuilder);
modelBuilder.Entity<TodoItem>(b =>
{
b.Property(x => x.Text).IsRequired().HasMaxLength(128);
});
}
}
}

29
test/DistEvents/DistDemoApp/TodoDbContextFactory.cs

@ -0,0 +1,29 @@
using System.IO;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Design;
using Microsoft.Extensions.Configuration;
namespace DistDemoApp
{
public class TodoDbContextFactory : IDesignTimeDbContextFactory<TodoDbContext>
{
public TodoDbContext CreateDbContext(string[] args)
{
var configuration = BuildConfiguration();
var builder = new DbContextOptionsBuilder<TodoDbContext>()
.UseSqlServer(configuration.GetConnectionString("Default"));
return new TodoDbContext(builder.Options);
}
private static IConfigurationRoot BuildConfiguration()
{
var builder = new ConfigurationBuilder()
.SetBasePath(Directory.GetCurrentDirectory())
.AddJsonFile("appsettings.json", optional: false);
return builder.Build();
}
}
}

65
test/DistEvents/DistDemoApp/TodoEventHandler.cs

@ -0,0 +1,65 @@
using System;
using System.Threading.Tasks;
using Volo.Abp.DependencyInjection;
using Volo.Abp.Domain.Entities.Events.Distributed;
using Volo.Abp.Domain.Repositories;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Uow;
namespace DistDemoApp
{
public class TodoEventHandler :
IDistributedEventHandler<EntityCreatedEto<TodoItemEto>>,
IDistributedEventHandler<EntityDeletedEto<TodoItemEto>>,
ITransientDependency
{
private readonly IRepository<TodoSummary, int> _todoSummaryRepository;
public TodoEventHandler(IRepository<TodoSummary, int> todoSummaryRepository)
{
_todoSummaryRepository = todoSummaryRepository;
}
[UnitOfWork]
public virtual async Task HandleEventAsync(EntityCreatedEto<TodoItemEto> eventData)
{
var dateTime = eventData.Entity.CreationTime;
var todoSummary = await _todoSummaryRepository.FindAsync(
x => x.Year == dateTime.Year &&
x.Month == dateTime.Month &&
x.Day == dateTime.Day
);
if (todoSummary == null)
{
todoSummary = await _todoSummaryRepository.InsertAsync(new TodoSummary(dateTime));
}
else
{
todoSummary.Increase();
await _todoSummaryRepository.UpdateAsync(todoSummary);
}
Console.WriteLine("Increased total count: " + todoSummary);
throw new ApplicationException("Thrown to rollback the UOW!");
}
public async Task HandleEventAsync(EntityDeletedEto<TodoItemEto> eventData)
{
var dateTime = eventData.Entity.CreationTime;
var todoSummary = await _todoSummaryRepository.FirstOrDefaultAsync(
x => x.Year == dateTime.Year &&
x.Month == dateTime.Month &&
x.Day == dateTime.Day
);
if (todoSummary != null)
{
todoSummary.Decrease();
await _todoSummaryRepository.UpdateAsync(todoSummary);
Console.WriteLine("Decreased total count: " + todoSummary);
}
}
}
}

15
test/DistEvents/DistDemoApp/TodoItem.cs

@ -0,0 +1,15 @@
using System;
using Volo.Abp.Domain.Entities.Auditing;
namespace DistDemoApp
{
public class TodoItem : CreationAuditedAggregateRoot<Guid>
{
public string Text { get; set; }
public override string ToString()
{
return $"{base.ToString()}, Text = {Text}";
}
}
}

12
test/DistEvents/DistDemoApp/TodoItemEto.cs

@ -0,0 +1,12 @@
using System;
using Volo.Abp.EventBus;
namespace DistDemoApp
{
[EventName("todo-item")]
public class TodoItemEto
{
public DateTime CreationTime { get; set; }
public string Text { get; set; }
}
}

24
test/DistEvents/DistDemoApp/TodoItemObjectMapper.cs

@ -0,0 +1,24 @@
using Volo.Abp.DependencyInjection;
using Volo.Abp.ObjectMapping;
namespace DistDemoApp
{
public class TodoItemObjectMapper : IObjectMapper<TodoItem, TodoItemEto>, ISingletonDependency
{
public TodoItemEto Map(TodoItem source)
{
return new TodoItemEto
{
Text = source.Text,
CreationTime = source.CreationTime
};
}
public TodoItemEto Map(TodoItem source, TodoItemEto destination)
{
destination.Text = source.Text;
destination.CreationTime = source.CreationTime;
return destination;
}
}
}

41
test/DistEvents/DistDemoApp/TodoSummary.cs

@ -0,0 +1,41 @@
using System;
using Volo.Abp.Domain.Entities;
namespace DistDemoApp
{
public class TodoSummary : AggregateRoot<int>
{
public int Year { get; private set; }
public byte Month { get; private set; }
public byte Day { get; private set; }
public int TotalCount { get; private set; }
private TodoSummary()
{
}
public TodoSummary(DateTime dateTime, int initialCount = 1)
{
Year = dateTime.Year;
Month = (byte)dateTime.Month;
Day = (byte)dateTime.Day;
TotalCount = initialCount;
}
public void Increase(int amount = 1)
{
TotalCount += amount;
}
public void Decrease(int amount = 1)
{
TotalCount -= amount;
}
public override string ToString()
{
return $"{base.ToString()}, {Year}-{Month:00}-{Day:00}: {TotalCount}";
}
}
}

5
test/DistEvents/DistDemoApp/appsettings.json

@ -0,0 +1,5 @@
{
"ConnectionStrings": {
"Default": "Server=(LocalDb)\\MSSQLLocalDB;Database=DistEventsDemo;Trusted_Connection=True"
}
}

16
test/DistEvents/DistEventsDemo.sln

@ -0,0 +1,16 @@

Microsoft Visual Studio Solution File, Format Version 12.00
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "DistDemoApp", "DistDemoApp\DistDemoApp.csproj", "{10DBC6BC-1269-4C68-9F6C-12209A3FBF5B}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Release|Any CPU = Release|Any CPU
EndGlobalSection
GlobalSection(ProjectConfigurationPlatforms) = postSolution
{10DBC6BC-1269-4C68-9F6C-12209A3FBF5B}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{10DBC6BC-1269-4C68-9F6C-12209A3FBF5B}.Debug|Any CPU.Build.0 = Debug|Any CPU
{10DBC6BC-1269-4C68-9F6C-12209A3FBF5B}.Release|Any CPU.ActiveCfg = Release|Any CPU
{10DBC6BC-1269-4C68-9F6C-12209A3FBF5B}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
EndGlobal
Loading…
Cancel
Save