From 8af7ccdbaf499004a9506ac643bcba4566505eb9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Halil=20=C4=B0brahim=20Kalkan?= Date: Thu, 9 Sep 2021 21:26:14 +0300 Subject: [PATCH] Implemented initial inbox processing logic --- .../DistributedEvents/DbContextEventInbox.cs | 21 ++- .../DistributedEvents/DbContextEventOutbox.cs | 1 + .../DistributedEvents/IncomingEventRecord.cs | 10 ++ .../EventBus/Boxes/AbpEventBusBoxesModule.cs | 1 + .../Abp/EventBus/Boxes/IInboxProcessor.cs | 13 ++ .../Volo/Abp/EventBus/Boxes/IOutboxSender.cs | 1 + .../Abp/EventBus/Boxes/InboxProcessManager.cs | 48 ++++++ .../Volo/Abp/EventBus/Boxes/InboxProcessor.cs | 111 +++++++++++++ .../Abp/EventBus/Boxes/OutboxSenderManager.cs | 4 +- .../Kafka/KafkaDistributedEventBus.cs | 17 ++ .../RabbitMq/RabbitMqDistributedEventBus.cs | 21 ++- .../Rebus/RebusDistributedEventBus.cs | 6 + .../Distributed/DistributedEventBusBase.cs | 2 + .../Abp/EventBus/Distributed/IEventInbox.cs | 5 +- .../Distributed/IRawEventPublisher.cs | 7 +- .../Volo/Abp/EventBus/EventBusBase.cs | 13 ++ ...51_Added_Inbox_Process_Columns.Designer.cs | 155 ++++++++++++++++++ ...10909182251_Added_Inbox_Process_Columns.cs | 35 ++++ .../Migrations/TodoDbContextModelSnapshot.cs | 6 + 19 files changed, 470 insertions(+), 7 deletions(-) create mode 100644 framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IInboxProcessor.cs create mode 100644 framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessManager.cs create mode 100644 framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs create mode 100644 test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.Designer.cs create mode 100644 test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.cs diff --git a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs index aeaeebe948..df45d50ae5 100644 --- a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs +++ b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs @@ -1,8 +1,10 @@ -using System.Collections.Generic; +using System; +using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using Microsoft.EntityFrameworkCore; using Volo.Abp.EventBus.Distributed; +using Volo.Abp.Timing; using Volo.Abp.Uow; namespace Volo.Abp.EntityFrameworkCore.DistributedEvents @@ -11,11 +13,14 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents where TDbContext : IHasEventInbox { protected IDbContextProvider DbContextProvider { get; } + protected IClock Clock { get; } public DbContextEventInbox( - IDbContextProvider dbContextProvider) + IDbContextProvider dbContextProvider, + IClock clock) { DbContextProvider = dbContextProvider; + Clock = clock; } [UnitOfWork] @@ -35,6 +40,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents var outgoingEventRecords = await dbContext .IncomingEvents .AsNoTracking() + .Where(x => !x.Processed) .OrderBy(x => x.CreationTime) .Take(maxCount) .ToListAsync(); @@ -43,5 +49,16 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents .Select(x => x.ToIncomingEventInfo()) .ToList(); } + + public async Task MarkAsProcessedAsync(Guid id) + { + //TODO: Optimize? + var dbContext = (IHasEventInbox) await DbContextProvider.GetDbContextAsync(); + var incomingEvent = await dbContext.IncomingEvents.FindAsync(id); + if (incomingEvent != null) + { + incomingEvent.MarkAsProcessed(Clock.Now); + } + } } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventOutbox.cs b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventOutbox.cs index 89e10bbd75..a785ad7ce5 100644 --- a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventOutbox.cs +++ b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventOutbox.cs @@ -48,6 +48,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents [UnitOfWork] public virtual async Task DeleteAsync(Guid id) { + //TODO: Optimize? var dbContext = (IHasEventOutbox) await DbContextProvider.GetDbContextAsync(); var outgoingEvent = await dbContext.OutgoingEvents.FindAsync(id); if (outgoingEvent != null) diff --git a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/IncomingEventRecord.cs b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/IncomingEventRecord.cs index 95bc5fa9d4..62c2781b48 100644 --- a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/IncomingEventRecord.cs +++ b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/IncomingEventRecord.cs @@ -20,6 +20,10 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents public byte[] EventData { get; private set; } public DateTime CreationTime { get; private set; } + + public bool Processed { get; set; } + + public DateTime? ProcessedTime { get; set; } protected IncomingEventRecord() { @@ -48,5 +52,11 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents CreationTime ); } + + public void MarkAsProcessed(DateTime processedTime) + { + Processed = true; + ProcessedTime = processedTime; + } } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs index 9445f1de08..8e5eeb2c60 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs @@ -12,6 +12,7 @@ namespace Volo.Abp.EventBus.Boxes public override void OnApplicationInitialization(ApplicationInitializationContext context) { context.AddBackgroundWorker(); + context.AddBackgroundWorker(); } } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IInboxProcessor.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IInboxProcessor.cs new file mode 100644 index 0000000000..e93ff0cc2e --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IInboxProcessor.cs @@ -0,0 +1,13 @@ +using System.Threading; +using System.Threading.Tasks; +using Volo.Abp.EventBus.Distributed; + +namespace Volo.Abp.EventBus.Boxes +{ + public interface IInboxProcessor + { + Task StartAsync(InboxConfig inboxConfig, CancellationToken cancellationToken = default); + + Task StopAsync(CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IOutboxSender.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IOutboxSender.cs index b9a78060df..4a700eb823 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IOutboxSender.cs +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IOutboxSender.cs @@ -7,6 +7,7 @@ namespace Volo.Abp.EventBus.Boxes public interface IOutboxSender { Task StartAsync(OutboxConfig outboxConfig, CancellationToken cancellationToken = default); + Task StopAsync(CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessManager.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessManager.cs new file mode 100644 index 0000000000..310cad2886 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessManager.cs @@ -0,0 +1,48 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Volo.Abp.BackgroundWorkers; +using Volo.Abp.EventBus.Distributed; + +namespace Volo.Abp.EventBus.Boxes +{ + public class InboxProcessManager : IBackgroundWorker + { + protected AbpDistributedEventBusOptions Options { get; } + protected IServiceProvider ServiceProvider { get; } + protected List Processors { get; } + + public InboxProcessManager( + IOptions options, + IServiceProvider serviceProvider) + { + ServiceProvider = serviceProvider; + Options = options.Value; + Processors = new List(); + } + + public async Task StartAsync(CancellationToken cancellationToken = default) + { + foreach (var inboxConfig in Options.Inboxes.Values) + { + if (inboxConfig.IsProcessingEnabled) + { + var processor = ServiceProvider.GetRequiredService(); + await processor.StartAsync(inboxConfig, cancellationToken); + Processors.Add(processor); + } + } + } + + public async Task StopAsync(CancellationToken cancellationToken = default) + { + foreach (var processor in Processors) + { + await processor.StopAsync(cancellationToken); + } + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs new file mode 100644 index 0000000000..4725e85d51 --- /dev/null +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs @@ -0,0 +1,111 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Medallion.Threading; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Volo.Abp.DependencyInjection; +using Volo.Abp.EventBus.Distributed; +using Volo.Abp.Threading; +using Volo.Abp.Uow; + +namespace Volo.Abp.EventBus.Boxes +{ + public class InboxProcessor : IInboxProcessor, ITransientDependency + { + protected IServiceProvider ServiceProvider { get; } + protected AbpTimer Timer { get; } + protected IDistributedEventBus DistributedEventBus { get; } + protected IDistributedLockProvider DistributedLockProvider { get; } + protected IUnitOfWorkManager UnitOfWorkManager { get; } + protected IEventInbox Inbox { get; private set; } + protected InboxConfig InboxConfig { get; private set; } + + protected string DistributedLockName => "Inbox_" + InboxConfig.Name; + public ILogger Logger { get; set; } + + public InboxProcessor( + IServiceProvider serviceProvider, + AbpTimer timer, + IDistributedEventBus distributedEventBus, + IDistributedLockProvider distributedLockProvider, + IUnitOfWorkManager unitOfWorkManager) + { + ServiceProvider = serviceProvider; + Timer = timer; + DistributedEventBus = distributedEventBus; + DistributedLockProvider = distributedLockProvider; + UnitOfWorkManager = unitOfWorkManager; + Timer.Period = 2000; //TODO: Config? + Timer.Elapsed += TimerOnElapsed; + Logger = NullLogger.Instance; + } + + private void TimerOnElapsed(object sender, EventArgs e) + { + AsyncHelper.RunSync(RunAsync); + } + + public Task StartAsync(InboxConfig inboxConfig, CancellationToken cancellationToken = default) + { + InboxConfig = inboxConfig; + Inbox = (IEventInbox)ServiceProvider.GetRequiredService(inboxConfig.ImplementationType); + Timer.Start(cancellationToken); + return Task.CompletedTask; + } + + public Task StopAsync(CancellationToken cancellationToken = default) + { + Timer.Stop(cancellationToken); + return Task.CompletedTask; + } + + protected virtual async Task RunAsync() + { + await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName)) + { + if (handle != null) + { + Logger.LogDebug("Obtained the distributed lock: " + DistributedLockName); + + while (true) + { + var waitingEvents = await Inbox.GetWaitingEventsAsync(1000); //TODO: Config? + if (waitingEvents.Count <= 0) + { + break; + } + + Logger.LogInformation($"Found {waitingEvents.Count} events in the inbox."); + + foreach (var waitingEvent in waitingEvents) + { + using (var uow = UnitOfWorkManager.Begin(isTransactional: true, requiresNew: true)) + { + await DistributedEventBus + .AsRawEventPublisher() + .ProcessRawAsync(waitingEvent.EventName, waitingEvent.EventData); + + /* + await DistributedEventBus + .AsRawEventPublisher() + .PublishRawAsync(waitingEvent.Id, waitingEvent.EventName, waitingEvent.EventData); + */ + await Inbox.MarkAsProcessedAsync(waitingEvent.Id); + await uow.CompleteAsync(); + } + + Logger.LogInformation($"Processed the incoming event with id = {waitingEvent.Id:N}"); + } + } + } + else + { + Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName); + await Task.Delay(7000); //TODO: Can we pass a cancellation token to cancel on shutdown? (Config?) + } + } + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSenderManager.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSenderManager.cs index 617913e42f..6d403bd243 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSenderManager.cs +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSenderManager.cs @@ -31,7 +31,7 @@ namespace Volo.Abp.EventBus.Boxes if (outboxConfig.IsSendingEnabled) { var sender = ServiceProvider.GetRequiredService(); - await sender.StartAsync(outboxConfig); + await sender.StartAsync(outboxConfig, cancellationToken); Senders.Add(sender); } } @@ -41,7 +41,7 @@ namespace Volo.Abp.EventBus.Boxes { foreach (var sender in Senders) { - await sender.StopAsync(); + await sender.StopAsync(cancellationToken); } } } diff --git a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs index 5977cf159d..ae20ab1dd0 100644 --- a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs @@ -211,6 +211,23 @@ namespace Volo.Abp.EventBus.Kafka ); } + public override async Task ProcessRawAsync(string eventName, byte[] eventDataBytes) + { + var eventType = EventTypes.GetOrDefault(eventName); + if (eventType == null) + { + return; + } + + var eventData = Serializer.Deserialize(eventDataBytes, eventType); + var exceptions = new List(); + await TriggerHandlersAsync(eventType, eventData, exceptions); + if (exceptions.Any()) + { + ThrowOriginalExceptions(eventType, exceptions); + } + } + protected override byte[] Serialize(object eventData) { return Serializer.Serialize(eventData); diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs index 59fe5e4b99..d3d74a048c 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs @@ -123,7 +123,7 @@ namespace Volo.Abp.EventBus.RabbitMq retryAttempt = (int)ea.BasicProperties.Headers[EventErrorHandlerBase.RetryAttemptKey]; } - errorContext.EventData = Serializer.Deserialize(ea.Body.ToArray(), eventType); + errorContext.EventData = Serializer.Deserialize(eventBytes, eventType); errorContext.SetProperty(EventErrorHandlerBase.HeadersKey, ea.BasicProperties); errorContext.SetProperty(EventErrorHandlerBase.RetryAttemptKey, retryAttempt); }); @@ -217,6 +217,25 @@ namespace Volo.Abp.EventBus.RabbitMq return PublishAsync(eventName, eventData, null, eventId: eventId); } + public override async Task ProcessRawAsync(string eventName, byte[] eventDataBytes) + { + //TODO: We have a duplication in logic and also with the kafka side! + + var eventType = EventTypes.GetOrDefault(eventName); + if (eventType == null) + { + return; + } + + var eventData = Serializer.Deserialize(eventDataBytes, eventType); + var exceptions = new List(); + await TriggerHandlersAsync(eventType, eventData, exceptions); + if (exceptions.Any()) + { + ThrowOriginalExceptions(eventType, exceptions); + } + } + protected override byte[] Serialize(object eventData) { return Serializer.Serialize(eventData); diff --git a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs index 07995d67dd..bbe80d9866 100644 --- a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs @@ -194,6 +194,12 @@ namespace Volo.Abp.EventBus.Rebus throw new NotImplementedException(); } + public override Task ProcessRawAsync(string eventName, byte[] eventDataBytes) + { + /* TODO: IMPLEMENT! */ + throw new NotImplementedException(); + } + protected override byte[] Serialize(object eventData) { return Serializer.Serialize(eventData); diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs index dd4ca189ab..a436294e88 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs @@ -81,6 +81,7 @@ namespace Volo.Abp.EventBus.Distributed } public abstract Task PublishRawAsync(Guid eventId, string eventName, byte[] eventData); + public abstract Task ProcessRawAsync(string eventName, byte[] eventDataBytes); private async Task AddToOutboxAsync(Type eventType, object eventData) { @@ -128,6 +129,7 @@ namespace Volo.Abp.EventBus.Distributed if (inboxConfig.EventSelector == null || inboxConfig.EventSelector(eventType)) { var eventInbox = (IEventInbox) scope.ServiceProvider.GetRequiredService(inboxConfig.ImplementationType); + //TODO: Check if event was received before!! await eventInbox.EnqueueAsync( new IncomingEventInfo( GuidGenerator.Create(), diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventInbox.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventInbox.cs index dddedab2c0..137a410afa 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventInbox.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventInbox.cs @@ -1,4 +1,5 @@ -using System.Collections.Generic; +using System; +using System.Collections.Generic; using System.Threading.Tasks; namespace Volo.Abp.EventBus.Distributed @@ -8,5 +9,7 @@ namespace Volo.Abp.EventBus.Distributed Task EnqueueAsync(IncomingEventInfo incomingEvent); Task> GetWaitingEventsAsync(int maxCount); + + Task MarkAsProcessedAsync(Guid id); } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs index f53eb2b78b..5f96b818b7 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs @@ -3,11 +3,16 @@ using System.Threading.Tasks; namespace Volo.Abp.EventBus.Distributed { - public interface IRawEventPublisher + public interface IRawEventPublisher //TODO: Rename: ISupportsEventBoxes { Task PublishRawAsync( Guid eventId, string eventName, byte[] eventData); + + Task ProcessRawAsync( + string eventName, + byte[] eventDataBytes + ); } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs index c8b7eb6790..c4a61b24a6 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -162,6 +162,19 @@ namespace Volo.Abp.EventBus } } } + + protected void ThrowOriginalExceptions(Type eventType, List exceptions) + { + if (exceptions.Count == 1) + { + exceptions[0].ReThrow(); + } + + throw new AggregateException( + "More than one error has occurred while triggering the event: " + eventType, + exceptions + ); + } protected virtual void SubscribeHandlers(ITypeList handlers) { diff --git a/test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.Designer.cs b/test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.Designer.cs new file mode 100644 index 0000000000..02404db39f --- /dev/null +++ b/test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.Designer.cs @@ -0,0 +1,155 @@ +// +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("20210909182251_Added_Inbox_Process_Columns")] + partial class Added_Inbox_Process_Columns + { + 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("Id") + .HasColumnType("uniqueidentifier"); + + b.Property("ConcurrencyStamp") + .IsConcurrencyToken() + .HasMaxLength(40) + .HasColumnType("nvarchar(40)") + .HasColumnName("ConcurrencyStamp"); + + b.Property("CreationTime") + .HasColumnType("datetime2") + .HasColumnName("CreationTime"); + + b.Property("CreatorId") + .HasColumnType("uniqueidentifier") + .HasColumnName("CreatorId"); + + b.Property("ExtraProperties") + .HasColumnType("nvarchar(max)") + .HasColumnName("ExtraProperties"); + + b.Property("Text") + .IsRequired() + .HasMaxLength(128) + .HasColumnType("nvarchar(128)"); + + b.HasKey("Id"); + + b.ToTable("TodoItems"); + }); + + modelBuilder.Entity("DistDemoApp.TodoSummary", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("int") + .HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn); + + b.Property("ConcurrencyStamp") + .IsConcurrencyToken() + .HasMaxLength(40) + .HasColumnType("nvarchar(40)") + .HasColumnName("ConcurrencyStamp"); + + b.Property("Day") + .HasColumnType("tinyint"); + + b.Property("ExtraProperties") + .HasColumnType("nvarchar(max)") + .HasColumnName("ExtraProperties"); + + b.Property("Month") + .HasColumnType("tinyint"); + + b.Property("TotalCount") + .HasColumnType("int"); + + b.Property("Year") + .HasColumnType("int"); + + b.HasKey("Id"); + + b.ToTable("TodoSummaries"); + }); + + modelBuilder.Entity("Volo.Abp.EntityFrameworkCore.DistributedEvents.IncomingEventRecord", b => + { + b.Property("Id") + .HasColumnType("uniqueidentifier"); + + b.Property("CreationTime") + .HasColumnType("datetime2") + .HasColumnName("CreationTime"); + + b.Property("EventData") + .IsRequired() + .HasColumnType("varbinary(max)"); + + b.Property("EventName") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("nvarchar(256)"); + + b.Property("ExtraProperties") + .HasColumnType("nvarchar(max)") + .HasColumnName("ExtraProperties"); + + b.Property("Processed") + .HasColumnType("bit"); + + b.Property("ProcessedTime") + .HasColumnType("datetime2"); + + b.HasKey("Id"); + + b.ToTable("AbpEventInbox"); + }); + + modelBuilder.Entity("Volo.Abp.EntityFrameworkCore.DistributedEvents.OutgoingEventRecord", b => + { + b.Property("Id") + .HasColumnType("uniqueidentifier"); + + b.Property("CreationTime") + .HasColumnType("datetime2") + .HasColumnName("CreationTime"); + + b.Property("EventData") + .IsRequired() + .HasColumnType("varbinary(max)"); + + b.Property("EventName") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("nvarchar(256)"); + + b.Property("ExtraProperties") + .HasColumnType("nvarchar(max)") + .HasColumnName("ExtraProperties"); + + b.HasKey("Id"); + + b.ToTable("AbpEventOutbox"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.cs b/test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.cs new file mode 100644 index 0000000000..7da910a1af --- /dev/null +++ b/test/DistEvents/DistDemoApp/Migrations/20210909182251_Added_Inbox_Process_Columns.cs @@ -0,0 +1,35 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +namespace DistDemoApp.Migrations +{ + public partial class Added_Inbox_Process_Columns : Migration + { + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.AddColumn( + name: "Processed", + table: "AbpEventInbox", + type: "bit", + nullable: false, + defaultValue: false); + + migrationBuilder.AddColumn( + name: "ProcessedTime", + table: "AbpEventInbox", + type: "datetime2", + nullable: true); + } + + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropColumn( + name: "Processed", + table: "AbpEventInbox"); + + migrationBuilder.DropColumn( + name: "ProcessedTime", + table: "AbpEventInbox"); + } + } +} diff --git a/test/DistEvents/DistDemoApp/Migrations/TodoDbContextModelSnapshot.cs b/test/DistEvents/DistDemoApp/Migrations/TodoDbContextModelSnapshot.cs index 3423520512..87cf72bfd0 100644 --- a/test/DistEvents/DistDemoApp/Migrations/TodoDbContextModelSnapshot.cs +++ b/test/DistEvents/DistDemoApp/Migrations/TodoDbContextModelSnapshot.cs @@ -110,6 +110,12 @@ namespace DistDemoApp.Migrations .HasColumnType("nvarchar(max)") .HasColumnName("ExtraProperties"); + b.Property("Processed") + .HasColumnType("bit"); + + b.Property("ProcessedTime") + .HasColumnType("datetime2"); + b.HasKey("Id"); b.ToTable("AbpEventInbox");