mirror of https://github.com/Squidex/squidex.git
13 changed files with 554 additions and 194 deletions
@ -0,0 +1,110 @@ |
|||
// ==========================================================================
|
|||
// Actor.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.Tasks; |
|||
|
|||
#pragma warning disable SA1401 // Fields must be private
|
|||
|
|||
namespace Squidex.Infrastructure.Actors |
|||
{ |
|||
public abstract class Actor : IActor, IDisposable |
|||
{ |
|||
private readonly ActionBlock<IMessage> block; |
|||
private readonly ReaderWriterLockSlim slimLock = new ReaderWriterLockSlim(); |
|||
private volatile bool isStopped; |
|||
|
|||
private sealed class StopMessage : IMessage |
|||
{ |
|||
} |
|||
|
|||
private sealed class ErrorMessage : IMessage |
|||
{ |
|||
public Exception Exception; |
|||
} |
|||
|
|||
protected Actor() |
|||
{ |
|||
block = new ActionBlock<IMessage>(Handle, new ExecutionDataflowBlockOptions { BoundedCapacity = 100 }); |
|||
} |
|||
|
|||
public void Dispose() |
|||
{ |
|||
StopAsync().Wait(); |
|||
} |
|||
|
|||
public async Task StopAsync() |
|||
{ |
|||
isStopped = true; |
|||
|
|||
await block.SendAsync(new StopMessage()); |
|||
|
|||
await block.Completion; |
|||
} |
|||
|
|||
public Task SendAsync(IMessage message) |
|||
{ |
|||
Guard.NotNull(message, nameof(message)); |
|||
|
|||
return block.SendAsync(message); |
|||
} |
|||
|
|||
public Task SendAsync(Exception exception) |
|||
{ |
|||
Guard.NotNull(exception, nameof(exception)); |
|||
|
|||
return block.SendAsync(new ErrorMessage { Exception = exception }); |
|||
} |
|||
|
|||
protected virtual Task OnStop() |
|||
{ |
|||
return TaskHelper.Done; |
|||
} |
|||
|
|||
protected virtual Task OnError(Exception exception) |
|||
{ |
|||
return TaskHelper.Done; |
|||
} |
|||
|
|||
protected virtual Task OnMessage(IMessage message) |
|||
{ |
|||
return TaskHelper.Done; |
|||
} |
|||
|
|||
private async Task Handle(IMessage message) |
|||
{ |
|||
try |
|||
{ |
|||
if (message is StopMessage) |
|||
{ |
|||
block.Complete(); |
|||
|
|||
await OnStop(); |
|||
} |
|||
else if (message is ErrorMessage errorMessage) |
|||
{ |
|||
await OnError(errorMessage.Exception); |
|||
} |
|||
else |
|||
{ |
|||
await OnMessage(message); |
|||
} |
|||
} |
|||
catch (Exception ex) |
|||
{ |
|||
if (!(message is ErrorMessage)) |
|||
{ |
|||
await block.SendAsync(new ErrorMessage { Exception = ex }); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,14 @@ |
|||
using System; |
|||
using System.Threading.Tasks; |
|||
|
|||
namespace Squidex.Infrastructure.Actors |
|||
{ |
|||
public interface IActor |
|||
{ |
|||
Task SendAsync(IMessage message); |
|||
|
|||
Task SendAsync(Exception exception); |
|||
|
|||
Task StopAsync(); |
|||
} |
|||
} |
|||
@ -0,0 +1,14 @@ |
|||
// ==========================================================================
|
|||
// IMessage.cs
|
|||
// Squidex Headless CMS
|
|||
// ==========================================================================
|
|||
// Copyright (c) Squidex Group
|
|||
// All rights reserved.
|
|||
// ==========================================================================
|
|||
|
|||
namespace Squidex.Infrastructure.Actors |
|||
{ |
|||
public interface IMessage |
|||
{ |
|||
} |
|||
} |
|||
@ -0,0 +1,215 @@ |
|||
// ==========================================================================
|
|||
// EventReceiver.cs
|
|||
// Squidex Headless CMS
|
|||
// ==========================================================================
|
|||
// Copyright (c) Squidex Group
|
|||
// All rights reserved.
|
|||
// ==========================================================================
|
|||
|
|||
using System; |
|||
using System.Threading.Tasks; |
|||
using Squidex.Infrastructure.Actors; |
|||
using Squidex.Infrastructure.CQRS.Events.Actors.Messages; |
|||
using Squidex.Infrastructure.Log; |
|||
using Squidex.Infrastructure.Tasks; |
|||
|
|||
namespace Squidex.Infrastructure.CQRS.Events.Receivers |
|||
{ |
|||
public sealed class EventConsumerActor : Actor |
|||
{ |
|||
private readonly EventDataFormatter formatter; |
|||
private readonly IEventStore eventStore; |
|||
private readonly IEventConsumerInfoRepository eventConsumerInfoRepository; |
|||
private readonly ISemanticLog log; |
|||
private IEventSubscription eventSubscription; |
|||
private IEventConsumer eventConsumer; |
|||
private string position; |
|||
|
|||
public EventConsumerActor( |
|||
EventDataFormatter formatter, |
|||
IEventStore eventStore, |
|||
IEventConsumerInfoRepository eventConsumerInfoRepository, |
|||
ISemanticLog log) |
|||
{ |
|||
Guard.NotNull(log, nameof(log)); |
|||
Guard.NotNull(formatter, nameof(formatter)); |
|||
Guard.NotNull(eventStore, nameof(eventStore)); |
|||
Guard.NotNull(eventConsumerInfoRepository, nameof(eventConsumerInfoRepository)); |
|||
|
|||
this.log = log; |
|||
|
|||
this.formatter = formatter; |
|||
this.eventStore = eventStore; |
|||
this.eventConsumerInfoRepository = eventConsumerInfoRepository; |
|||
} |
|||
|
|||
public void Subscribe(IEventConsumer eventConsumer) |
|||
{ |
|||
Guard.NotNull(eventConsumer, nameof(eventConsumer)); |
|||
|
|||
this.eventConsumer = eventConsumer; |
|||
} |
|||
|
|||
protected override async Task OnStop() |
|||
{ |
|||
if (eventSubscription != null) |
|||
{ |
|||
await eventSubscription.StopAsync(); |
|||
} |
|||
} |
|||
|
|||
protected override Task OnError(Exception exception) |
|||
{ |
|||
return StopAsync(exception); |
|||
} |
|||
|
|||
protected override async Task OnMessage(IMessage message) |
|||
{ |
|||
switch (message) |
|||
{ |
|||
case StopReceiverMessage stopReceiver: |
|||
{ |
|||
await StopAsync(stopReceiver.Exception); |
|||
|
|||
break; |
|||
} |
|||
|
|||
case StartReceiverMessage startReceiver: |
|||
{ |
|||
await StartAsync(); |
|||
|
|||
break; |
|||
} |
|||
|
|||
case ResetReceiverMessage resetReceiver: |
|||
{ |
|||
await StopAsync(); |
|||
await ResetAsync(); |
|||
await StartAsync(); |
|||
|
|||
break; |
|||
} |
|||
|
|||
case ReceiveEventMessage receiveEvent: |
|||
{ |
|||
await DispatchConsumerAsync(ParseEvent(receiveEvent.Event)); |
|||
|
|||
break; |
|||
} |
|||
} |
|||
} |
|||
|
|||
private async Task StartAsync() |
|||
{ |
|||
await eventConsumerInfoRepository.CreateAsync(eventConsumer.Name); |
|||
|
|||
position = (await eventConsumerInfoRepository.FindAsync(eventConsumer.Name)).Position; |
|||
|
|||
eventSubscription = eventStore.CreateSubscription(eventConsumer.EventsFilter, position); |
|||
eventSubscription.SendAsync(new SubscribeMessage { Parent = this }).Forget(); |
|||
} |
|||
|
|||
private async Task StopAsync(Exception exception = null) |
|||
{ |
|||
if (eventSubscription != null) |
|||
{ |
|||
await eventSubscription.StopAsync(); |
|||
} |
|||
|
|||
await eventConsumerInfoRepository.StopAsync(eventConsumer.Name, exception?.Message); |
|||
} |
|||
|
|||
private async Task ResetAsync() |
|||
{ |
|||
var actionId = Guid.NewGuid().ToString(); |
|||
try |
|||
{ |
|||
log.LogInformation(w => w |
|||
.WriteProperty("action", "EventConsumerReset") |
|||
.WriteProperty("actionId", actionId) |
|||
.WriteProperty("state", "Started") |
|||
.WriteProperty("eventConsumer", eventConsumer.Name)); |
|||
|
|||
await eventConsumer.ClearAsync(); |
|||
await eventConsumerInfoRepository.SetPositionAsync(eventConsumer.Name, null, true); |
|||
|
|||
log.LogInformation(w => w |
|||
.WriteProperty("action", "EventConsumerReset") |
|||
.WriteProperty("actionId", actionId) |
|||
.WriteProperty("state", "Completed") |
|||
.WriteProperty("eventConsumer", eventConsumer.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 DispatchConsumerAsync(Envelope<IEvent> @event) |
|||
{ |
|||
var eventId = @event.Headers.EventId().ToString(); |
|||
var eventType = @event.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", eventConsumer.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.Name)); |
|||
} |
|||
catch (Exception ex) |
|||
{ |
|||
log.LogError(ex, w => w |
|||
.WriteProperty("action", "HandleEvent") |
|||
.WriteProperty("actionId", eventId) |
|||
.WriteProperty("state", "Started") |
|||
.WriteProperty("eventId", eventId) |
|||
.WriteProperty("eventType", eventType) |
|||
.WriteProperty("eventConsumer", eventConsumer.Name)); |
|||
|
|||
throw; |
|||
} |
|||
} |
|||
|
|||
private Envelope<IEvent> ParseEvent(StoredEvent message) |
|||
{ |
|||
try |
|||
{ |
|||
var @event = formatter.Parse(message.Data); |
|||
|
|||
@event.SetEventPosition(message.EventPosition); |
|||
@event.SetEventStreamNumber(message.EventStreamNumber); |
|||
|
|||
return @event; |
|||
} |
|||
catch (Exception ex) |
|||
{ |
|||
log.LogFatal(ex, w => w |
|||
.WriteProperty("action", "ParseEvent") |
|||
.WriteProperty("state", "Failed") |
|||
.WriteProperty("eventId", message.Data.EventId.ToString()) |
|||
.WriteProperty("eventPosition", message.EventPosition)); |
|||
|
|||
throw; |
|||
} |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,9 @@ |
|||
using Squidex.Infrastructure.Actors; |
|||
|
|||
namespace Squidex.Infrastructure.CQRS.Events.Actors.Messages |
|||
{ |
|||
public sealed class ReceiveEventMessage : IMessage |
|||
{ |
|||
public StoredEvent Event { get; set; } |
|||
} |
|||
} |
|||
@ -0,0 +1,8 @@ |
|||
using Squidex.Infrastructure.Actors; |
|||
|
|||
namespace Squidex.Infrastructure.CQRS.Events.Actors.Messages |
|||
{ |
|||
public sealed class ResetReceiverMessage : IMessage |
|||
{ |
|||
} |
|||
} |
|||
@ -0,0 +1,8 @@ |
|||
using Squidex.Infrastructure.Actors; |
|||
|
|||
namespace Squidex.Infrastructure.CQRS.Events.Actors.Messages |
|||
{ |
|||
public sealed class StartReceiverMessage : IMessage |
|||
{ |
|||
} |
|||
} |
|||
@ -0,0 +1,10 @@ |
|||
using System; |
|||
using Squidex.Infrastructure.Actors; |
|||
|
|||
namespace Squidex.Infrastructure.CQRS.Events.Actors.Messages |
|||
{ |
|||
public sealed class StopReceiverMessage : IMessage |
|||
{ |
|||
public Exception Exception { get; set; } |
|||
} |
|||
} |
|||
@ -0,0 +1,9 @@ |
|||
using Squidex.Infrastructure.Actors; |
|||
|
|||
namespace Squidex.Infrastructure.CQRS.Events.Actors.Messages |
|||
{ |
|||
public sealed class SubscribeMessage : IMessage |
|||
{ |
|||
public IActor Parent { get; set; } |
|||
} |
|||
} |
|||
Loading…
Reference in new issue