Browse Source

Code fixed

pull/130/merge
Sebastian Stehle 9 years ago
parent
commit
a08367e419
  1. 20
      src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs
  2. 29
      tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs

20
src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs

@ -31,11 +31,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors
public IEventConsumer EventConsumer { get; set; }
}
private sealed class StopFailed
{
public Exception Exception { get; set; }
}
private abstract class SubscriptionMessage
{
public IEventSubscription Subscription { get; set; }
@ -84,16 +79,16 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors
}
}
protected override Task OnError(Exception exception)
protected override async Task OnError(Exception exception)
{
log.LogError(exception, w => w
.WriteProperty("action", "HandleEvent")
.WriteProperty("state", "Failed")
.WriteProperty("eventConsumer", eventConsumer.Name));
DispatchAsync(new StopFailed { Exception = exception }).Forget();
await StopAsync(exception);
return TaskHelper.Done;
isRunning = false;
}
Task IEventSubscriber.OnEventAsync(IEventSubscription subscription, StoredEvent @event)
@ -135,15 +130,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors
break;
}
case StopFailed stopFailed when isSetup && isRunning:
{
await StopAsync(stopFailed.Exception);
isRunning = false;
break;
}
case StopConsumerMessage stopConsumer when isSetup && isRunning:
{
await StopAsync();

29
tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs

@ -196,9 +196,34 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors
var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData);
await SubscribeAsync();
await sutSubscriber.OnEventAsync(eventSubscription, @event);
sut.Dispose();
A.CallTo(() => eventConsumer.On(envelope))
.MustHaveHappened();
A.CallTo(() => eventConsumerInfoRepository.SetPositionAsync(consumerName, @event.EventPosition, false))
.MustNotHaveHappened();
A.CallTo(() => eventConsumerInfoRepository.StopAsync(consumerName, exception.ToString()))
.MustHaveHappened();
}
[Fact]
public async Task Should_start_after_stop_when_handling_failed()
{
var exception = new InvalidOperationException("Exception");
A.CallTo(() => eventConsumer.On(envelope))
.Throws(exception);
var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData);
await SubscribeAsync();
await sutSubscriber.OnEventAsync(eventSubscription, @event);
sutActor.Tell(new StartConsumerMessage());
sut.Dispose();
A.CallTo(() => eventConsumer.On(envelope))
@ -209,6 +234,9 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors
A.CallTo(() => eventConsumerInfoRepository.StopAsync(consumerName, exception.ToString()))
.MustHaveHappened();
A.CallTo(() => eventConsumerInfoRepository.StartAsync(consumerName))
.MustHaveHappened();
}
[Fact]
@ -222,7 +250,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors
var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData);
await SubscribeAsync();
await sutSubscriber.OnEventAsync(eventSubscription, @event);
sut.Dispose();

Loading…
Cancel
Save