Browse Source

Complete the Outbox & Inbox Patterns feature

pull/10159/head
liangshiwei 5 years ago
parent
commit
0cf5d248fd
  1. 19
      framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs
  2. 15
      framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventOutbox.cs
  3. 43
      framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesOptions.cs
  4. 16
      framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs
  5. 16
      framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs
  6. 0
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/IRebusSerializer.cs
  7. 37
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs
  8. 2
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventHandlerAdapter.cs
  9. 13
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs
  10. 11
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventInbox.cs
  11. 9
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventOutbox.cs
  12. 60
      framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventInbox.cs
  13. 61
      framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventOutbox.cs
  14. 2
      test/DistEvents/DistDemoApp.EfCoreRabbitMq/appsettings.json
  15. 12
      test/DistEvents/DistDemoApp.MongoDbKafka/DistDemoApp.MongoDbKafka.csproj
  16. 38
      test/DistEvents/DistDemoApp.MongoDbKafka/DistDemoAppMongoDbKafkaModule.cs
  17. 53
      test/DistEvents/DistDemoApp.MongoDbKafka/Program.cs
  18. 26
      test/DistEvents/DistDemoApp.MongoDbKafka/TodoMongoDbContext.cs
  19. 19
      test/DistEvents/DistDemoApp.MongoDbKafka/appsettings.json
  20. 21
      test/DistEvents/DistDemoApp.MongoDbRebus/DistDemoApp.MongoDbRebus.csproj
  21. 53
      test/DistEvents/DistDemoApp.MongoDbRebus/DistDemoAppMongoDbRebusModule.cs
  22. 57
      test/DistEvents/DistDemoApp.MongoDbRebus/Program.cs
  23. 26
      test/DistEvents/DistDemoApp.MongoDbRebus/TodoMongoDbContext.cs
  24. 19
      test/DistEvents/DistDemoApp.MongoDbRebus/appsettings.json
  25. 6
      test/DistEvents/DistEventsDemo.sln

19
framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs

@ -1,26 +1,31 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Options;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Timing;
using Volo.Abp.Uow;
namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
{
public class DbContextEventInbox<TDbContext> : IDbContextEventInbox<TDbContext>
public class DbContextEventInbox<TDbContext> : IDbContextEventInbox<TDbContext>
where TDbContext : IHasEventInbox
{
protected IDbContextProvider<TDbContext> DbContextProvider { get; }
protected AbpDistributedEventBusOptions DistributedEventsOptions { get; }
protected IClock Clock { get; }
public DbContextEventInbox(
IDbContextProvider<TDbContext> dbContextProvider,
IClock clock)
IClock clock,
IOptions<AbpDistributedEventBusOptions> distributedEventsOptions)
{
DbContextProvider = dbContextProvider;
Clock = clock;
DistributedEventsOptions = distributedEventsOptions.Value;
}
[UnitOfWork]
@ -34,7 +39,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
}
[UnitOfWork]
public virtual async Task<List<IncomingEventInfo>> GetWaitingEventsAsync(int maxCount)
public virtual async Task<List<IncomingEventInfo>> GetWaitingEventsAsync(int maxCount, CancellationToken cancellationToken = default)
{
var dbContext = await DbContextProvider.GetDbContextAsync();
@ -44,8 +49,8 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
.Where(x => !x.Processed)
.OrderBy(x => x.CreationTime)
.Take(maxCount)
.ToListAsync();
.ToListAsync(cancellationToken: cancellationToken);
return outgoingEventRecords
.Select(x => x.ToIncomingEventInfo())
.ToList();
@ -76,11 +81,11 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
{
//TODO: Optimize
var dbContext = await DbContextProvider.GetDbContextAsync();
var timeToKeepEvents = Clock.Now.AddHours(-2); //TODO: Config?
var timeToKeepEvents = Clock.Now.Add(DistributedEventsOptions.InboxKeepEventTimeSpan);
var oldEvents = await dbContext.IncomingEvents
.Where(x => x.Processed && x.CreationTime < timeToKeepEvents)
.ToListAsync();
dbContext.IncomingEvents.RemoveRange(oldEvents);
}
}
}
}

15
framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventOutbox.cs

@ -1,6 +1,7 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
using Volo.Abp.EventBus.Distributed;
@ -8,7 +9,7 @@ using Volo.Abp.Uow;
namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
{
public class DbContextEventOutbox<TDbContext> : IDbContextEventOutbox<TDbContext>
public class DbContextEventOutbox<TDbContext> : IDbContextEventOutbox<TDbContext>
where TDbContext : IHasEventOutbox
{
protected IDbContextProvider<TDbContext> DbContextProvider { get; }
@ -18,7 +19,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
{
DbContextProvider = dbContextProvider;
}
[UnitOfWork]
public virtual async Task EnqueueAsync(OutgoingEventInfo outgoingEvent)
{
@ -29,17 +30,17 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
}
[UnitOfWork]
public virtual async Task<List<OutgoingEventInfo>> GetWaitingEventsAsync(int maxCount)
public virtual async Task<List<OutgoingEventInfo>> GetWaitingEventsAsync(int maxCount, CancellationToken cancellationToken = default)
{
var dbContext = (IHasEventOutbox) await DbContextProvider.GetDbContextAsync();
var outgoingEventRecords = await dbContext
.OutgoingEvents
.AsNoTracking()
.OrderBy(x => x.CreationTime)
.Take(maxCount)
.ToListAsync();
.ToListAsync(cancellationToken: cancellationToken);
return outgoingEventRecords
.Select(x => x.ToOutgoingEventInfo())
.ToList();
@ -57,4 +58,4 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents
}
}
}
}
}

43
framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesOptions.cs

@ -0,0 +1,43 @@
using System;
namespace Volo.Abp.EventBus.Boxes
{
public class AbpEventBusBoxesOptions
{
/// <summary>
/// Default: 6 hours
/// </summary>
public TimeSpan CleanOldEventTimeIntervalSpan { get; set; }
/// <summary>
/// Default: 1000
/// </summary>
public int InboxWaitingEventMaxCount { get; set; }
/// <summary>
/// Default: 1000
/// </summary>
public int OutboxWaitingEventMaxCount { get; set; }
/// <summary>
/// Period time of <see cref="InboxProcessor"/> and <see cref="OutboxSender"/>
/// Default: 2 seconds
/// </summary>
public TimeSpan PeriodTimeSpan { get; set; }
/// <summary>
/// Delay time of <see cref="InboxProcessor"/> and <see cref="OutboxSender"/>
/// Default: 15 seconds
/// </summary>
public TimeSpan DelayTimeSpan { get; set; }
public AbpEventBusBoxesOptions()
{
CleanOldEventTimeIntervalSpan = TimeSpan.FromHours(6);
InboxWaitingEventMaxCount = 1000;
OutboxWaitingEventMaxCount = 1000;
PeriodTimeSpan = TimeSpan.FromSeconds(2);
DelayTimeSpan = TimeSpan.FromSeconds(15);
}
}
}

16
framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs

@ -5,6 +5,7 @@ using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using Volo.Abp.DependencyInjection;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Threading;
@ -23,6 +24,7 @@ namespace Volo.Abp.EventBus.Boxes
protected IClock Clock { get; }
protected IEventInbox Inbox { get; private set; }
protected InboxConfig InboxConfig { get; private set; }
protected AbpEventBusBoxesOptions EventBusBoxesOptions { get; }
protected DateTime? LastCleanTime { get; set; }
@ -37,7 +39,8 @@ namespace Volo.Abp.EventBus.Boxes
IDistributedEventBus distributedEventBus,
IDistributedLockProvider distributedLockProvider,
IUnitOfWorkManager unitOfWorkManager,
IClock clock)
IClock clock,
IOptions<AbpEventBusBoxesOptions> eventBusBoxesOptions)
{
ServiceProvider = serviceProvider;
Timer = timer;
@ -45,7 +48,8 @@ namespace Volo.Abp.EventBus.Boxes
DistributedLockProvider = distributedLockProvider;
UnitOfWorkManager = unitOfWorkManager;
Clock = clock;
Timer.Period = 2000; //TODO: Config?
EventBusBoxesOptions = eventBusBoxesOptions.Value;
Timer.Period = EventBusBoxesOptions.PeriodTimeSpan.Seconds;
Timer.Elapsed += TimerOnElapsed;
Logger = NullLogger<InboxProcessor>.Instance;
StoppingTokenSource = new CancellationTokenSource();
@ -79,7 +83,7 @@ namespace Volo.Abp.EventBus.Boxes
{
return;
}
await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName, cancellationToken: StoppingToken))
{
if (handle != null)
@ -88,7 +92,7 @@ namespace Volo.Abp.EventBus.Boxes
while (true)
{
var waitingEvents = await Inbox.GetWaitingEventsAsync(1000); //TODO: Config? Pass StoppingToken!
var waitingEvents = await Inbox.GetWaitingEventsAsync(EventBusBoxesOptions.InboxWaitingEventMaxCount, StoppingToken);
if (waitingEvents.Count <= 0)
{
break;
@ -116,14 +120,14 @@ namespace Volo.Abp.EventBus.Boxes
else
{
Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName);
await TaskDelayHelper.DelayAsync(15000, StoppingToken); //TODO: Config?
await TaskDelayHelper.DelayAsync(EventBusBoxesOptions.DelayTimeSpan.Milliseconds, StoppingToken);
}
}
}
protected virtual async Task DeleteOldEventsAsync()
{
if (LastCleanTime != null && LastCleanTime > Clock.Now.AddHours(6)) //TODO: Config?
if (LastCleanTime != null && LastCleanTime > Clock.Now.Add(EventBusBoxesOptions.CleanOldEventTimeIntervalSpan))
{
return;
}

16
framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs

@ -5,6 +5,7 @@ using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using Volo.Abp.DependencyInjection;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Threading;
@ -19,9 +20,10 @@ namespace Volo.Abp.EventBus.Boxes
protected IDistributedLockProvider DistributedLockProvider { get; }
protected IEventOutbox Outbox { get; private set; }
protected OutboxConfig OutboxConfig { get; private set; }
protected AbpEventBusBoxesOptions EventBusBoxesOptions { get; }
protected string DistributedLockName => "Outbox_" + OutboxConfig.Name;
public ILogger<OutboxSender> Logger { get; set; }
protected CancellationTokenSource StoppingTokenSource { get; }
protected CancellationToken StoppingToken { get; }
@ -29,13 +31,15 @@ namespace Volo.Abp.EventBus.Boxes
IServiceProvider serviceProvider,
AbpAsyncTimer timer,
IDistributedEventBus distributedEventBus,
IDistributedLockProvider distributedLockProvider)
IDistributedLockProvider distributedLockProvider,
IOptions<AbpEventBusBoxesOptions> eventBusBoxesOptions)
{
ServiceProvider = serviceProvider;
Timer = timer;
DistributedEventBus = distributedEventBus;
DistributedLockProvider = distributedLockProvider;
Timer.Period = 2000; //TODO: Config?
EventBusBoxesOptions = eventBusBoxesOptions.Value;
Timer.Period = EventBusBoxesOptions.PeriodTimeSpan.Seconds;
Timer.Elapsed += TimerOnElapsed;
Logger = NullLogger<OutboxSender>.Instance;
StoppingTokenSource = new CancellationTokenSource();
@ -65,13 +69,13 @@ namespace Volo.Abp.EventBus.Boxes
protected virtual async Task RunAsync()
{
await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName))
await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName, cancellationToken: StoppingToken))
{
if (handle != null)
{
while (true)
{
var waitingEvents = await Outbox.GetWaitingEventsAsync(1000); //TODO: Config?
var waitingEvents = await Outbox.GetWaitingEventsAsync(EventBusBoxesOptions.OutboxWaitingEventMaxCount, StoppingToken);
if (waitingEvents.Count <= 0)
{
break;
@ -96,7 +100,7 @@ namespace Volo.Abp.EventBus.Boxes
else
{
Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName);
await TaskDelayHelper.DelayAsync(15000, StoppingToken); //TODO: Config?
await TaskDelayHelper.DelayAsync(EventBusBoxesOptions.DelayTimeSpan.Milliseconds, StoppingToken);
}
}
}

0
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/IRabbitMqSerializer.cs → framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/IRebusSerializer.cs

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

@ -6,6 +6,7 @@ using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Rebus.Bus;
using Rebus.Pipeline;
using Volo.Abp.DependencyInjection;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Guids;
@ -134,6 +135,19 @@ namespace Volo.Abp.EventBus.Rebus
Rebus.Unsubscribe(eventType);
}
public async Task ProcessEventAsync(Type eventType, object eventData)
{
var messageId = MessageContext.Current.TransportMessage.GetMessageId();
var eventName = EventNameAttribute.GetNameOrDefault(eventType);
if (await AddToInboxAsync(messageId, eventName, eventType, MessageContext.Current.TransportMessage.Body))
{
return;
}
await TriggerHandlersAsync(eventType, eventData);
}
protected override async Task PublishToEventBusAsync(Type eventType, object eventData)
{
await AbpRebusEventBusOptions.Publish(Rebus, eventType, eventData);
@ -192,16 +206,29 @@ namespace Volo.Abp.EventBus.Rebus
OutgoingEventInfo outgoingEvent,
OutboxConfig outboxConfig)
{
/* TODO: IMPLEMENT! */
throw new NotImplementedException();
var eventType = EventTypes.GetOrDefault(outgoingEvent.EventName);
var eventData = Serializer.Deserialize(outgoingEvent.EventData, eventType);
return PublishToEventBusAsync(eventType, eventData);
}
public override Task ProcessFromInboxAsync(
public override async Task ProcessFromInboxAsync(
IncomingEventInfo incomingEvent,
InboxConfig inboxConfig)
{
/* TODO: IMPLEMENT! */
throw new NotImplementedException();
var eventType = EventTypes.GetOrDefault(incomingEvent.EventName);
if (eventType == null)
{
return;
}
var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType);
var exceptions = new List<Exception>();
await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig);
if (exceptions.Any())
{
ThrowOriginalExceptions(eventType, exceptions);
}
}
protected override byte[] Serialize(object eventData)

2
framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventHandlerAdapter.cs

@ -14,7 +14,7 @@ namespace Volo.Abp.EventBus.Rebus
public async Task Handle(TEventData message)
{
await RebusDistributedEventBus.TriggerHandlersAsync(typeof(TEventData), message);
await RebusDistributedEventBus.ProcessEventAsync(message.GetType(), message);
}
}
}

13
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs

@ -1,3 +1,4 @@
using System;
using Volo.Abp.Collections;
namespace Volo.Abp.EventBus.Distributed
@ -5,16 +6,22 @@ namespace Volo.Abp.EventBus.Distributed
public class AbpDistributedEventBusOptions
{
public ITypeList<IEventHandler> Handlers { get; }
public OutboxConfigDictionary Outboxes { get; }
public InboxConfigDictionary Inboxes { get; }
/// <summary>
/// Default: -2 hours
/// </summary>
public TimeSpan InboxKeepEventTimeSpan { get; set; }
public AbpDistributedEventBusOptions()
{
Handlers = new TypeList<IEventHandler>();
Outboxes = new OutboxConfigDictionary();
Inboxes = new InboxConfigDictionary();
InboxKeepEventTimeSpan = TimeSpan.FromHours(-2);
}
}
}
}

11
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventInbox.cs

@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
namespace Volo.Abp.EventBus.Distributed
@ -7,13 +8,13 @@ namespace Volo.Abp.EventBus.Distributed
public interface IEventInbox
{
Task EnqueueAsync(IncomingEventInfo incomingEvent);
Task<List<IncomingEventInfo>> GetWaitingEventsAsync(int maxCount);
Task<List<IncomingEventInfo>> GetWaitingEventsAsync(int maxCount, CancellationToken cancellationToken = default);
Task MarkAsProcessedAsync(Guid id);
Task<bool> ExistsByMessageIdAsync(string messageId);
Task DeleteOldEventsAsync();
}
}
}

9
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IEventOutbox.cs

@ -1,5 +1,6 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
namespace Volo.Abp.EventBus.Distributed
@ -7,9 +8,9 @@ namespace Volo.Abp.EventBus.Distributed
public interface IEventOutbox
{
Task EnqueueAsync(OutgoingEventInfo outgoingEvent);
Task<List<OutgoingEventInfo>> GetWaitingEventsAsync(int maxCount);
Task<List<OutgoingEventInfo>> GetWaitingEventsAsync(int maxCount, CancellationToken cancellationToken = default);
Task DeleteAsync(Guid id);
}
}
}

60
framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventInbox.cs

@ -1,7 +1,9 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Options;
using MongoDB.Driver;
using MongoDB.Driver.Linq;
using Volo.Abp.EventBus.Distributed;
@ -10,21 +12,24 @@ using Volo.Abp.Uow;
namespace Volo.Abp.MongoDB.DistributedEvents
{
public class MongoDbContextEventInbox<TMongoDbContext> : IMongoDbContextEventInbox<TMongoDbContext>
public class MongoDbContextEventInbox<TMongoDbContext> : IMongoDbContextEventInbox<TMongoDbContext>
where TMongoDbContext : IHasEventInbox
{
protected IMongoDbContextProvider<TMongoDbContext> DbContextProvider { get; }
protected AbpDistributedEventBusOptions DistributedEventsOptions { get; }
protected IClock Clock { get; }
public MongoDbContextEventInbox(
IMongoDbContextProvider<TMongoDbContext> dbContextProvider,
IClock clock)
IClock clock,
IOptions<AbpDistributedEventBusOptions> distributedEventsOptions)
{
DbContextProvider = dbContextProvider;
Clock = clock;
DistributedEventsOptions = distributedEventsOptions.Value;
}
[UnitOfWork]
public virtual async Task EnqueueAsync(IncomingEventInfo incomingEvent)
{
@ -45,9 +50,9 @@ namespace Volo.Abp.MongoDB.DistributedEvents
}
[UnitOfWork]
public virtual async Task<List<IncomingEventInfo>> GetWaitingEventsAsync(int maxCount)
public virtual async Task<List<IncomingEventInfo>> GetWaitingEventsAsync(int maxCount, CancellationToken cancellationToken = default)
{
var dbContext = await DbContextProvider.GetDbContextAsync();
var dbContext = await DbContextProvider.GetDbContextAsync(cancellationToken);
var outgoingEventRecords = await dbContext
.IncomingEvents
@ -55,8 +60,8 @@ namespace Volo.Abp.MongoDB.DistributedEvents
.Where(x => !x.Processed)
.OrderBy(x => x.CreationTime)
.Take(maxCount)
.ToListAsync();
.ToListAsync(cancellationToken: cancellationToken);
return outgoingEventRecords
.Select(x => x.ToIncomingEventInfo())
.ToList();
@ -65,17 +70,44 @@ namespace Volo.Abp.MongoDB.DistributedEvents
[UnitOfWork]
public async Task MarkAsProcessedAsync(Guid id)
{
throw new NotImplementedException();
var dbContext = await DbContextProvider.GetDbContextAsync();
var incomingEvent = await dbContext.IncomingEvents.Find(x => x.Id.Equals(id)).FirstOrDefaultAsync();
if (incomingEvent != null)
{
incomingEvent.MarkAsProcessed(Clock.Now);
if (dbContext.SessionHandle != null)
{
await dbContext.IncomingEvents.ReplaceOneAsync(dbContext.SessionHandle, Builders<IncomingEventRecord>.Filter.Eq(e => e.Id, incomingEvent.Id), incomingEvent);
}
else
{
await dbContext.IncomingEvents.ReplaceOneAsync(Builders<IncomingEventRecord>.Filter.Eq(e => e.Id, incomingEvent.Id), incomingEvent);
}
}
}
public Task<bool> ExistsByMessageIdAsync(string messageId)
[UnitOfWork]
public async Task<bool> ExistsByMessageIdAsync(string messageId)
{
throw new NotImplementedException();
var dbContext = await DbContextProvider.GetDbContextAsync();
return await dbContext.IncomingEvents.AsQueryable().AnyAsync(x => x.MessageId == messageId);
}
public Task DeleteOldEventsAsync()
[UnitOfWork]
public async Task DeleteOldEventsAsync()
{
throw new NotImplementedException();
var dbContext = await DbContextProvider.GetDbContextAsync();
var timeToKeepEvents = Clock.Now.Add(DistributedEventsOptions.InboxKeepEventTimeSpan);
if (dbContext.SessionHandle != null)
{
await dbContext.IncomingEvents.DeleteManyAsync(dbContext.SessionHandle, x => x.Processed && x.CreationTime < timeToKeepEvents);
}
else
{
await dbContext.IncomingEvents.DeleteManyAsync(x => x.Processed && x.CreationTime < timeToKeepEvents);
}
}
}
}
}

61
framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventOutbox.cs

@ -1,27 +1,72 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MongoDB.Driver;
using MongoDB.Driver.Linq;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Uow;
namespace Volo.Abp.MongoDB.DistributedEvents
{
public class MongoDbContextEventOutbox<TMongoDbContext> : IMongoDbContextEventOutbox<TMongoDbContext>
public class MongoDbContextEventOutbox<TMongoDbContext> : IMongoDbContextEventOutbox<TMongoDbContext>
where TMongoDbContext : IHasEventOutbox
{
public Task EnqueueAsync(OutgoingEventInfo outgoingEvent)
protected IMongoDbContextProvider<TMongoDbContext> MongoDbContextProvider { get; }
public MongoDbContextEventOutbox(IMongoDbContextProvider<TMongoDbContext> mongoDbContextProvider)
{
MongoDbContextProvider = mongoDbContextProvider;
}
[UnitOfWork]
public async Task EnqueueAsync(OutgoingEventInfo outgoingEvent)
{
throw new NotImplementedException();
var dbContext = (IHasEventOutbox) await MongoDbContextProvider.GetDbContextAsync();
if (dbContext.SessionHandle != null)
{
await dbContext.OutgoingEvents.InsertOneAsync(
dbContext.SessionHandle,
new OutgoingEventRecord(outgoingEvent)
);
}
else
{
await dbContext.OutgoingEvents.InsertOneAsync(
new OutgoingEventRecord(outgoingEvent)
);
}
}
public Task<List<OutgoingEventInfo>> GetWaitingEventsAsync(int maxCount)
[UnitOfWork]
public async Task<List<OutgoingEventInfo>> GetWaitingEventsAsync(int maxCount, CancellationToken cancellationToken = default)
{
throw new NotImplementedException();
var dbContext = (IHasEventOutbox) await MongoDbContextProvider.GetDbContextAsync(cancellationToken);
var outgoingEventRecords = await dbContext
.OutgoingEvents.AsQueryable()
.OrderBy(x => x.CreationTime)
.Take(maxCount)
.ToListAsync(cancellationToken: cancellationToken);
return outgoingEventRecords
.Select(x => x.ToOutgoingEventInfo())
.ToList();
}
public Task DeleteAsync(Guid id)
[UnitOfWork]
public async Task DeleteAsync(Guid id)
{
throw new NotImplementedException();
var dbContext = (IHasEventOutbox) await MongoDbContextProvider.GetDbContextAsync();
if (dbContext.SessionHandle != null)
{
await dbContext.OutgoingEvents.DeleteOneAsync(dbContext.SessionHandle, x => x.Id.Equals(id));
}
else
{
await dbContext.OutgoingEvents.DeleteOneAsync(x => x.Id.Equals(id));
}
}
}
}
}

2
test/DistEvents/DistDemoApp.EfCoreRabbitMq/appsettings.json

@ -16,4 +16,4 @@
"Redis": {
"Configuration": "127.0.0.1"
}
}
}

12
test/DistEvents/DistDemoApp.MongoDbKafka/DistDemoApp.MongoDbKafka.csproj

@ -6,4 +6,16 @@
<RootNamespace>DistDemoApp</RootNamespace>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\framework\src\Volo.Abp.MongoDB\Volo.Abp.MongoDB.csproj" />
<ProjectReference Include="..\..\..\framework\src\Volo.Abp.EventBus.Kafka\Volo.Abp.EventBus.Kafka.csproj" />
<ProjectReference Include="..\DistDemoApp.Shared\DistDemoApp.Shared.csproj" />
</ItemGroup>
<ItemGroup>
<None Update="appsettings.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
</ItemGroup>
</Project>

38
test/DistEvents/DistDemoApp.MongoDbKafka/DistDemoAppMongoDbKafkaModule.cs

@ -0,0 +1,38 @@
using Microsoft.Extensions.DependencyInjection;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Kafka;
using Volo.Abp.Modularity;
using Volo.Abp.MongoDB;
using Volo.Abp.MongoDB.DistributedEvents;
namespace DistDemoApp
{
[DependsOn(
typeof(AbpMongoDbModule),
typeof(AbpEventBusKafkaModule),
typeof(DistDemoAppSharedModule)
)]
public class DistDemoAppMongoDbKafkaModule : AbpModule
{
public override void ConfigureServices(ServiceConfigurationContext context)
{
context.Services.AddMongoDbContext<TodoMongoDbContext>(options =>
{
options.AddDefaultRepositories();
});
Configure<AbpDistributedEventBusOptions>(options =>
{
options.Outboxes.Configure(config =>
{
config.UseMongoDbContext<TodoMongoDbContext>();
});
options.Inboxes.Configure(config =>
{
config.UseMongoDbContext<TodoMongoDbContext>();
});
});
}
}
}

53
test/DistEvents/DistDemoApp.MongoDbKafka/Program.cs

@ -1,12 +1,57 @@
using System;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Serilog;
using Serilog.Events;
namespace DistDemoApp
{
class Program
public class Program
{
static void Main(string[] args)
public static async Task<int> Main(string[] args)
{
Console.WriteLine("Hello World!");
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<DistDemoAppMongoDbKafkaModule>();
});
}
}
}

26
test/DistEvents/DistDemoApp.MongoDbKafka/TodoMongoDbContext.cs

@ -0,0 +1,26 @@
using MongoDB.Driver;
using Volo.Abp.Data;
using Volo.Abp.MongoDB;
using Volo.Abp.MongoDB.DistributedEvents;
namespace DistDemoApp
{
[ConnectionStringName("Default")]
public class TodoMongoDbContext : AbpMongoDbContext, IHasEventOutbox, IHasEventInbox
{
public IMongoCollection<TodoItem> TodoItems => Collection<TodoItem>();
public IMongoCollection<TodoSummary> TodoSummaries => Collection<TodoSummary>();
public IMongoCollection<OutgoingEventRecord> OutgoingEvents
{
get => Collection<OutgoingEventRecord>();
set {}
}
public IMongoCollection<IncomingEventRecord> IncomingEvents
{
get => Collection<IncomingEventRecord>();
set {}
}
}
}

19
test/DistEvents/DistDemoApp.MongoDbKafka/appsettings.json

@ -0,0 +1,19 @@
{
"ConnectionStrings": {
"Default": "mongodb://localhost:27018,localhost:27019,localhost:27020/DistEventsDemo"
},
"Kafka": {
"Connections": {
"Default": {
"BootstrapServers": "localhost:9092"
}
},
"EventBus": {
"GroupId": "DistDemoApp",
"TopicName": "DistDemoTopic"
}
},
"Redis": {
"Configuration": "127.0.0.1"
}
}

21
test/DistEvents/DistDemoApp.MongoDbRebus/DistDemoApp.MongoDbRebus.csproj

@ -0,0 +1,21 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net5.0</TargetFramework>
<RootNamespace>DistDemoApp</RootNamespace>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\framework\src\Volo.Abp.MongoDB\Volo.Abp.MongoDB.csproj" />
<ProjectReference Include="..\..\..\framework\src\Volo.Abp.EventBus.Rebus\Volo.Abp.EventBus.Rebus.csproj" />
<ProjectReference Include="..\DistDemoApp.Shared\DistDemoApp.Shared.csproj" />
</ItemGroup>
<ItemGroup>
<None Update="appsettings.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
</ItemGroup>
</Project>

53
test/DistEvents/DistDemoApp.MongoDbRebus/DistDemoAppMongoDbRebusModule.cs

@ -0,0 +1,53 @@
using Microsoft.Extensions.DependencyInjection;
using Rebus.Persistence.InMem;
using Rebus.Transport.InMem;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Rebus;
using Volo.Abp.Modularity;
using Volo.Abp.MongoDB;
using Volo.Abp.MongoDB.DistributedEvents;
namespace DistDemoApp
{
[DependsOn(
typeof(AbpMongoDbModule),
typeof(AbpEventBusRebusModule),
typeof(DistDemoAppSharedModule)
)]
public class DistDemoAppMongoDbRebusModule : AbpModule
{
public override void PreConfigureServices(ServiceConfigurationContext context)
{
PreConfigure<AbpRebusEventBusOptions>(options =>
{
options.InputQueueName = "eventbus";
options.Configurer = rebusConfigurer =>
{
rebusConfigurer.Transport(t => t.UseInMemoryTransport(new InMemNetwork(), "eventbus"));
rebusConfigurer.Subscriptions(s => s.StoreInMemory());
};
});
}
public override void ConfigureServices(ServiceConfigurationContext context)
{
context.Services.AddMongoDbContext<TodoMongoDbContext>(options =>
{
options.AddDefaultRepositories();
});
Configure<AbpDistributedEventBusOptions>(options =>
{
options.Outboxes.Configure(config =>
{
config.UseMongoDbContext<TodoMongoDbContext>();
});
options.Inboxes.Configure(config =>
{
config.UseMongoDbContext<TodoMongoDbContext>();
});
});
}
}
}

57
test/DistEvents/DistDemoApp.MongoDbRebus/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<DistDemoAppMongoDbRebusModule>();
});
}
}

26
test/DistEvents/DistDemoApp.MongoDbRebus/TodoMongoDbContext.cs

@ -0,0 +1,26 @@
using MongoDB.Driver;
using Volo.Abp.Data;
using Volo.Abp.MongoDB;
using Volo.Abp.MongoDB.DistributedEvents;
namespace DistDemoApp
{
[ConnectionStringName("Default")]
public class TodoMongoDbContext : AbpMongoDbContext, IHasEventOutbox, IHasEventInbox
{
public IMongoCollection<TodoItem> TodoItems => Collection<TodoItem>();
public IMongoCollection<TodoSummary> TodoSummaries => Collection<TodoSummary>();
public IMongoCollection<OutgoingEventRecord> OutgoingEvents
{
get => Collection<OutgoingEventRecord>();
set {}
}
public IMongoCollection<IncomingEventRecord> IncomingEvents
{
get => Collection<IncomingEventRecord>();
set {}
}
}
}

19
test/DistEvents/DistDemoApp.MongoDbRebus/appsettings.json

@ -0,0 +1,19 @@
{
"ConnectionStrings": {
"Default": "mongodb://localhost:27018,localhost:27019,localhost:27020/DistEventsDemo"
},
"Kafka": {
"Connections": {
"Default": {
"BootstrapServers": "localhost:9092"
}
},
"EventBus": {
"GroupId": "DistDemoApp",
"TopicName": "DistDemoTopic"
}
},
"Redis": {
"Configuration": "127.0.0.1"
}
}

6
test/DistEvents/DistEventsDemo.sln

@ -6,6 +6,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "DistDemoApp.MongoDbKafka",
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "DistDemoApp.Shared", "DistDemoApp.Shared\DistDemoApp.Shared.csproj", "{C515F4E2-0ED3-4561-BC58-FC633B50E2EB}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "DistDemoApp.MongoDbRebus", "DistDemoApp.MongoDbRebus\DistDemoApp.MongoDbRebus.csproj", "{4FB63540-4CC5-4A7B-900B-F5FCD907456E}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -24,5 +26,9 @@ Global
{C515F4E2-0ED3-4561-BC58-FC633B50E2EB}.Debug|Any CPU.Build.0 = Debug|Any CPU
{C515F4E2-0ED3-4561-BC58-FC633B50E2EB}.Release|Any CPU.ActiveCfg = Release|Any CPU
{C515F4E2-0ED3-4561-BC58-FC633B50E2EB}.Release|Any CPU.Build.0 = Release|Any CPU
{4FB63540-4CC5-4A7B-900B-F5FCD907456E}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{4FB63540-4CC5-4A7B-900B-F5FCD907456E}.Debug|Any CPU.Build.0 = Debug|Any CPU
{4FB63540-4CC5-4A7B-900B-F5FCD907456E}.Release|Any CPU.ActiveCfg = Release|Any CPU
{4FB63540-4CC5-4A7B-900B-F5FCD907456E}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
EndGlobal

Loading…
Cancel
Save