From fc7f3992741dd63d452b350a1a30ba6553f43d53 Mon Sep 17 00:00:00 2001 From: Sebastian Stehle Date: Fri, 7 Apr 2017 21:14:09 +0200 Subject: [PATCH] 1) EventReceiver simplified. 2) Filter EventConsumers with regex. (#2) --- .../EventStore/MongoEventStore.cs | 18 +- .../CQRS/Events/CompoundEventConsumer.cs | 7 +- .../CQRS/Events/EventReceiver.cs | 177 +++++++++++++----- .../CQRS/Events/IEventConsumer.cs | 2 + .../CQRS/Events/IEventStore.cs | 4 +- .../Events/Internal/DispatchEventBlock.cs | 115 ------------ .../Events/Internal/IEventReceiverBlock.cs | 21 --- .../CQRS/Events/Internal/ParseEventBlock.cs | 98 ---------- .../CQRS/Events/Internal/QueryEventsBlock.cs | 161 ---------------- .../CQRS/Events/Internal/UpdateStateBlock.cs | 86 --------- .../Squidex.Infrastructure.csproj | 1 - .../Apps/MongoAppRepository_EventHandling.cs | 5 + .../MongoContentRepository_EventHandling.cs | 5 + .../History/MongoHistoryEventRepository.cs | 15 +- .../MongoSchemaRepository_EventHandling.cs | 5 + .../Implementations/CachingAppProvider.cs | 5 + .../Implementations/CachingSchemaProvider.cs | 5 + .../CQRS/Events/CompoundEventConsumerTests.cs | 18 +- .../CQRS/Events/EventReceiverTests.cs | 6 +- 19 files changed, 207 insertions(+), 547 deletions(-) delete mode 100644 src/Squidex.Infrastructure/CQRS/Events/Internal/DispatchEventBlock.cs delete mode 100644 src/Squidex.Infrastructure/CQRS/Events/Internal/IEventReceiverBlock.cs delete mode 100644 src/Squidex.Infrastructure/CQRS/Events/Internal/ParseEventBlock.cs delete mode 100644 src/Squidex.Infrastructure/CQRS/Events/Internal/QueryEventsBlock.cs delete mode 100644 src/Squidex.Infrastructure/CQRS/Events/Internal/UpdateStateBlock.cs diff --git a/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs b/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs index 84ea77e20..7e6dac229 100644 --- a/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs +++ b/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs @@ -18,6 +18,7 @@ using NodaTime; using Squidex.Infrastructure.CQRS.Events; using Squidex.Infrastructure.Reflection; +// ReSharper disable ConvertIfStatementToConditionalTernaryExpression // ReSharper disable ClassNeverInstantiated.Local // ReSharper disable UnusedMember.Local // ReSharper disable InvertIf @@ -64,7 +65,7 @@ namespace Squidex.Infrastructure.MongoDb.EventStore eventsOffsetIndex = indexNames[0]; } - public IObservable GetEventsAsync(string streamName, long lastReceivedEventNumber = -1) + public IObservable GetEventsAsync(string streamFilter, long lastReceivedEventNumber = -1) { return Observable.Create((observer, ct) => { @@ -73,11 +74,11 @@ namespace Squidex.Infrastructure.MongoDb.EventStore observer.OnNext(storedEvent); return Tasks.TaskHelper.Done; - }, ct, streamName, lastReceivedEventNumber); + }, ct, streamFilter, lastReceivedEventNumber); }); } - public async Task GetEventsAsync(Func callback, CancellationToken cancellationToken, string streamName = null, long lastReceivedEventNumber = -1) + public async Task GetEventsAsync(Func callback, CancellationToken cancellationToken, string streamFilter = null, long lastReceivedEventNumber = -1) { Guard.NotNull(callback, nameof(callback)); @@ -90,9 +91,16 @@ namespace Squidex.Infrastructure.MongoDb.EventStore filters.Add(Filter.Gte(x => x.EventsOffset, commitOffset)); } - if (!string.IsNullOrWhiteSpace(streamName)) + if (!string.IsNullOrWhiteSpace(streamFilter) && !string.Equals(streamFilter, "*", StringComparison.OrdinalIgnoreCase)) { - filters.Add(Filter.Eq(x => x.EventStream, streamName)); + if (streamFilter.StartsWith("^")) + { + filters.Add(Filter.Regex(x => x.EventStream, streamFilter)); + } + else + { + filters.Add(Filter.Eq(x => x.EventStream, streamFilter)); + } } FilterDefinition filter = new BsonDocument(); diff --git a/src/Squidex.Infrastructure/CQRS/Events/CompoundEventConsumer.cs b/src/Squidex.Infrastructure/CQRS/Events/CompoundEventConsumer.cs index 83c1843ba..11d044079 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/CompoundEventConsumer.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/CompoundEventConsumer.cs @@ -17,6 +17,11 @@ namespace Squidex.Infrastructure.CQRS.Events public string Name { get; } + public string StreamFilter + { + get { return inners.FirstOrDefault()?.StreamFilter; } + } + public CompoundEventConsumer(IEventConsumer first, params IEventConsumer[] inners) { Guard.NotNull(first, nameof(first)); @@ -24,7 +29,7 @@ namespace Squidex.Infrastructure.CQRS.Events this.inners = new[] { first }.Union(inners).ToArray(); - Name = first.GetType().Name; + Name = first.Name; } public CompoundEventConsumer(string name, params IEventConsumer[] inners) diff --git a/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs b/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs index db2848efe..cd364b6b3 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs @@ -7,12 +7,14 @@ // ========================================================================== using System; -using System.Threading; -using Squidex.Infrastructure.CQRS.Events.Internal; +using System.Reactive.Linq; +using System.Threading.Tasks; using Squidex.Infrastructure.Log; +using Squidex.Infrastructure.Timers; +// ReSharper disable MethodSupportsCancellation +// ReSharper disable ConvertIfStatementToConditionalTernaryExpression // ReSharper disable InvertIf -// ReSharper disable UseObjectOrCollectionInitializer namespace Squidex.Infrastructure.CQRS.Events { @@ -23,16 +25,11 @@ namespace Squidex.Infrastructure.CQRS.Events private readonly IEventNotifier eventNotifier; private readonly IEventConsumerInfoRepository eventConsumerInfoRepository; private readonly ISemanticLog log; - private QueryEventsBlock queryEventsBlock; - private DispatchEventBlock dispatchEventBlock; - private UpdateStateBlock updateStateBlock; - private ParseEventBlock parseEventBlock; - private IEventReceiverBlock[] blocks; - private Timer timer; + private CompletionTimer timer; public EventReceiver( EventDataFormatter formatter, - IEventStore eventStore, + IEventStore eventStore, IEventNotifier eventNotifier, IEventConsumerInfoRepository eventConsumerInfoRepository, ISemanticLog log) @@ -56,11 +53,7 @@ namespace Squidex.Infrastructure.CQRS.Events { try { - queryEventsBlock?.Complete(); - timer?.Dispose(); - - updateStateBlock?.Completion.Wait(); } catch (Exception ex) { @@ -73,71 +66,159 @@ namespace Squidex.Infrastructure.CQRS.Events public void Next() { - queryEventsBlock.NextOrThrowAway(); + timer?.Trigger(); } - public void Subscribe(IEventConsumer eventConsumer, int delay = 5000, bool autoTrigger = true) + public void Subscribe(IEventConsumer eventConsumer, int delay = 5000) { Guard.NotNull(eventConsumer, nameof(eventConsumer)); - if (updateStateBlock != null) + if (timer != null) { return; } - var onError = new Action(ex => Stop(ex, eventConsumer)); + var consumerName = eventConsumer.Name; + var consumerStarted = false; - updateStateBlock = new UpdateStateBlock(eventConsumerInfoRepository, eventConsumer, log); - updateStateBlock.OnError = onError; + timer = new CompletionTimer(delay, async ct => + { + if (!consumerStarted) + { + await eventConsumerInfoRepository.CreateAsync(consumerName); - dispatchEventBlock = new DispatchEventBlock(eventConsumer, log); - dispatchEventBlock.OnError = onError; - dispatchEventBlock.LinkTo(updateStateBlock.Target); + consumerStarted = true; + } - parseEventBlock = new ParseEventBlock(formatter, log); - parseEventBlock.OnError = onError; - parseEventBlock.LinkTo(dispatchEventBlock.Target); + try + { + var status = await eventConsumerInfoRepository.FindAsync(consumerName); - queryEventsBlock = new QueryEventsBlock(eventConsumerInfoRepository, eventConsumer, eventStore, log); - queryEventsBlock.OnEvent = parseEventBlock.NextAsync; - queryEventsBlock.OnReset = Reset; - queryEventsBlock.OnError = onError; - queryEventsBlock.Completion.ContinueWith(x => parseEventBlock.Complete()); + var lastHandledEventNumber = status.LastHandledEventNumber; - blocks = new IEventReceiverBlock[] { updateStateBlock, dispatchEventBlock, parseEventBlock, queryEventsBlock }; + if (status.IsResetting) + { + await ResetAsync(eventConsumer, consumerName); - if (autoTrigger) - { - timer = new Timer(x => queryEventsBlock.NextOrThrowAway(), null, 0, delay); - } + lastHandledEventNumber = -1; + } + else if (status.IsStopped) + { + return; + } + + await eventStore.GetEventsAsync(se => HandleEventAsync(eventConsumer, se, consumerName), ct, + eventConsumer.StreamFilter, lastHandledEventNumber); + } + catch (Exception ex) + { + log.LogFatal(ex, w => w.WriteProperty("action", "EventHandlingFailed")); + + await eventConsumerInfoRepository.StopAsync(consumerName, ex.ToString()); + } + }); - eventNotifier.Subscribe(() => queryEventsBlock.NextOrThrowAway()); + eventNotifier.Subscribe(timer.Trigger); } - private void Stop(Exception ex, IEventConsumer eventConsumer) + private async Task HandleEventAsync(IEventConsumer eventConsumer, StoredEvent storedEvent, string consumerName) { - foreach (var block in blocks) + var @event = ParseEvent(storedEvent); + + await DispatchConsumer(@event, eventConsumer); + await eventConsumerInfoRepository.SetLastHandledEventNumberAsync(consumerName, storedEvent.EventNumber); + } + + private async Task ResetAsync(IEventConsumer eventConsumer, string consumerName) + { + var actionId = Guid.NewGuid().ToString(); + try { - block.Stop(); + log.LogInformation(w => w + .WriteProperty("action", "EventConsumerReset") + .WriteProperty("actionId", actionId) + .WriteProperty("state", "Started") + .WriteProperty("eventConsumer", eventConsumer.GetType().Name)); + + await eventConsumer.ClearAsync(); + await eventConsumerInfoRepository.SetLastHandledEventNumberAsync(consumerName, -1); + + log.LogInformation(w => w + .WriteProperty("action", "EventConsumerReset") + .WriteProperty("actionId", actionId) + .WriteProperty("state", "Completed") + .WriteProperty("eventConsumer", eventConsumer.GetType().Name)); } + catch (Exception ex) + { + log.LogFatal(ex, w => w + .WriteProperty("action", "EventConsumerReset") + .WriteProperty("actionId", actionId) + .WriteProperty("state", "Completed") + .WriteProperty("eventConsumer", eventConsumer.GetType().Name)); + + throw; + } + } + private async Task DispatchConsumer(Envelope @event, IEventConsumer eventConsumer) + { + var eventId = @event.Headers.EventId().ToString(); + var eventType = @event.Payload.GetType().Name; try { - eventConsumerInfoRepository.StopAsync(eventConsumer.Name, ex.ToString()).Wait(); + log.LogInformation(w => w + .WriteProperty("action", "HandleEvent") + .WriteProperty("actionId", eventId) + .WriteProperty("state", "Started") + .WriteProperty("eventId", eventId) + .WriteProperty("eventType", eventType) + .WriteProperty("eventConsumer", eventConsumer.GetType().Name)); + + await eventConsumer.On(@event); + + log.LogInformation(w => w + .WriteProperty("action", "HandleEvent") + .WriteProperty("actionId", eventId) + .WriteProperty("state", "Completed") + .WriteProperty("eventId", eventId) + .WriteProperty("eventType", eventType) + .WriteProperty("eventConsumer", eventConsumer.GetType().Name)); } - catch (Exception ex2) + catch (Exception ex) { - log.LogFatal(ex2, w => w - .WriteProperty("action", "StopConsumer") - .WriteProperty("state", "Failed")); + log.LogError(ex, w => w + .WriteProperty("action", "HandleEvent") + .WriteProperty("actionId", eventId) + .WriteProperty("state", "Started") + .WriteProperty("eventId", eventId) + .WriteProperty("eventType", eventType) + .WriteProperty("eventConsumer", eventConsumer.GetType().Name)); + + throw; } } - private void Reset() + private Envelope ParseEvent(StoredEvent storedEvent) { - foreach (var block in blocks) + try { - block.Reset(); + var @event = formatter.Parse(storedEvent.Data); + + @event.SetEventNumber(storedEvent.EventNumber); + @event.SetEventStreamNumber(storedEvent.EventStreamNumber); + + return @event; + } + catch (Exception ex) + { + log.LogFatal(ex, w => w + .WriteProperty("action", "ParseEvent") + .WriteProperty("state", "Failed") + .WriteProperty("eventId", storedEvent.Data.EventId.ToString()) + .WriteProperty("eventNumber", storedEvent.EventNumber)); + + throw; } } } diff --git a/src/Squidex.Infrastructure/CQRS/Events/IEventConsumer.cs b/src/Squidex.Infrastructure/CQRS/Events/IEventConsumer.cs index da0eedbb8..936cd37cd 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/IEventConsumer.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/IEventConsumer.cs @@ -14,6 +14,8 @@ namespace Squidex.Infrastructure.CQRS.Events { string Name { get; } + string StreamFilter { get; } + Task ClearAsync(); Task On(Envelope @event); diff --git a/src/Squidex.Infrastructure/CQRS/Events/IEventStore.cs b/src/Squidex.Infrastructure/CQRS/Events/IEventStore.cs index 33648ffdd..ef69f6b03 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/IEventStore.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/IEventStore.cs @@ -15,9 +15,9 @@ namespace Squidex.Infrastructure.CQRS.Events { public interface IEventStore { - IObservable GetEventsAsync(string streamName = null, long lastReceivedEventNumber = -1); + IObservable GetEventsAsync(string streamFilter = null, long lastReceivedEventNumber = -1); - Task GetEventsAsync(Func callback, CancellationToken cancellationToken, string streamName = null, long lastReceivedEventNumber = -1); + Task GetEventsAsync(Func callback, CancellationToken cancellationToken, string streamFilter = null, long lastReceivedEventNumber = -1); Task AppendEventsAsync(Guid commitId, string streamName, int expectedVersion, IEnumerable events); } diff --git a/src/Squidex.Infrastructure/CQRS/Events/Internal/DispatchEventBlock.cs b/src/Squidex.Infrastructure/CQRS/Events/Internal/DispatchEventBlock.cs deleted file mode 100644 index 022314177..000000000 --- a/src/Squidex.Infrastructure/CQRS/Events/Internal/DispatchEventBlock.cs +++ /dev/null @@ -1,115 +0,0 @@ -// ========================================================================== -// DispatchEventBlock.cs -// Squidex Headless CMS -// ========================================================================== -// Copyright (c) Squidex Group -// All rights reserved. -// ========================================================================== - -using System; -using System.Threading.Tasks; -using System.Threading.Tasks.Dataflow; -using Squidex.Infrastructure.Log; - -namespace Squidex.Infrastructure.CQRS.Events.Internal -{ - internal sealed class DispatchEventBlock : IEventReceiverBlock - { - private readonly ISemanticLog log; - private readonly IEventConsumer eventConsumer; - private readonly TransformBlock, Envelope> transformBlock; - private long lastReceivedEventNumber = -1; - private bool isRunning = true; - - public Action OnError { get; set; } - - public ITargetBlock> Target - { - get { return transformBlock; } - } - - public DispatchEventBlock(IEventConsumer eventConsumer, ISemanticLog log) - { - this.eventConsumer = eventConsumer; - - this.log = log; - - var nullHandler = new ActionBlock>(x => { }); - - transformBlock = - new TransformBlock, Envelope>(x => HandleAsync(x), - new ExecutionDataflowBlockOptions { BoundedCapacity = 1 }); - transformBlock.LinkTo(nullHandler, x => x == null); - } - - public void LinkTo(ITargetBlock> target) - { - transformBlock.LinkTo(target, new DataflowLinkOptions { PropagateCompletion = true }); - } - - public void Stop() - { - isRunning = false; - } - - public void Reset() - { - isRunning = true; - - lastReceivedEventNumber = -1; - } - - private async Task> HandleAsync(Envelope input) - { - var eventNumber = input.Headers.EventNumber(); - - if (eventNumber <= lastReceivedEventNumber || !isRunning) - { - return null; - } - - var consumerName = eventConsumer.Name; - - var eventId = input.Headers.EventId().ToString(); - var eventType = input.Payload.GetType().Name; - try - { - log.LogInformation(w => w - .WriteProperty("action", "HandleEvent") - .WriteProperty("actionId", eventId) - .WriteProperty("state", "Started") - .WriteProperty("eventId", eventId) - .WriteProperty("eventType", eventType) - .WriteProperty("eventConsumer", consumerName)); - - await eventConsumer.On(input); - - log.LogInformation(w => w - .WriteProperty("action", "HandleEvent") - .WriteProperty("actionId", eventId) - .WriteProperty("state", "Completed") - .WriteProperty("eventId", eventId) - .WriteProperty("eventType", eventType) - .WriteProperty("eventConsumer", consumerName)); - - lastReceivedEventNumber = eventNumber; - - return input; - } - catch (Exception ex) - { - OnError?.Invoke(ex); - - log.LogError(ex, w => w - .WriteProperty("action", "HandleEvent") - .WriteProperty("actionId", eventId) - .WriteProperty("state", "Started") - .WriteProperty("eventId", eventId) - .WriteProperty("eventType", eventType) - .WriteProperty("eventConsumer", consumerName)); - - return null; - } - } - } -} diff --git a/src/Squidex.Infrastructure/CQRS/Events/Internal/IEventReceiverBlock.cs b/src/Squidex.Infrastructure/CQRS/Events/Internal/IEventReceiverBlock.cs deleted file mode 100644 index 4cb8d0d00..000000000 --- a/src/Squidex.Infrastructure/CQRS/Events/Internal/IEventReceiverBlock.cs +++ /dev/null @@ -1,21 +0,0 @@ -// ========================================================================== -// IEventReceiverBlock.cs -// Squidex Headless CMS -// ========================================================================== -// Copyright (c) Squidex Group -// All rights reserved. -// ========================================================================== - -using System; - -namespace Squidex.Infrastructure.CQRS.Events.Internal -{ - public interface IEventReceiverBlock - { - Action OnError { get; set; } - - void Reset(); - - void Stop(); - } -} diff --git a/src/Squidex.Infrastructure/CQRS/Events/Internal/ParseEventBlock.cs b/src/Squidex.Infrastructure/CQRS/Events/Internal/ParseEventBlock.cs deleted file mode 100644 index 9c2e498dd..000000000 --- a/src/Squidex.Infrastructure/CQRS/Events/Internal/ParseEventBlock.cs +++ /dev/null @@ -1,98 +0,0 @@ -// ========================================================================== -// ParseEventBlock.cs -// Squidex Headless CMS -// ========================================================================== -// Copyright (c) Squidex Group -// All rights reserved. -// ========================================================================== - -using System; -using System.Threading.Tasks; -using System.Threading.Tasks.Dataflow; -using Squidex.Infrastructure.Log; - -namespace Squidex.Infrastructure.CQRS.Events.Internal -{ - internal sealed class ParseEventBlock : IEventReceiverBlock - { - private readonly EventDataFormatter formatter; - private readonly ISemanticLog log; - private readonly TransformBlock> transformBlock; - private long lastReceivedEventNumber = -1; - private bool isRunning = true; - - public Action OnError { get; set; } - - public ParseEventBlock(EventDataFormatter formatter, ISemanticLog log) - { - this.formatter = formatter; - this.log = log; - - var nullHandler = new ActionBlock>(x => { }); - - transformBlock = - new TransformBlock>(x => HandleAsync(x), - new ExecutionDataflowBlockOptions { BoundedCapacity = 1 }); - transformBlock.LinkTo(nullHandler, x => x == null); - } - - public Task NextAsync(StoredEvent input) - { - return transformBlock.SendAsync(input); - } - - public void LinkTo(ITargetBlock> target) - { - transformBlock.LinkTo(target, new DataflowLinkOptions { PropagateCompletion = true }); - } - - public void Complete() - { - transformBlock.Complete(); - } - - public void Stop() - { - isRunning = false; - } - - public void Reset() - { - isRunning = true; - - lastReceivedEventNumber = -1; - } - - private Envelope HandleAsync(StoredEvent input) - { - var eventNumber = input.EventNumber; - - if (eventNumber <= lastReceivedEventNumber || !isRunning) - { - return null; - } - - try - { - var result = formatter.Parse(input.Data); - - result.SetEventNumber(input.EventNumber); - result.SetEventStreamNumber(input.EventStreamNumber); - - return result; - } - catch (Exception ex) - { - OnError?.Invoke(ex); - - log.LogFatal(ex, w => w - .WriteProperty("action", "ParseEvent") - .WriteProperty("state", "Failed") - .WriteProperty("eventId", input.Data.EventId.ToString()) - .WriteProperty("eventNumber", input.EventNumber)); - - return null; - } - } - } -} diff --git a/src/Squidex.Infrastructure/CQRS/Events/Internal/QueryEventsBlock.cs b/src/Squidex.Infrastructure/CQRS/Events/Internal/QueryEventsBlock.cs deleted file mode 100644 index 7e3f986ff..000000000 --- a/src/Squidex.Infrastructure/CQRS/Events/Internal/QueryEventsBlock.cs +++ /dev/null @@ -1,161 +0,0 @@ -// ========================================================================== -// QueryEventsBlock.cs -// Squidex Headless CMS -// ========================================================================== -// Copyright (c) Squidex Group -// All rights reserved. -// ========================================================================== - -using System; -using System.Threading; -using System.Threading.Tasks; -using System.Threading.Tasks.Dataflow; -using Squidex.Infrastructure.Log; - -namespace Squidex.Infrastructure.CQRS.Events.Internal -{ - internal sealed class QueryEventsBlock : IEventReceiverBlock - { - private readonly ISemanticLog log; - private readonly IEventConsumerInfoRepository eventConsumerInfoRepository; - private readonly IEventConsumer eventConsumer; - private readonly IEventStore eventStore; - private readonly ActionBlock> actionBlock; - private bool isRunning = true; - private bool isStarted; - - public Action OnError { get; set; } - - public Action OnReset { get; set; } - - public Func OnEvent { get; set; } - - public ITargetBlock> Target - { - get { return actionBlock; } - } - - public Task Completion - { - get { return actionBlock.Completion; } - } - - public QueryEventsBlock(IEventConsumerInfoRepository eventConsumerInfoRepository, IEventConsumer eventConsumer, IEventStore eventStore, ISemanticLog log) - { - this.eventConsumerInfoRepository = eventConsumerInfoRepository; - this.eventConsumer = eventConsumer; - this.eventStore = eventStore; - - this.log = log; - - actionBlock = - new ActionBlock>(HandleAsync, - new ExecutionDataflowBlockOptions { BoundedCapacity = 1 }); - } - - public void Complete() - { - actionBlock.Complete(); - } - - public void NextOrThrowAway() - { - actionBlock.Post(null); - } - - public void Stop() - { - isRunning = false; - } - - public void Reset() - { - isRunning = true; - } - - private async Task HandleAsync(object input) - { - try - { - if (!isStarted) - { - await eventConsumerInfoRepository.CreateAsync(eventConsumer.Name); - - isStarted = true; - } - - var status = await eventConsumerInfoRepository.FindAsync(eventConsumer.Name); - - var lastReceivedEventNumber = status.LastHandledEventNumber; - - if (status.IsResetting) - { - await ResetAsync(); - - Reset(); - } - - if (!status.IsStopped || !isRunning) - { - var ct = new CancellationTokenSource(); - - await eventStore.GetEventsAsync(async storedEvent => - { - if (!isRunning) - { - ct.Cancel(); - } - - var onEvent = OnEvent; - - if (onEvent != null) - { - await onEvent(storedEvent); - } - }, ct.Token, null, lastReceivedEventNumber); - } - } - catch (OperationCanceledException) - { - } - catch (Exception ex) - { - OnError?.Invoke(ex); - } - } - - private async Task ResetAsync() - { - var consumerName = eventConsumer.Name; - - var actionId = Guid.NewGuid().ToString(); - try - { - log.LogInformation(w => w - .WriteProperty("action", "EventConsumerReset") - .WriteProperty("actionId", actionId) - .WriteProperty("state", "Started") - .WriteProperty("eventConsumer", consumerName)); - - await eventConsumer.ClearAsync(); - await eventConsumerInfoRepository.SetLastHandledEventNumberAsync(consumerName, -1); - - log.LogInformation(w => w - .WriteProperty("action", "EventConsumerReset") - .WriteProperty("actionId", actionId) - .WriteProperty("state", "Completed") - .WriteProperty("eventConsumer", consumerName)); - - OnReset?.Invoke(); - } - catch (Exception ex) - { - log.LogFatal(ex, w => w - .WriteProperty("action", "EventConsumerReset") - .WriteProperty("actionId", actionId) - .WriteProperty("state", "Completed") - .WriteProperty("eventConsumer", consumerName)); - } - } - } -} diff --git a/src/Squidex.Infrastructure/CQRS/Events/Internal/UpdateStateBlock.cs b/src/Squidex.Infrastructure/CQRS/Events/Internal/UpdateStateBlock.cs deleted file mode 100644 index 1b1b63697..000000000 --- a/src/Squidex.Infrastructure/CQRS/Events/Internal/UpdateStateBlock.cs +++ /dev/null @@ -1,86 +0,0 @@ -// ========================================================================== -// UpdateStateBlock.cs -// Squidex Headless CMS -// ========================================================================== -// Copyright (c) Squidex Group -// All rights reserved. -// ========================================================================== - -using System; -using System.Threading.Tasks; -using System.Threading.Tasks.Dataflow; -using Squidex.Infrastructure.Log; - -namespace Squidex.Infrastructure.CQRS.Events.Internal -{ - internal sealed class UpdateStateBlock : IEventReceiverBlock - { - private readonly ISemanticLog log; - private readonly IEventConsumerInfoRepository eventConsumerInfoRepository; - private readonly IEventConsumer eventConsumer; - private readonly ActionBlock> actionBlock; - private long lastReceivedEventNumber = -1; - private bool isRunning = true; - - public Action OnError { get; set; } - - public ITargetBlock> Target - { - get { return actionBlock; } - } - - public Task Completion - { - get { return actionBlock.Completion; } - } - - public UpdateStateBlock(IEventConsumerInfoRepository eventConsumerInfoRepository, IEventConsumer eventConsumer, ISemanticLog log) - { - this.eventConsumerInfoRepository = eventConsumerInfoRepository; - this.eventConsumer = eventConsumer; - - this.log = log; - - actionBlock = - new ActionBlock>(HandleAsync, - new ExecutionDataflowBlockOptions { BoundedCapacity = 1 }); - } - - public void Stop() - { - isRunning = false; - } - - public void Reset() - { - isRunning = true; - - lastReceivedEventNumber = -1; - } - - private async Task HandleAsync(Envelope input) - { - var eventNumber = input.Headers.EventNumber(); - - if (eventNumber <= lastReceivedEventNumber || !isRunning) - { - return; - } - - try - { - await eventConsumerInfoRepository.SetLastHandledEventNumberAsync(eventConsumer.Name, eventNumber); - } - catch (Exception ex) - { - OnError?.Invoke(ex); - - log.LogFatal(ex, w => w - .WriteProperty("action", "UpdateState") - .WriteProperty("state", "Failed") - .WriteProperty("eventId", input.Headers.EventId().ToString()) - .WriteProperty("eventNumber", input.Headers.EventNumber())); - } - } - } -} diff --git a/src/Squidex.Infrastructure/Squidex.Infrastructure.csproj b/src/Squidex.Infrastructure/Squidex.Infrastructure.csproj index 92d0ccac6..2319a1bb8 100644 --- a/src/Squidex.Infrastructure/Squidex.Infrastructure.csproj +++ b/src/Squidex.Infrastructure/Squidex.Infrastructure.csproj @@ -16,6 +16,5 @@ - diff --git a/src/Squidex.Read.MongoDb/Apps/MongoAppRepository_EventHandling.cs b/src/Squidex.Read.MongoDb/Apps/MongoAppRepository_EventHandling.cs index ff92bae3b..5dfb421a8 100644 --- a/src/Squidex.Read.MongoDb/Apps/MongoAppRepository_EventHandling.cs +++ b/src/Squidex.Read.MongoDb/Apps/MongoAppRepository_EventHandling.cs @@ -23,6 +23,11 @@ namespace Squidex.Read.MongoDb.Apps get { return GetType().Name; } } + public string StreamFilter + { + get { return "^app-"; } + } + public Task On(Envelope @event) { return this.DispatchActionAsync(@event.Payload, @event.Headers); diff --git a/src/Squidex.Read.MongoDb/Contents/MongoContentRepository_EventHandling.cs b/src/Squidex.Read.MongoDb/Contents/MongoContentRepository_EventHandling.cs index 4d13b2e00..481f4a52e 100644 --- a/src/Squidex.Read.MongoDb/Contents/MongoContentRepository_EventHandling.cs +++ b/src/Squidex.Read.MongoDb/Contents/MongoContentRepository_EventHandling.cs @@ -37,6 +37,11 @@ namespace Squidex.Read.MongoDb.Contents get { return GetType().Name; } } + public string StreamFilter + { + get { return "^content-"; } + } + public async Task ClearAsync() { using (var collections = await database.ListCollectionsAsync()) diff --git a/src/Squidex.Read.MongoDb/History/MongoHistoryEventRepository.cs b/src/Squidex.Read.MongoDb/History/MongoHistoryEventRepository.cs index da4376ba3..5925db0ff 100644 --- a/src/Squidex.Read.MongoDb/History/MongoHistoryEventRepository.cs +++ b/src/Squidex.Read.MongoDb/History/MongoHistoryEventRepository.cs @@ -27,6 +27,16 @@ namespace Squidex.Read.MongoDb.History private readonly Dictionary texts = new Dictionary(); private int sessionEventCount; + public string Name + { + get { return GetType().Name; } + } + + public string StreamFilter + { + get { return "*"; } + } + public MongoHistoryEventRepository(IMongoDatabase database, IEnumerable creators) : base(database) { @@ -67,11 +77,6 @@ namespace Squidex.Read.MongoDb.History return entities.Select(x => (IHistoryEventEntity)new ParsedHistoryEvent(x, texts)).ToList(); } - public string Name - { - get { return GetType().Name; } - } - public async Task On(Envelope @event) { foreach (var creator in creators) diff --git a/src/Squidex.Read.MongoDb/Schemas/MongoSchemaRepository_EventHandling.cs b/src/Squidex.Read.MongoDb/Schemas/MongoSchemaRepository_EventHandling.cs index 93e2229b4..d35a916fb 100644 --- a/src/Squidex.Read.MongoDb/Schemas/MongoSchemaRepository_EventHandling.cs +++ b/src/Squidex.Read.MongoDb/Schemas/MongoSchemaRepository_EventHandling.cs @@ -26,6 +26,11 @@ namespace Squidex.Read.MongoDb.Schemas get { return GetType().Name; } } + public string StreamFilter + { + get { return "^schema-"; } + } + public Task On(Envelope @event) { return this.DispatchActionAsync(@event.Payload, @event.Headers); diff --git a/src/Squidex.Read/Apps/Services/Implementations/CachingAppProvider.cs b/src/Squidex.Read/Apps/Services/Implementations/CachingAppProvider.cs index 18223174e..c5ec9c8cd 100644 --- a/src/Squidex.Read/Apps/Services/Implementations/CachingAppProvider.cs +++ b/src/Squidex.Read/Apps/Services/Implementations/CachingAppProvider.cs @@ -32,6 +32,11 @@ namespace Squidex.Read.Apps.Services.Implementations get { return GetType().Name; } } + public string StreamFilter + { + get { return "*"; } + } + public CachingAppProvider(IMemoryCache cache, IAppRepository repository) : base(cache) { diff --git a/src/Squidex.Read/Schemas/Services/Implementations/CachingSchemaProvider.cs b/src/Squidex.Read/Schemas/Services/Implementations/CachingSchemaProvider.cs index 860d1764f..a667fe8cb 100644 --- a/src/Squidex.Read/Schemas/Services/Implementations/CachingSchemaProvider.cs +++ b/src/Squidex.Read/Schemas/Services/Implementations/CachingSchemaProvider.cs @@ -32,6 +32,11 @@ namespace Squidex.Read.Schemas.Services.Implementations get { return GetType().Name; } } + public string StreamFilter + { + get { return "*"; } + } + public CachingSchemaProvider(IMemoryCache cache, ISchemaRepository repository) : base(cache) { diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Events/CompoundEventConsumerTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Events/CompoundEventConsumerTests.cs index e01b23c43..7e49ac0fa 100644 --- a/tests/Squidex.Infrastructure.Tests/CQRS/Events/CompoundEventConsumerTests.cs +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Events/CompoundEventConsumerTests.cs @@ -33,9 +33,25 @@ namespace Squidex.Infrastructure.CQRS.Events [Fact] public void Should_return_first_inner_name() { + const string name = "my-inner-consumer"; + + consumer1.Setup(x => x.Name).Returns(name); + + var sut = new CompoundEventConsumer(consumer1.Object, consumer2.Object); + + Assert.Equal(name, sut.Name); + } + + [Fact] + public void Should_return_first_inner_filter() + { + const string filter = "my-inner-filter"; + + consumer1.Setup(x => x.StreamFilter).Returns(filter); + var sut = new CompoundEventConsumer(consumer1.Object, consumer2.Object); - Assert.Equal(consumer1.Object.GetType().Name, sut.Name); + Assert.Equal(filter, sut.StreamFilter); } [Fact] diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Events/EventReceiverTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Events/EventReceiverTests.cs index 47444876a..c8c6844a2 100644 --- a/tests/Squidex.Infrastructure.Tests/CQRS/Events/EventReceiverTests.cs +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Events/EventReceiverTests.cs @@ -45,7 +45,7 @@ namespace Squidex.Infrastructure.CQRS.Events this.storedEvents = storedEvents; } - public async Task GetEventsAsync(Func callback, CancellationToken cancellationToken, string streamName = null, long lastReceivedEventNumber = -1) + public async Task GetEventsAsync(Func callback, CancellationToken cancellationToken, string streamFilter = null, long lastReceivedEventNumber = -1) { foreach (var @event in storedEvents) { @@ -53,7 +53,7 @@ namespace Squidex.Infrastructure.CQRS.Events } } - public IObservable GetEventsAsync(string streamName = null, long lastReceivedEventNumber = -1) + public IObservable GetEventsAsync(string streamFilter = null, long lastReceivedEventNumber = -1) { throw new NotSupportedException(); } @@ -170,7 +170,7 @@ namespace Squidex.Infrastructure.CQRS.Events consumerInfo.IsResetting = true; consumerInfo.LastHandledEventNumber = 2L; - sut.Subscribe(eventConsumer.Object, autoTrigger: false); + sut.Subscribe(eventConsumer.Object); sut.Next(); sut.Dispose();