mirror of https://github.com/abpframework/abp.git
16 changed files with 207 additions and 28 deletions
@ -0,0 +1,21 @@ |
|||
using JetBrains.Annotations; |
|||
using Microsoft.EntityFrameworkCore; |
|||
using Volo.Abp.Data; |
|||
using Volo.Abp.EntityFrameworkCore.Modeling; |
|||
|
|||
namespace Volo.Abp.EntityFrameworkCore.DistributedEvents |
|||
{ |
|||
public static class EventOutboxDbContextModelBuilderExtensions |
|||
{ |
|||
public static void ConfigureEventOutbox([NotNull] this ModelBuilder builder) |
|||
{ |
|||
builder.Entity<OutgoingEventRecord>(b => |
|||
{ |
|||
b.ToTable(AbpCommonDbProperties.DbTablePrefix + "EventOutbox", AbpCommonDbProperties.DbSchema); |
|||
b.ConfigureByConvention(); |
|||
b.Property(x => x.EventName).IsRequired().HasMaxLength(OutgoingEventRecord.MaxEventNameLength); |
|||
b.Property(x => x.EventData).IsRequired(); |
|||
}); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,9 @@ |
|||
using Microsoft.EntityFrameworkCore; |
|||
|
|||
namespace Volo.Abp.EntityFrameworkCore.DistributedEvents |
|||
{ |
|||
public interface IHasEventOutbox |
|||
{ |
|||
DbSet<OutgoingEventRecord> OutgoingEventRecords { get; set; } |
|||
} |
|||
} |
|||
@ -0,0 +1,29 @@ |
|||
using System; |
|||
using Volo.Abp.Data; |
|||
using Volo.Abp.Domain.Entities; |
|||
|
|||
namespace Volo.Abp.EntityFrameworkCore.DistributedEvents |
|||
{ |
|||
public class OutgoingEventRecord : BasicAggregateRoot<Guid>, IHasExtraProperties |
|||
{ |
|||
public static int MaxEventNameLength { get; set; } = 256; |
|||
|
|||
public ExtraPropertyDictionary ExtraProperties { get; protected set; } |
|||
|
|||
public string EventName { get; set; } |
|||
public byte[] EventData { get; set; } |
|||
|
|||
protected OutgoingEventRecord() |
|||
{ |
|||
ExtraProperties = new ExtraPropertyDictionary(); |
|||
this.SetDefaultsForExtraProperties(); |
|||
} |
|||
|
|||
public OutgoingEventRecord(Guid id) |
|||
: base(id) |
|||
{ |
|||
ExtraProperties = new ExtraPropertyDictionary(); |
|||
this.SetDefaultsForExtraProperties(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,69 @@ |
|||
using System; |
|||
using System.Threading.Tasks; |
|||
using Microsoft.Extensions.DependencyInjection; |
|||
using Volo.Abp.MultiTenancy; |
|||
using Volo.Abp.Uow; |
|||
|
|||
namespace Volo.Abp.EventBus.Distributed |
|||
{ |
|||
public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventBus |
|||
{ |
|||
protected DistributedEventBusBase( |
|||
IServiceScopeFactory serviceScopeFactory, |
|||
ICurrentTenant currentTenant, |
|||
IUnitOfWorkManager unitOfWorkManager, |
|||
IEventErrorHandler errorHandler |
|||
) : base( |
|||
serviceScopeFactory, |
|||
currentTenant, |
|||
unitOfWorkManager, |
|||
errorHandler) |
|||
{ |
|||
} |
|||
|
|||
public IDisposable Subscribe<TEvent>(IDistributedEventHandler<TEvent> handler) where TEvent : class |
|||
{ |
|||
return Subscribe(typeof(TEvent), handler); |
|||
} |
|||
|
|||
public Task PublishAsync<TEvent>( |
|||
TEvent eventData, |
|||
bool onUnitOfWorkComplete = true, |
|||
bool useOutbox = true) |
|||
where TEvent : class |
|||
{ |
|||
return PublishAsync(typeof(TEvent), eventData, onUnitOfWorkComplete, useOutbox); |
|||
} |
|||
|
|||
public async Task PublishAsync( |
|||
Type eventType, |
|||
object eventData, |
|||
bool onUnitOfWorkComplete = true, |
|||
bool useOutbox = true) |
|||
{ |
|||
if (onUnitOfWorkComplete && UnitOfWorkManager.Current != null) |
|||
{ |
|||
AddToUnitOfWork( |
|||
UnitOfWorkManager.Current, |
|||
new UnitOfWorkEventRecord(eventType, eventData, EventOrderGenerator.GetNext(), useOutbox) |
|||
); |
|||
return; |
|||
} |
|||
|
|||
if (useOutbox) |
|||
{ |
|||
if (await AddToOutboxAsync(eventType, eventData)) |
|||
{ |
|||
return; |
|||
} |
|||
} |
|||
|
|||
await PublishToEventBusAsync(eventType, eventData); |
|||
} |
|||
|
|||
private async Task<bool> AddToOutboxAsync(Type eventType, object eventData) |
|||
{ |
|||
return false; |
|||
} |
|||
} |
|||
} |
|||
Loading…
Reference in new issue