Browse Source

`AddOrReplaceEvent` when unit of work `CompleteAsync`.

pull/21211/head
maliming 2 years ago
parent
commit
f252f5681f
No known key found for this signature in database GPG Key ID: A646B9CB645ECEA4
  1. 47
      framework/src/Volo.Abp.Uow/Volo/Abp/Uow/UnitOfWork.cs

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

@ -34,7 +34,11 @@ public class UnitOfWork : IUnitOfWork, ITransientDependency
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<KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>> DistributedEventWithPredicates { get; } = new List<KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>>();
protected List<UnitOfWorkEventRecord> DistributedEvents { get; } = new List<UnitOfWorkEventRecord>(); protected List<UnitOfWorkEventRecord> DistributedEvents { get; } = new List<UnitOfWorkEventRecord>();
protected List<KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>> LocalEventWithPredicates { get; } = new List<KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>>();
protected List<UnitOfWorkEventRecord> LocalEvents { get; } = new List<UnitOfWorkEventRecord>(); protected List<UnitOfWorkEventRecord> LocalEvents { get; } = new List<UnitOfWorkEventRecord>();
public event EventHandler<UnitOfWorkFailedEventArgs> Failed = default!; public event EventHandler<UnitOfWorkFailedEventArgs> Failed = default!;
@ -135,24 +139,31 @@ public class UnitOfWork : IUnitOfWork, ITransientDependency
_isCompleting = true; _isCompleting = true;
await SaveChangesAsync(cancellationToken); await SaveChangesAsync(cancellationToken);
DistributedEvents.AddRange(GetEventsRecordsByPredicate(DistributedEventWithPredicates));
LocalEvents.AddRange(GetEventsRecordsByPredicate(LocalEventWithPredicates));
while (LocalEvents.Any() || DistributedEvents.Any()) while (LocalEvents.Any() || DistributedEvents.Any())
{ {
if (LocalEvents.Any()) if (LocalEvents.Any())
{ {
var localEventsToBePublished = LocalEvents.OrderBy(e => e.EventOrder).ToArray(); var localEventsToBePublished = LocalEvents.OrderBy(e => e.EventOrder).ToArray();
LocalEventWithPredicates.Clear();
LocalEvents.Clear(); LocalEvents.Clear();
await UnitOfWorkEventPublisher.PublishLocalEventsAsync( await UnitOfWorkEventPublisher.PublishLocalEventsAsync(
localEventsToBePublished localEventsToBePublished
); );
LocalEvents.AddRange(GetEventsRecordsByPredicate(LocalEventWithPredicates));
} }
if (DistributedEvents.Any()) if (DistributedEvents.Any())
{ {
var distributedEventsToBePublished = DistributedEvents.OrderBy(e => e.EventOrder).ToArray(); var distributedEventsToBePublished = DistributedEvents.OrderBy(e => e.EventOrder).ToArray();
DistributedEventWithPredicates.Clear();
DistributedEvents.Clear(); DistributedEvents.Clear();
await UnitOfWorkEventPublisher.PublishDistributedEventsAsync( await UnitOfWorkEventPublisher.PublishDistributedEventsAsync(
distributedEventsToBePublished distributedEventsToBePublished
); );
DistributedEvents.AddRange(GetEventsRecordsByPredicate(DistributedEventWithPredicates));
} }
await SaveChangesAsync(cancellationToken); await SaveChangesAsync(cancellationToken);
@ -244,38 +255,44 @@ public class UnitOfWork : IUnitOfWork, ITransientDependency
UnitOfWorkEventRecord eventRecord, UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord>? replacementSelector = null) Predicate<UnitOfWorkEventRecord>? replacementSelector = null)
{ {
AddOrReplaceEvent(LocalEvents, eventRecord, replacementSelector); LocalEventWithPredicates.Add(new KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>(eventRecord, replacementSelector));
} }
public virtual void AddOrReplaceDistributedEvent( public virtual void AddOrReplaceDistributedEvent(
UnitOfWorkEventRecord eventRecord, UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord>? replacementSelector = null) Predicate<UnitOfWorkEventRecord>? replacementSelector = null)
{ {
AddOrReplaceEvent(DistributedEvents, eventRecord, replacementSelector); DistributedEventWithPredicates.Add(new KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>(eventRecord, replacementSelector));
} }
public virtual void AddOrReplaceEvent( protected virtual List<UnitOfWorkEventRecord> GetEventsRecordsByPredicate(List<KeyValuePair<UnitOfWorkEventRecord, Predicate<UnitOfWorkEventRecord>?>> eventWithPredicates)
List<UnitOfWorkEventRecord> eventRecords,
UnitOfWorkEventRecord eventRecord,
Predicate<UnitOfWorkEventRecord>? replacementSelector = null)
{ {
if (replacementSelector == null) var eventRecords = new List<UnitOfWorkEventRecord>();
{ foreach (var eventWithPredicate in eventWithPredicates)
eventRecords.Add(eventRecord);
}
else
{ {
var foundIndex = eventRecords.FindIndex(replacementSelector); var eventRecord = eventWithPredicate.Key;
if (foundIndex < 0) var replacementSelector = eventWithPredicate.Value;
if (replacementSelector == null)
{ {
eventRecords.Add(eventRecord); eventRecords.Add(eventRecord);
} }
else else
{ {
eventRecord.SetOrder(eventRecords[foundIndex].EventOrder); var foundIndex = eventRecords.FindIndex(replacementSelector);
eventRecords[foundIndex] = eventRecord; if (foundIndex < 0)
{
eventRecords.Add(eventRecord);
}
else
{
eventRecord.SetOrder(eventRecords[foundIndex].EventOrder);
eventRecords[foundIndex] = eventRecord;
}
} }
} }
return eventRecords;
} }
protected virtual async Task OnCompletedAsync() protected virtual async Task OnCompletedAsync()

Loading…
Cancel
Save