From a66dec9a75dc828aaf124595e4c81330399efaf8 Mon Sep 17 00:00:00 2001 From: Sebastian Stehle Date: Tue, 21 Nov 2017 17:44:08 +0100 Subject: [PATCH] Dispatching simplified. --- .../Rules/MongoRuleEventEntity.cs | 4 - .../Rules/MongoRuleEventRepository.cs | 10 +- .../Orleans/Grains/IRuleDequeuerGrain.cs | 3 + .../Implementation/RuleDequeuerGrain.cs | 150 ++++++------ .../Repositories/IRuleEventRepository.cs | 2 - .../Orleans/Grains/IEventConsumerGrain.cs | 9 +- .../Implementation/EventConsumerGrain.cs | 226 +++++++++--------- .../Implementation/WrapperSubscription.cs | 48 ++++ .../Tasks/SingleThreadedDispatcher.cs | 44 +--- src/Squidex/Program.cs | 2 + .../Rules/RuleDequeuerGrainTests.cs | 63 ++--- .../Events/Grains/EventConsumerGrainTests.cs | 28 ++- 12 files changed, 277 insertions(+), 312 deletions(-) create mode 100644 src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/WrapperSubscription.cs diff --git a/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventEntity.cs b/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventEntity.cs index f060b9294..e48aa1543 100644 --- a/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventEntity.cs +++ b/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventEntity.cs @@ -38,10 +38,6 @@ namespace Squidex.Domain.Apps.Read.MongoDb.Rules [BsonElement] public Instant? NextAttempt { get; set; } - [BsonRequired] - [BsonElement] - public bool IsSending { get; set; } - [BsonRequired] [BsonElement] public RuleResult Result { get; set; } diff --git a/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventRepository.cs b/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventRepository.cs index db39b6675..99b22d125 100644 --- a/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventRepository.cs +++ b/src/Squidex.Domain.Apps.Read.MongoDb/Rules/MongoRuleEventRepository.cs @@ -36,14 +36,14 @@ namespace Squidex.Domain.Apps.Read.MongoDb.Rules protected override Task SetupCollectionAsync(IMongoCollection collection) { return Task.WhenAll( - collection.Indexes.CreateOneAsync(Index.Ascending(x => x.NextAttempt).Descending(x => x.IsSending)), + collection.Indexes.CreateOneAsync(Index.Ascending(x => x.NextAttempt)), collection.Indexes.CreateOneAsync(Index.Ascending(x => x.AppId).Descending(x => x.Created)), collection.Indexes.CreateOneAsync(Index.Ascending(x => x.Expires), new CreateIndexOptions { ExpireAfter = TimeSpan.Zero })); } public Task QueryPendingAsync(Instant now, Func callback, CancellationToken cancellationToken = default(CancellationToken)) { - return Collection.Find(x => x.NextAttempt < now && !x.IsSending).ForEachAsync(callback, cancellationToken); + return Collection.Find(x => x.NextAttempt < now).ForEachAsync(callback, cancellationToken); } public async Task> QueryByAppAsync(Guid appId, int skip = 0, int take = 20) @@ -81,18 +81,12 @@ namespace Squidex.Domain.Apps.Read.MongoDb.Rules return Collection.InsertOneIfNotExistsAsync(entity); } - public Task MarkSendingAsync(Guid jobId) - { - return Collection.UpdateOneAsync(x => x.Id == jobId, Update.Set(x => x.IsSending, true)); - } - public Task MarkSentAsync(Guid jobId, string dump, RuleResult result, RuleJobResult jobResult, TimeSpan elapsed, Instant? nextAttempt) { return Collection.UpdateOneAsync(x => x.Id == jobId, Update.Set(x => x.Result, result) .Set(x => x.LastDump, dump) .Set(x => x.JobResult, jobResult) - .Set(x => x.IsSending, false) .Set(x => x.NextAttempt, nextAttempt) .Inc(x => x.NumCalls, 1)); } diff --git a/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/IRuleDequeuerGrain.cs b/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/IRuleDequeuerGrain.cs index bf30781c0..eed563b25 100644 --- a/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/IRuleDequeuerGrain.cs +++ b/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/IRuleDequeuerGrain.cs @@ -8,11 +8,14 @@ using System.Threading.Tasks; using Orleans; +using Orleans.Concurrency; namespace Squidex.Domain.Apps.Read.Rules.Orleans.Grains { public interface IRuleDequeuerGrain : IGrainWithStringKey { Task ActivateAsync(); + + Task HandleAsync(Immutable @event); } } diff --git a/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/Implementation/RuleDequeuerGrain.cs b/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/Implementation/RuleDequeuerGrain.cs index 0e02f503e..2880b1c42 100644 --- a/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/Implementation/RuleDequeuerGrain.cs +++ b/src/Squidex.Domain.Apps.Read/Rules/Orleans/Grains/Implementation/RuleDequeuerGrain.cs @@ -7,14 +7,15 @@ // ========================================================================== using System; -using System.Threading; +using System.Collections.Generic; using System.Threading.Tasks; -using System.Threading.Tasks.Dataflow; using NodaTime; using Orleans; +using Orleans.Concurrency; using Orleans.Core; using Orleans.Runtime; using Squidex.Domain.Apps.Core.HandleRules; +using Squidex.Domain.Apps.Core.Rules; using Squidex.Domain.Apps.Read.Rules.Repositories; using Squidex.Infrastructure; using Squidex.Infrastructure.Log; @@ -22,14 +23,15 @@ using Squidex.Infrastructure.Tasks; namespace Squidex.Domain.Apps.Read.Rules.Orleans.Grains.Implementation { + [Reentrant] public class RuleDequeuerGrain : Grain, IRuleDequeuerGrain, IRemindable { private readonly IRuleEventRepository ruleEventRepository; private readonly RuleService ruleService; private readonly IClock clock; private readonly ISemanticLog log; - private ActionBlock requestBlock; - private TransformBlock blockBlock; + private readonly HashSet executing = new HashSet(); + private TaskFactory scheduler; public RuleDequeuerGrain(RuleService ruleService, IRuleEventRepository ruleEventRepository, ISemanticLog log, IClock clock) : this(ruleService, ruleEventRepository, log, clock, null, null) @@ -56,39 +58,19 @@ namespace Squidex.Domain.Apps.Read.Rules.Orleans.Grains.Implementation public override Task OnActivateAsync() { + scheduler = new TaskFactory(TaskScheduler.Current ?? TaskScheduler.Default); + DelayDeactivation(TimeSpan.FromDays(1)); RegisterOrUpdateReminder("Default", TimeSpan.Zero, TimeSpan.FromMinutes(10)); - RegisterTimer(x => QueryAsync(), null, TimeSpan.Zero, TimeSpan.FromSeconds(10)); - - requestBlock = - new ActionBlock(MakeRequestAsync, - new ExecutionDataflowBlockOptions - { - TaskScheduler = TaskScheduler.Current, - MaxMessagesPerTask = 1, - MaxDegreeOfParallelism = 32, - BoundedCapacity = 32 - }); - - blockBlock = - new TransformBlock(x => BlockAsync(x), - new ExecutionDataflowBlockOptions - { - TaskScheduler = TaskScheduler.Current, - MaxMessagesPerTask = 1, - MaxDegreeOfParallelism = 1, - BoundedCapacity = 1 - }); - - blockBlock.LinkTo(requestBlock, new DataflowLinkOptions { PropagateCompletion = true }); + RegisterTimer(x => QueryAsync(), null, TimeSpan.Zero, TimeSpan.FromSeconds(2)); return base.OnActivateAsync(); } public Task ReceiveReminder(string reminderName, TickStatus status) { - return QueryAsync(); + return TaskHelper.Done; } public Task ActivateAsync() @@ -100,9 +82,14 @@ namespace Squidex.Domain.Apps.Read.Rules.Orleans.Grains.Implementation { try { - var now = clock.GetCurrentInstant(); + var self = GetSelf(); + + await ruleEventRepository.QueryPendingAsync(clock.GetCurrentInstant(), x => + { + scheduler.StartNew(() => self.HandleAsync(x.AsImmutable()).Forget()).Forget(); - await ruleEventRepository.QueryPendingAsync(now, blockBlock.SendAsync, CancellationToken.None); + return TaskHelper.Done; + }); } catch (Exception ex) { @@ -112,78 +99,75 @@ namespace Squidex.Domain.Apps.Read.Rules.Orleans.Grains.Implementation } } - private async Task BlockAsync(IRuleEventEntity @event) + public async Task HandleAsync(Immutable @event) { + if (!executing.Add(@event.Value.Id)) + { + return; + } + try { - await ruleEventRepository.MarkSendingAsync(@event.Id); + var job = @event.Value.Job; + + var response = await ruleService.InvokeAsync(job.ActionName, job.ActionData); - return @event; + var jobInvoke = ComputeJobInvoke(response.Result, @event, job); + var jobResult = ComputeJobResult(response.Result, jobInvoke); + + await ruleEventRepository.MarkSentAsync(@event.Value.Id, response.Dump, response.Result, jobResult, response.Elapsed, jobInvoke); } catch (Exception ex) { log.LogError(ex, w => w - .WriteProperty("action", "BlockWebhookEvent") + .WriteProperty("action", "SendWebhookEvent") .WriteProperty("status", "Failed")); - - throw; + } + finally + { + executing.Remove(@event.Value.Id); } } - private async Task MakeRequestAsync(IRuleEventEntity @event) + private static RuleJobResult ComputeJobResult(RuleResult result, Instant? nextCall) { - try + if (result != RuleResult.Success && !nextCall.HasValue) { - var job = @event.Job; - - var response = await ruleService.InvokeAsync(job.ActionName, job.ActionData); - - Instant? nextCall = null; - - if (response.Result != RuleResult.Success) - { - switch (@event.NumCalls) - { - case 0: - nextCall = job.Created.Plus(Duration.FromMinutes(5)); - break; - case 1: - nextCall = job.Created.Plus(Duration.FromHours(1)); - break; - case 2: - nextCall = job.Created.Plus(Duration.FromHours(6)); - break; - case 3: - nextCall = job.Created.Plus(Duration.FromHours(12)); - break; - } - } - - RuleJobResult jobResult; + return RuleJobResult.Failed; + } + else if (result != RuleResult.Success && nextCall.HasValue) + { + return RuleJobResult.Retry; + } + else + { + return RuleJobResult.Success; + } + } - if (response.Result != RuleResult.Success && !nextCall.HasValue) - { - jobResult = RuleJobResult.Failed; - } - else if (response.Result != RuleResult.Success && nextCall.HasValue) - { - jobResult = RuleJobResult.Retry; - } - else + private static Instant? ComputeJobInvoke(RuleResult result, Immutable @event, RuleJob job) + { + if (result != RuleResult.Success) + { + switch (@event.Value.NumCalls) { - jobResult = RuleJobResult.Success; + case 0: + return job.Created.Plus(Duration.FromMinutes(5)); + case 1: + return job.Created.Plus(Duration.FromHours(1)); + case 2: + return job.Created.Plus(Duration.FromHours(6)); + case 3: + return job.Created.Plus(Duration.FromHours(12)); } - - await ruleEventRepository.MarkSentAsync(@event.Id, response.Dump, response.Result, jobResult, response.Elapsed, nextCall); } - catch (Exception ex) - { - log.LogError(ex, w => w - .WriteProperty("action", "SendWebhookEvent") - .WriteProperty("status", "Failed")); - throw; - } + return null; + } + + protected virtual IRuleDequeuerGrain GetSelf() + { + return this.AsReference(); } } } diff --git a/src/Squidex.Domain.Apps.Read/Rules/Repositories/IRuleEventRepository.cs b/src/Squidex.Domain.Apps.Read/Rules/Repositories/IRuleEventRepository.cs index 256aa9b71..2b05e6741 100644 --- a/src/Squidex.Domain.Apps.Read/Rules/Repositories/IRuleEventRepository.cs +++ b/src/Squidex.Domain.Apps.Read/Rules/Repositories/IRuleEventRepository.cs @@ -22,8 +22,6 @@ namespace Squidex.Domain.Apps.Read.Rules.Repositories Task EnqueueAsync(Guid id, Instant nextAttempt); - Task MarkSendingAsync(Guid jobId); - Task MarkSentAsync(Guid jobId, string dump, RuleResult result, RuleJobResult jobResult, TimeSpan elapsed, Instant? nextCall); Task QueryPendingAsync(Instant now, Func callback, CancellationToken cancellationToken = default(CancellationToken)); diff --git a/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/IEventConsumerGrain.cs b/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/IEventConsumerGrain.cs index 159155bcf..b554a2fc6 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/IEventConsumerGrain.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/IEventConsumerGrain.cs @@ -6,13 +6,14 @@ // All rights reserved. // ========================================================================== +using System; using System.Threading.Tasks; using Orleans; using Orleans.Concurrency; namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains { - public interface IEventConsumerGrain : IGrainWithStringKey, IEventSubscriber + public interface IEventConsumerGrain : IGrainWithStringKey { Task> GetStateAsync(); @@ -23,5 +24,11 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains Task StartAsync(); Task ResetAsync(); + + Task OnEventAsync(Immutable subscription, Immutable storedEvent); + + Task OnErrorAsync(Immutable subscription, Immutable exception); + + Task OnClosedAsync(Immutable subscription); } } diff --git a/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/EventConsumerGrain.cs b/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/EventConsumerGrain.cs index 79dd419b1..0d6570280 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/EventConsumerGrain.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/EventConsumerGrain.cs @@ -18,15 +18,20 @@ using Squidex.Infrastructure.Tasks; namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation { - public class EventConsumerGrain : GrainV2, IEventSubscriber, IEventConsumerGrain + public class EventConsumerGrain : GrainV2, IEventConsumerGrain { private readonly EventDataFormatter eventFormatter; private readonly EventConsumerFactory eventConsumerFactory; private readonly IEventStore eventStore; private readonly ISemanticLog log; - private IEventSubscription currentSubscription; + private TaskScheduler scheduler; private IEventConsumer eventConsumer; - private SingleThreadedDispatcher dispatcher; + private IEventSubscription eventSubscription; + + protected IEventStore EventStore + { + get { return eventStore; } + } public EventConsumerGrain( EventDataFormatter eventFormatter, @@ -62,7 +67,7 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation public override Task OnActivateAsync() { - dispatcher = new SingleThreadedDispatcher(1, TaskScheduler.Current); + scheduler = TaskScheduler.Current; eventConsumer = eventConsumerFactory(this.GetPrimaryKeyString()); @@ -71,53 +76,32 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation public Task ActivateAsync() { - return dispatcher.DispatchAndUnwrapAsync(() => - { - if (!State.IsStopped) - { - Subscribe(State.Position); - } - - return TaskHelper.Done; - }); - } - - private Task HandleEventAsync(IEventSubscription subscription, StoredEvent storedEvent) - { - if (subscription != currentSubscription) + if (!State.IsStopped) { - return TaskHelper.Done; + Subscribe(State.Position); } - return DoAndUpdateStateAsync(async () => - { - var @event = ParseKnownEvent(storedEvent); - - if (@event != null) - { - await DispatchConsumerAsync(@event); - } - - State = EventConsumerGrainState.Handled(storedEvent.EventPosition); - }); + return TaskHelper.Done; } - private Task HandleClosedAsync(IEventSubscription subscription) + public Task StartAsync() { - if (subscription != currentSubscription) + if (!State.IsStopped) { return TaskHelper.Done; } return DoAndUpdateStateAsync(() => { - Unsubscribe(); + Subscribe(State.Position); + + State = State.Started(); }); } - private Task HandleErrorAsync(IEventSubscription subscription, Exception exception) + public Task StopAsync() { - if (subscription != currentSubscription) + if (State.IsStopped) { return TaskHelper.Done; } @@ -126,76 +110,70 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation { Unsubscribe(); - State = State.Failed(exception); + State = State.Stopped(); }); } - public Task StartAsync() + public Task ResetAsync() { - return dispatcher.DispatchAndUnwrapAsync(() => + return DoAndUpdateStateAsync(async () => { - if (!State.IsStopped) - { - return TaskHelper.Done; - } + Unsubscribe(); - return DoAndUpdateStateAsync(() => - { - Subscribe(State.Position); + await ClearAsync(); + + Subscribe(null); - State = State.Started(); - }); + State = EventConsumerGrainState.Initial(); }); } - public Task StopAsync() + public Task OnEventAsync(Immutable subscription, Immutable storedEvent) { - return dispatcher.DispatchAndUnwrapAsync(() => + if (subscription.Value != eventSubscription) { - if (State.IsStopped) - { - return TaskHelper.Done; - } + return TaskHelper.Done; + } - return DoAndUpdateStateAsync(() => + return DoAndUpdateStateAsync(async () => + { + var @event = ParseKnownEvent(storedEvent.Value); + + if (@event != null) { - Unsubscribe(); + await DispatchConsumerAsync(@event); + } - State = State.Stopped(); - }); + State = EventConsumerGrainState.Handled(storedEvent.Value.EventPosition); }); } - public Task ResetAsync() + public Task OnErrorAsync(Immutable subscription, Immutable exception) { - return dispatcher.DispatchAndUnwrapAsync(() => + if (subscription.Value != eventSubscription) { - return DoAndUpdateStateAsync(async () => - { - Unsubscribe(); - - await ClearAsync(); + return TaskHelper.Done; + } - Subscribe(null); + return DoAndUpdateStateAsync(() => + { + Unsubscribe(); - State = EventConsumerGrainState.Initial(); - }); + State = State.Failed(exception.Value); }); } - Task IEventSubscriber.OnEventAsync(IEventSubscription subscription, StoredEvent storedEvent) + public Task OnClosedAsync(Immutable subscription) { - return dispatcher.DispatchAndUnwrapAsync(() => HandleEventAsync(subscription, storedEvent)); - } - - Task IEventSubscriber.OnErrorAsync(IEventSubscription subscription, Exception exception) - { - return dispatcher.DispatchAndUnwrapAsync(() => HandleErrorAsync(subscription, exception)); - } + if (subscription.Value != eventSubscription) + { + return TaskHelper.Done; + } - Task IEventSubscriber.OnClosedAsync(IEventSubscription subscription) - { - return dispatcher.DispatchAndUnwrapAsync(() => HandleClosedAsync(subscription)); + return DoAndUpdateStateAsync(() => + { + Unsubscribe(); + }); } public Task> GetStateAsync() @@ -203,39 +181,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation return Task.FromResult(new Immutable(State.ToInfo(this.GetPrimaryKeyString()))); } - private Task DoAndUpdateStateAsync(Action action) - { - return DoAndUpdateStateAsync(() => { action(); return TaskHelper.Done; }); - } - - private async Task DoAndUpdateStateAsync(Func action) - { - try - { - await action(); - } - catch (Exception ex) - { - try - { - Unsubscribe(); - } - catch (Exception unsubscribeException) - { - ex = new AggregateException(ex, unsubscribeException); - } - - log.LogFatal(ex, w => w - .WriteProperty("action", "HandleEvent") - .WriteProperty("state", "Failed") - .WriteProperty("eventConsumer", eventConsumer.Name)); - - State = State.Failed(ex); - } - - await WriteStateAsync(); - } - private async Task ClearAsync() { var actionId = Guid.NewGuid().ToString(); @@ -283,19 +228,19 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation private void Unsubscribe() { - if (currentSubscription != null) + if (eventSubscription != null) { - currentSubscription.StopAsync().Forget(); - currentSubscription = null; + eventSubscription.StopAsync().Forget(); + eventSubscription = null; } } private void Subscribe(string position) { - if (currentSubscription == null) + if (eventSubscription == null) { - currentSubscription?.StopAsync().Forget(); - currentSubscription = CreateSubscription(eventStore, eventConsumer.EventsFilter, position); + eventSubscription?.StopAsync().Forget(); + eventSubscription = CreateSubscription(eventConsumer.EventsFilter, position); } } @@ -318,9 +263,52 @@ namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation } } - protected virtual IEventSubscription CreateSubscription(IEventStore eventStore, string streamFilter, string position) + protected virtual IEventConsumerGrain GetSelf() + { + return this.AsReference(); + } + + protected virtual IEventSubscription CreateSubscription(IEventSubscriber subscriber, string streamFilter, string position) + { + return new RetrySubscription(EventStore, subscriber, streamFilter, position); + } + + private IEventSubscription CreateSubscription(string streamFilter, string position) + { + return CreateSubscription(new WrapperSubscription(GetSelf(), scheduler), streamFilter, position); + } + + private Task DoAndUpdateStateAsync(Action action) + { + return DoAndUpdateStateAsync(() => { action(); return TaskHelper.Done; }); + } + + private async Task DoAndUpdateStateAsync(Func action) { - return new RetrySubscription(eventStore, this, streamFilter, position); + try + { + await action(); + } + catch (Exception ex) + { + try + { + Unsubscribe(); + } + catch (Exception unsubscribeException) + { + ex = new AggregateException(ex, unsubscribeException); + } + + log.LogFatal(ex, w => w + .WriteProperty("action", "HandleEvent") + .WriteProperty("state", "Failed") + .WriteProperty("eventConsumer", eventConsumer.Name)); + + State = State.Failed(ex); + } + + await WriteStateAsync(); } } } \ No newline at end of file diff --git a/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/WrapperSubscription.cs b/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/WrapperSubscription.cs new file mode 100644 index 000000000..e1ebe67f8 --- /dev/null +++ b/src/Squidex.Infrastructure/CQRS/Events/Orleans/Grains/Implementation/WrapperSubscription.cs @@ -0,0 +1,48 @@ +// ========================================================================== +// WrapperSubscription.cs +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex Group +// All rights reserved. +// ========================================================================== + +using System; +using System.Threading; +using System.Threading.Tasks; +using Orleans.Concurrency; + +namespace Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation +{ + internal sealed class WrapperSubscription : IEventSubscriber + { + private readonly IEventConsumerGrain grain; + private readonly TaskScheduler scheduler; + + public WrapperSubscription(IEventConsumerGrain grain, TaskScheduler scheduler) + { + this.grain = grain; + + this.scheduler = scheduler ?? TaskScheduler.Default; + } + + public Task OnEventAsync(IEventSubscription subscription, StoredEvent storedEvent) + { + return Dispatch(() => grain.OnEventAsync(subscription.AsImmutable(), storedEvent.AsImmutable())); + } + + public Task OnErrorAsync(IEventSubscription subscription, Exception exception) + { + return Dispatch(() => grain.OnErrorAsync(subscription.AsImmutable(), exception.AsImmutable())); + } + + public Task OnClosedAsync(IEventSubscription subscription) + { + return Dispatch(() => grain.OnClosedAsync(subscription.AsImmutable())); + } + + private Task Dispatch(Func task) + { + return Task.Factory.StartNew(() => task(), CancellationToken.None, TaskCreationOptions.None, scheduler).Unwrap(); + } + } +} diff --git a/src/Squidex.Infrastructure/Tasks/SingleThreadedDispatcher.cs b/src/Squidex.Infrastructure/Tasks/SingleThreadedDispatcher.cs index e1558acb0..c382270aa 100644 --- a/src/Squidex.Infrastructure/Tasks/SingleThreadedDispatcher.cs +++ b/src/Squidex.Infrastructure/Tasks/SingleThreadedDispatcher.cs @@ -17,59 +17,20 @@ namespace Squidex.Infrastructure.Tasks private readonly ActionBlock> block; private bool isStopped; - public SingleThreadedDispatcher(int capacity = 1, TaskScheduler scheduler = null) + public SingleThreadedDispatcher(int capacity = 1) { var options = new ExecutionDataflowBlockOptions { BoundedCapacity = capacity, MaxMessagesPerTask = -1, - MaxDegreeOfParallelism = 1, - TaskScheduler = scheduler ?? TaskScheduler.Default + MaxDegreeOfParallelism = 1 }; block = new ActionBlock>(Handle, options); } - public Task DispatchAndUnwrapAsync(Func action) - { - return Task.CompletedTask; - Guard.NotNull(action, nameof(action)); - - var cts = new TaskCompletionSource(); - - block.SendAsync(async () => - { - try - { - await action(); - - cts.SetResult(true); - } - catch (Exception ex) - { - cts.TrySetException(ex); - } - }); - - return cts.Task; - } - - public Task DispatchAndUnwrapAsync(Action action) - { - return Task.CompletedTask; - Guard.NotNull(action, nameof(action)); - - return DispatchAndUnwrapAsync(() => - { - action(); - - return TaskHelper.Done; - }); - } - public Task DispatchAsync(Func action) { - return Task.CompletedTask; Guard.NotNull(action, nameof(action)); return block.SendAsync(action); @@ -77,7 +38,6 @@ namespace Squidex.Infrastructure.Tasks public Task DispatchAsync(Action action) { - return Task.CompletedTask; Guard.NotNull(action, nameof(action)); return block.SendAsync(() => { action(); return TaskHelper.Done; }); diff --git a/src/Squidex/Program.cs b/src/Squidex/Program.cs index 950d34614..d0a45c6bc 100644 --- a/src/Squidex/Program.cs +++ b/src/Squidex/Program.cs @@ -10,6 +10,8 @@ using System; using System.IO; using Microsoft.AspNetCore.Hosting; using Orleans; +using Orleans.Hosting; +using Orleans.Runtime.Configuration; using Squidex.Config.Orleans; using Squidex.Domain.Apps.Read.State.Orleans.Grains.Implementations; using Squidex.Domain.Users.DataProtection.Orleans.Grains.Implementations; diff --git a/tests/Squidex.Domain.Apps.Read.Tests/Rules/RuleDequeuerGrainTests.cs b/tests/Squidex.Domain.Apps.Read.Tests/Rules/RuleDequeuerGrainTests.cs index 0cd9684b1..a98b27a2b 100644 --- a/tests/Squidex.Domain.Apps.Read.Tests/Rules/RuleDequeuerGrainTests.cs +++ b/tests/Squidex.Domain.Apps.Read.Tests/Rules/RuleDequeuerGrainTests.cs @@ -7,14 +7,15 @@ // ========================================================================== using System; -using System.Threading; using System.Threading.Tasks; using FakeItEasy; using NodaTime; +using Orleans.Concurrency; using Orleans.Core; using Orleans.Runtime; using Squidex.Domain.Apps.Core.HandleRules; using Squidex.Domain.Apps.Core.Rules; +using Squidex.Domain.Apps.Read.Rules.Orleans.Grains; using Squidex.Domain.Apps.Read.Rules.Orleans.Grains.Implementation; using Squidex.Domain.Apps.Read.Rules.Repositories; using Squidex.Infrastructure.Log; @@ -31,6 +32,7 @@ namespace Squidex.Domain.Apps.Read.Rules private readonly IAppProvider appProvider = A.Fake(); private readonly IRuleEventRepository ruleEventRepository = A.Fake(); private readonly RuleService ruleService = A.Fake(); + private readonly MyRuleDequeuerGrain sut; private readonly Instant now = SystemClock.Instance.GetCurrentInstant(); public sealed class MyRuleDequeuerGrain : RuleDequeuerGrain @@ -41,11 +43,24 @@ namespace Squidex.Domain.Apps.Read.Rules : base(ruleService, ruleEventRepository, log, clock, identity, runtime) { } + + protected override IRuleDequeuerGrain GetSelf() + { + return this; + } } public RuleDequeuerGrainTests() { A.CallTo(() => clock.GetCurrentInstant()).Returns(now); + + sut = new MyRuleDequeuerGrain( + ruleService, + ruleEventRepository, + log, + clock, + A.Fake(), + A.Fake()); } [Theory] @@ -65,20 +80,8 @@ namespace Squidex.Domain.Apps.Read.Rules var requestElapsed = TimeSpan.FromMinutes(1); var requestDump = "Dump"; - SetupSender(@event, requestDump, result, requestElapsed); - SetupPendingEvents(@event); - - var sut = new MyRuleDequeuerGrain( - ruleService, - ruleEventRepository, - log, - clock, - A.Fake(), - A.Fake()); - - await sut.OnActivateAsync(); - await sut.QueryAsync(); - await sut.OnDeactivateAsync(); + A.CallTo(() => ruleService.InvokeAsync(@event.Job.ActionName, @event.Job.ActionData)) + .Returns((requestDump, result, requestElapsed)); Instant? nextCall = null; @@ -87,33 +90,11 @@ namespace Squidex.Domain.Apps.Read.Rules nextCall = now.Plus(Duration.FromMinutes(minutes)); } - VerifyRepositories(@event, requestDump, result, jobResult, requestElapsed, nextCall); - } - - private void SetupSender(IRuleEventEntity @event, string requestDump, RuleResult requestResult, TimeSpan requestTime) - { - A.CallTo(() => ruleService.InvokeAsync(@event.Job.ActionName, @event.Job.ActionData)) - .Returns((requestDump, requestResult, requestTime)); - } - - private void SetupPendingEvents(IRuleEventEntity @event) - { - A.CallTo(() => ruleEventRepository.QueryPendingAsync( - now, - A>.Ignored, - A.Ignored)) - .Invokes(async (Instant n, Func callback, CancellationToken ct) => - { - await callback(@event); - }); - } - - private void VerifyRepositories(IRuleEventEntity @event, string dump, RuleResult result, RuleJobResult jobResult, TimeSpan elapsed, Instant? nextCall) - { - A.CallTo(() => ruleEventRepository.MarkSendingAsync(@event.Id)) - .MustHaveHappened(); + await sut.OnActivateAsync(); + await sut.HandleAsync(@event.AsImmutable()); + await sut.OnDeactivateAsync(); - A.CallTo(() => ruleEventRepository.MarkSentAsync(@event.Id, dump, result, jobResult, elapsed, nextCall)) + A.CallTo(() => ruleEventRepository.MarkSentAsync(@event.Id, requestDump, result, jobResult, requestElapsed, nextCall)) .MustHaveHappened(); } diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Events/Grains/EventConsumerGrainTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Events/Grains/EventConsumerGrainTests.cs index 79f631a4c..780a7ec02 100644 --- a/tests/Squidex.Infrastructure.Tests/CQRS/Events/Grains/EventConsumerGrainTests.cs +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Events/Grains/EventConsumerGrainTests.cs @@ -9,8 +9,10 @@ using System; using System.Threading.Tasks; using FakeItEasy; +using Orleans.Concurrency; using Orleans.Core; using Orleans.Runtime; +using Squidex.Infrastructure.CQRS.Events.Orleans.Grains; using Squidex.Infrastructure.CQRS.Events.Orleans.Grains.Implementation; using Squidex.Infrastructure.Log; using Xunit; @@ -23,9 +25,9 @@ namespace Squidex.Infrastructure.CQRS.Events.Grains { } - public sealed class MyEventConsumerActor : EventConsumerGrain + public sealed class MyEventConsumerGrain : EventConsumerGrain { - public MyEventConsumerActor( + public MyEventConsumerGrain( EventDataFormatter formatter, EventConsumerFactory eventConsumerFactory, IEventStore eventStore, @@ -37,9 +39,14 @@ namespace Squidex.Infrastructure.CQRS.Events.Grains { } - protected override IEventSubscription CreateSubscription(IEventStore eventStore, string streamFilter, string position) + protected override IEventConsumerGrain GetSelf() { - return eventStore.CreateSubscription(this, streamFilter, position); + return this; + } + + protected override IEventSubscription CreateSubscription(IEventSubscriber subscriber, string streamFilter, string position) + { + return EventStore.CreateSubscription(subscriber, streamFilter, position); } } @@ -47,13 +54,12 @@ namespace Squidex.Infrastructure.CQRS.Events.Grains private readonly IEventStore eventStore = A.Fake(); private readonly IEventSubscription eventSubscription = A.Fake(); private readonly ISemanticLog log = A.Fake(); - private readonly IEventSubscriber sutSubscriber; private readonly IStorage storage = A.Fake>(); private readonly EventDataFormatter formatter = A.Fake(); private readonly EventData eventData = new EventData(); private readonly Envelope envelope = new Envelope(new MyEvent()); private readonly EventConsumerFactory factory; - private readonly MyEventConsumerActor sut; + private readonly MyEventConsumerGrain sut; private readonly string consumerName; private EventConsumerGrainState state = new EventConsumerGrainState(); @@ -72,7 +78,7 @@ namespace Squidex.Infrastructure.CQRS.Events.Grains A.CallTo(() => storage.State).ReturnsLazily(() => state); A.CallToSet(() => storage.State).Invokes(new Action(s => state = s)); - sut = new MyEventConsumerActor( + sut = new MyEventConsumerGrain( formatter, factory, eventStore, @@ -80,8 +86,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Grains A.Fake(), A.Fake(), storage); - - sutSubscriber = sut; } [Fact] @@ -388,17 +392,17 @@ namespace Squidex.Infrastructure.CQRS.Events.Grains private Task OnErrorAsync(IEventSubscription subscriber, Exception ex) { - return sutSubscriber.OnErrorAsync(subscriber, ex); + return sut.OnErrorAsync(subscriber.AsImmutable(), ex.AsImmutable()); } private Task OnEventAsync(IEventSubscription subscriber, StoredEvent ev) { - return sutSubscriber.OnEventAsync(subscriber, ev); + return sut.OnEventAsync(subscriber.AsImmutable(), ev.AsImmutable()); } private Task OnClosedAsync(IEventSubscription subscriber) { - return sutSubscriber.OnClosedAsync(subscriber); + return sut.OnClosedAsync(subscriber.AsImmutable()); } } } \ No newline at end of file