Browse Source

1) EventReceiver simplified.

2) Filter EventConsumers with regex. (#2)
pull/65/head
Sebastian Stehle 10 years ago
parent
commit
fc7f399274
  1. 18
      src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs
  2. 7
      src/Squidex.Infrastructure/CQRS/Events/CompoundEventConsumer.cs
  3. 177
      src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs
  4. 2
      src/Squidex.Infrastructure/CQRS/Events/IEventConsumer.cs
  5. 4
      src/Squidex.Infrastructure/CQRS/Events/IEventStore.cs
  6. 115
      src/Squidex.Infrastructure/CQRS/Events/Internal/DispatchEventBlock.cs
  7. 21
      src/Squidex.Infrastructure/CQRS/Events/Internal/IEventReceiverBlock.cs
  8. 98
      src/Squidex.Infrastructure/CQRS/Events/Internal/ParseEventBlock.cs
  9. 161
      src/Squidex.Infrastructure/CQRS/Events/Internal/QueryEventsBlock.cs
  10. 86
      src/Squidex.Infrastructure/CQRS/Events/Internal/UpdateStateBlock.cs
  11. 1
      src/Squidex.Infrastructure/Squidex.Infrastructure.csproj
  12. 5
      src/Squidex.Read.MongoDb/Apps/MongoAppRepository_EventHandling.cs
  13. 5
      src/Squidex.Read.MongoDb/Contents/MongoContentRepository_EventHandling.cs
  14. 15
      src/Squidex.Read.MongoDb/History/MongoHistoryEventRepository.cs
  15. 5
      src/Squidex.Read.MongoDb/Schemas/MongoSchemaRepository_EventHandling.cs
  16. 5
      src/Squidex.Read/Apps/Services/Implementations/CachingAppProvider.cs
  17. 5
      src/Squidex.Read/Schemas/Services/Implementations/CachingSchemaProvider.cs
  18. 18
      tests/Squidex.Infrastructure.Tests/CQRS/Events/CompoundEventConsumerTests.cs
  19. 6
      tests/Squidex.Infrastructure.Tests/CQRS/Events/EventReceiverTests.cs

18
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<StoredEvent> GetEventsAsync(string streamName, long lastReceivedEventNumber = -1)
public IObservable<StoredEvent> GetEventsAsync(string streamFilter, long lastReceivedEventNumber = -1)
{
return Observable.Create<StoredEvent>((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<StoredEvent, Task> callback, CancellationToken cancellationToken, string streamName = null, long lastReceivedEventNumber = -1)
public async Task GetEventsAsync(Func<StoredEvent, Task> 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<MongoEventCommit> filter = new BsonDocument();

7
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)

177
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<Exception>(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<IEvent> @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<IEvent> 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;
}
}
}

2
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<IEvent> @event);

4
src/Squidex.Infrastructure/CQRS/Events/IEventStore.cs

@ -15,9 +15,9 @@ namespace Squidex.Infrastructure.CQRS.Events
{
public interface IEventStore
{
IObservable<StoredEvent> GetEventsAsync(string streamName = null, long lastReceivedEventNumber = -1);
IObservable<StoredEvent> GetEventsAsync(string streamFilter = null, long lastReceivedEventNumber = -1);
Task GetEventsAsync(Func<StoredEvent, Task> callback, CancellationToken cancellationToken, string streamName = null, long lastReceivedEventNumber = -1);
Task GetEventsAsync(Func<StoredEvent, Task> callback, CancellationToken cancellationToken, string streamFilter = null, long lastReceivedEventNumber = -1);
Task AppendEventsAsync(Guid commitId, string streamName, int expectedVersion, IEnumerable<EventData> events);
}

115
src/Squidex.Infrastructure/CQRS/Events/Internal/DispatchEventBlock.cs

@ -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<IEvent>, Envelope<IEvent>> transformBlock;
private long lastReceivedEventNumber = -1;
private bool isRunning = true;
public Action<Exception> OnError { get; set; }
public ITargetBlock<Envelope<IEvent>> Target
{
get { return transformBlock; }
}
public DispatchEventBlock(IEventConsumer eventConsumer, ISemanticLog log)
{
this.eventConsumer = eventConsumer;
this.log = log;
var nullHandler = new ActionBlock<Envelope<IEvent>>(x => { });
transformBlock =
new TransformBlock<Envelope<IEvent>, Envelope<IEvent>>(x => HandleAsync(x),
new ExecutionDataflowBlockOptions { BoundedCapacity = 1 });
transformBlock.LinkTo(nullHandler, x => x == null);
}
public void LinkTo(ITargetBlock<Envelope<IEvent>> target)
{
transformBlock.LinkTo(target, new DataflowLinkOptions { PropagateCompletion = true });
}
public void Stop()
{
isRunning = false;
}
public void Reset()
{
isRunning = true;
lastReceivedEventNumber = -1;
}
private async Task<Envelope<IEvent>> HandleAsync(Envelope<IEvent> 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;
}
}
}
}

21
src/Squidex.Infrastructure/CQRS/Events/Internal/IEventReceiverBlock.cs

@ -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<Exception> OnError { get; set; }
void Reset();
void Stop();
}
}

98
src/Squidex.Infrastructure/CQRS/Events/Internal/ParseEventBlock.cs

@ -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<StoredEvent, Envelope<IEvent>> transformBlock;
private long lastReceivedEventNumber = -1;
private bool isRunning = true;
public Action<Exception> OnError { get; set; }
public ParseEventBlock(EventDataFormatter formatter, ISemanticLog log)
{
this.formatter = formatter;
this.log = log;
var nullHandler = new ActionBlock<Envelope<IEvent>>(x => { });
transformBlock =
new TransformBlock<StoredEvent, Envelope<IEvent>>(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<Envelope<IEvent>> 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<IEvent> 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;
}
}
}
}

161
src/Squidex.Infrastructure/CQRS/Events/Internal/QueryEventsBlock.cs

@ -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<Envelope<IEvent>> actionBlock;
private bool isRunning = true;
private bool isStarted;
public Action<Exception> OnError { get; set; }
public Action OnReset { get; set; }
public Func<StoredEvent, Task> OnEvent { get; set; }
public ITargetBlock<Envelope<IEvent>> 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<Envelope<IEvent>>(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));
}
}
}
}

86
src/Squidex.Infrastructure/CQRS/Events/Internal/UpdateStateBlock.cs

@ -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<Envelope<IEvent>> actionBlock;
private long lastReceivedEventNumber = -1;
private bool isRunning = true;
public Action<Exception> OnError { get; set; }
public ITargetBlock<Envelope<IEvent>> 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<Envelope<IEvent>>(HandleAsync,
new ExecutionDataflowBlockOptions { BoundedCapacity = 1 });
}
public void Stop()
{
isRunning = false;
}
public void Reset()
{
isRunning = true;
lastReceivedEventNumber = -1;
}
private async Task HandleAsync(Envelope<IEvent> 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()));
}
}
}
}

1
src/Squidex.Infrastructure/Squidex.Infrastructure.csproj

@ -16,6 +16,5 @@
<PackageReference Include="System.Reactive" Version="3.1.1" />
<PackageReference Include="System.Reflection.TypeExtensions" Version="4.3.0" />
<PackageReference Include="System.Security.Claims" Version="4.3.0" />
<PackageReference Include="System.Threading.Tasks.Dataflow" Version="4.7.0" />
</ItemGroup>
</Project>

5
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<IEvent> @event)
{
return this.DispatchActionAsync(@event.Payload, @event.Headers);

5
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())

15
src/Squidex.Read.MongoDb/History/MongoHistoryEventRepository.cs

@ -27,6 +27,16 @@ namespace Squidex.Read.MongoDb.History
private readonly Dictionary<string, string> texts = new Dictionary<string, string>();
private int sessionEventCount;
public string Name
{
get { return GetType().Name; }
}
public string StreamFilter
{
get { return "*"; }
}
public MongoHistoryEventRepository(IMongoDatabase database, IEnumerable<IHistoryEventsCreator> 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<IEvent> @event)
{
foreach (var creator in creators)

5
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<IEvent> @event)
{
return this.DispatchActionAsync(@event.Payload, @event.Headers);

5
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)
{

5
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)
{

18
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]

6
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<StoredEvent, Task> callback, CancellationToken cancellationToken, string streamName = null, long lastReceivedEventNumber = -1)
public async Task GetEventsAsync(Func<StoredEvent, Task> 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<StoredEvent> GetEventsAsync(string streamName = null, long lastReceivedEventNumber = -1)
public IObservable<StoredEvent> 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();

Loading…
Cancel
Save