mirror of https://github.com/abpframework/abp.git
19 changed files with 470 additions and 7 deletions
@ -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); |
|||
} |
|||
} |
|||
@ -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<IInboxProcessor> Processors { get; } |
|||
|
|||
public InboxProcessManager( |
|||
IOptions<AbpDistributedEventBusOptions> options, |
|||
IServiceProvider serviceProvider) |
|||
{ |
|||
ServiceProvider = serviceProvider; |
|||
Options = options.Value; |
|||
Processors = new List<IInboxProcessor>(); |
|||
} |
|||
|
|||
public async Task StartAsync(CancellationToken cancellationToken = default) |
|||
{ |
|||
foreach (var inboxConfig in Options.Inboxes.Values) |
|||
{ |
|||
if (inboxConfig.IsProcessingEnabled) |
|||
{ |
|||
var processor = ServiceProvider.GetRequiredService<IInboxProcessor>(); |
|||
await processor.StartAsync(inboxConfig, cancellationToken); |
|||
Processors.Add(processor); |
|||
} |
|||
} |
|||
} |
|||
|
|||
public async Task StopAsync(CancellationToken cancellationToken = default) |
|||
{ |
|||
foreach (var processor in Processors) |
|||
{ |
|||
await processor.StopAsync(cancellationToken); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
@ -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<InboxProcessor> 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<InboxProcessor>.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?)
|
|||
} |
|||
} |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,155 @@ |
|||
// <auto-generated />
|
|||
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<Guid>("Id") |
|||
.HasColumnType("uniqueidentifier"); |
|||
|
|||
b.Property<string>("ConcurrencyStamp") |
|||
.IsConcurrencyToken() |
|||
.HasMaxLength(40) |
|||
.HasColumnType("nvarchar(40)") |
|||
.HasColumnName("ConcurrencyStamp"); |
|||
|
|||
b.Property<DateTime>("CreationTime") |
|||
.HasColumnType("datetime2") |
|||
.HasColumnName("CreationTime"); |
|||
|
|||
b.Property<Guid?>("CreatorId") |
|||
.HasColumnType("uniqueidentifier") |
|||
.HasColumnName("CreatorId"); |
|||
|
|||
b.Property<string>("ExtraProperties") |
|||
.HasColumnType("nvarchar(max)") |
|||
.HasColumnName("ExtraProperties"); |
|||
|
|||
b.Property<string>("Text") |
|||
.IsRequired() |
|||
.HasMaxLength(128) |
|||
.HasColumnType("nvarchar(128)"); |
|||
|
|||
b.HasKey("Id"); |
|||
|
|||
b.ToTable("TodoItems"); |
|||
}); |
|||
|
|||
modelBuilder.Entity("DistDemoApp.TodoSummary", b => |
|||
{ |
|||
b.Property<int>("Id") |
|||
.ValueGeneratedOnAdd() |
|||
.HasColumnType("int") |
|||
.HasAnnotation("SqlServer:ValueGenerationStrategy", SqlServerValueGenerationStrategy.IdentityColumn); |
|||
|
|||
b.Property<string>("ConcurrencyStamp") |
|||
.IsConcurrencyToken() |
|||
.HasMaxLength(40) |
|||
.HasColumnType("nvarchar(40)") |
|||
.HasColumnName("ConcurrencyStamp"); |
|||
|
|||
b.Property<byte>("Day") |
|||
.HasColumnType("tinyint"); |
|||
|
|||
b.Property<string>("ExtraProperties") |
|||
.HasColumnType("nvarchar(max)") |
|||
.HasColumnName("ExtraProperties"); |
|||
|
|||
b.Property<byte>("Month") |
|||
.HasColumnType("tinyint"); |
|||
|
|||
b.Property<int>("TotalCount") |
|||
.HasColumnType("int"); |
|||
|
|||
b.Property<int>("Year") |
|||
.HasColumnType("int"); |
|||
|
|||
b.HasKey("Id"); |
|||
|
|||
b.ToTable("TodoSummaries"); |
|||
}); |
|||
|
|||
modelBuilder.Entity("Volo.Abp.EntityFrameworkCore.DistributedEvents.IncomingEventRecord", b => |
|||
{ |
|||
b.Property<Guid>("Id") |
|||
.HasColumnType("uniqueidentifier"); |
|||
|
|||
b.Property<DateTime>("CreationTime") |
|||
.HasColumnType("datetime2") |
|||
.HasColumnName("CreationTime"); |
|||
|
|||
b.Property<byte[]>("EventData") |
|||
.IsRequired() |
|||
.HasColumnType("varbinary(max)"); |
|||
|
|||
b.Property<string>("EventName") |
|||
.IsRequired() |
|||
.HasMaxLength(256) |
|||
.HasColumnType("nvarchar(256)"); |
|||
|
|||
b.Property<string>("ExtraProperties") |
|||
.HasColumnType("nvarchar(max)") |
|||
.HasColumnName("ExtraProperties"); |
|||
|
|||
b.Property<bool>("Processed") |
|||
.HasColumnType("bit"); |
|||
|
|||
b.Property<DateTime?>("ProcessedTime") |
|||
.HasColumnType("datetime2"); |
|||
|
|||
b.HasKey("Id"); |
|||
|
|||
b.ToTable("AbpEventInbox"); |
|||
}); |
|||
|
|||
modelBuilder.Entity("Volo.Abp.EntityFrameworkCore.DistributedEvents.OutgoingEventRecord", b => |
|||
{ |
|||
b.Property<Guid>("Id") |
|||
.HasColumnType("uniqueidentifier"); |
|||
|
|||
b.Property<DateTime>("CreationTime") |
|||
.HasColumnType("datetime2") |
|||
.HasColumnName("CreationTime"); |
|||
|
|||
b.Property<byte[]>("EventData") |
|||
.IsRequired() |
|||
.HasColumnType("varbinary(max)"); |
|||
|
|||
b.Property<string>("EventName") |
|||
.IsRequired() |
|||
.HasMaxLength(256) |
|||
.HasColumnType("nvarchar(256)"); |
|||
|
|||
b.Property<string>("ExtraProperties") |
|||
.HasColumnType("nvarchar(max)") |
|||
.HasColumnName("ExtraProperties"); |
|||
|
|||
b.HasKey("Id"); |
|||
|
|||
b.ToTable("AbpEventOutbox"); |
|||
}); |
|||
#pragma warning restore 612, 618
|
|||
} |
|||
} |
|||
} |
|||
@ -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<bool>( |
|||
name: "Processed", |
|||
table: "AbpEventInbox", |
|||
type: "bit", |
|||
nullable: false, |
|||
defaultValue: false); |
|||
|
|||
migrationBuilder.AddColumn<DateTime>( |
|||
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"); |
|||
} |
|||
} |
|||
} |
|||
Loading…
Reference in new issue