|
|
@ -1,5 +1,6 @@ |
|
|
using System; |
|
|
using System; |
|
|
using System.Threading.Tasks; |
|
|
using System.Threading.Tasks; |
|
|
|
|
|
using Medallion.Threading; |
|
|
using Microsoft.Extensions.DependencyInjection; |
|
|
using Microsoft.Extensions.DependencyInjection; |
|
|
using Microsoft.Extensions.Logging; |
|
|
using Microsoft.Extensions.Logging; |
|
|
using Microsoft.Extensions.Logging.Abstractions; |
|
|
using Microsoft.Extensions.Logging.Abstractions; |
|
|
@ -15,17 +16,22 @@ namespace Volo.Abp.EventBus.Boxes |
|
|
protected IServiceProvider ServiceProvider { get; } |
|
|
protected IServiceProvider ServiceProvider { get; } |
|
|
protected AbpTimer Timer { get; } |
|
|
protected AbpTimer Timer { get; } |
|
|
protected IDistributedEventBus DistributedEventBus { get; } |
|
|
protected IDistributedEventBus DistributedEventBus { get; } |
|
|
|
|
|
protected IDistributedLockProvider DistributedLockProvider { get; } |
|
|
protected IEventOutbox Outbox { get; private set; } |
|
|
protected IEventOutbox Outbox { get; private set; } |
|
|
|
|
|
protected OutboxConfig OutboxConfig { get; private set; } |
|
|
|
|
|
protected string DistributedLockName => "Outbox_" + OutboxConfig.Name; |
|
|
public ILogger<OutboxSender> Logger { get; set; } |
|
|
public ILogger<OutboxSender> Logger { get; set; } |
|
|
|
|
|
|
|
|
public OutboxSender( |
|
|
public OutboxSender( |
|
|
IServiceProvider serviceProvider, |
|
|
IServiceProvider serviceProvider, |
|
|
AbpTimer timer, |
|
|
AbpTimer timer, |
|
|
IDistributedEventBus distributedEventBus) |
|
|
IDistributedEventBus distributedEventBus, |
|
|
|
|
|
IDistributedLockProvider distributedLockProvider) |
|
|
{ |
|
|
{ |
|
|
ServiceProvider = serviceProvider; |
|
|
ServiceProvider = serviceProvider; |
|
|
Timer = timer; |
|
|
Timer = timer; |
|
|
DistributedEventBus = distributedEventBus; |
|
|
DistributedEventBus = distributedEventBus; |
|
|
|
|
|
DistributedLockProvider = distributedLockProvider; |
|
|
Timer.Period = 2000; //TODO: Config?
|
|
|
Timer.Period = 2000; //TODO: Config?
|
|
|
Timer.Elapsed += TimerOnElapsed; |
|
|
Timer.Elapsed += TimerOnElapsed; |
|
|
Logger = NullLogger<OutboxSender>.Instance; |
|
|
Logger = NullLogger<OutboxSender>.Instance; |
|
|
@ -33,6 +39,7 @@ namespace Volo.Abp.EventBus.Boxes |
|
|
|
|
|
|
|
|
public virtual Task StartAsync(OutboxConfig outboxConfig) |
|
|
public virtual Task StartAsync(OutboxConfig outboxConfig) |
|
|
{ |
|
|
{ |
|
|
|
|
|
OutboxConfig = outboxConfig; |
|
|
Outbox = (IEventOutbox)ServiceProvider.GetRequiredService(outboxConfig.ImplementationType); |
|
|
Outbox = (IEventOutbox)ServiceProvider.GetRequiredService(outboxConfig.ImplementationType); |
|
|
Timer.Start(); |
|
|
Timer.Start(); |
|
|
return Task.CompletedTask; |
|
|
return Task.CompletedTask; |
|
|
@ -51,25 +58,39 @@ namespace Volo.Abp.EventBus.Boxes |
|
|
|
|
|
|
|
|
protected virtual async Task RunAsync() |
|
|
protected virtual async Task RunAsync() |
|
|
{ |
|
|
{ |
|
|
while (true) |
|
|
await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName)) |
|
|
{ |
|
|
{ |
|
|
var waitingEvents = await Outbox.GetWaitingEventsAsync(1000); //TODO: Config?
|
|
|
if (handle != null) |
|
|
if (waitingEvents.Count <= 0) |
|
|
|
|
|
{ |
|
|
{ |
|
|
break; |
|
|
Logger.LogDebug("Obtained the distributed lock: " + DistributedLockName); |
|
|
} |
|
|
|
|
|
|
|
|
while (true) |
|
|
|
|
|
{ |
|
|
|
|
|
var waitingEvents = await Outbox.GetWaitingEventsAsync(1000); //TODO: Config?
|
|
|
|
|
|
if (waitingEvents.Count <= 0) |
|
|
|
|
|
{ |
|
|
|
|
|
break; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
Logger.LogInformation($"Found {waitingEvents.Count} events in the outbox."); |
|
|
Logger.LogInformation($"Found {waitingEvents.Count} events in the outbox."); |
|
|
|
|
|
|
|
|
foreach (var waitingEvent in waitingEvents) |
|
|
foreach (var waitingEvent in waitingEvents) |
|
|
{ |
|
|
{ |
|
|
await DistributedEventBus |
|
|
await DistributedEventBus |
|
|
.AsRawEventPublisher() |
|
|
.AsRawEventPublisher() |
|
|
.PublishRawAsync(waitingEvent.Id, waitingEvent.EventName, waitingEvent.EventData); |
|
|
.PublishRawAsync(waitingEvent.Id, waitingEvent.EventName, waitingEvent.EventData); |
|
|
|
|
|
|
|
|
await Outbox.DeleteAsync(waitingEvent.Id); |
|
|
await Outbox.DeleteAsync(waitingEvent.Id); |
|
|
|
|
|
|
|
|
Logger.LogInformation($"Sent the event to the message broker with id = {waitingEvent.Id:N}"); |
|
|
Logger.LogInformation($"Sent the event to the message broker with id = {waitingEvent.Id:N}"); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
await Task.Delay(30000); |
|
|
|
|
|
} |
|
|
|
|
|
else |
|
|
|
|
|
{ |
|
|
|
|
|
Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|