diff --git a/src/Squidex.Infrastructure.MongoDb/CQRS/Events/PollingSubscription.cs b/src/Squidex.Infrastructure.MongoDb/CQRS/Events/PollingSubscription.cs index dbb2af459..5f4e48bf8 100644 --- a/src/Squidex.Infrastructure.MongoDb/CQRS/Events/PollingSubscription.cs +++ b/src/Squidex.Infrastructure.MongoDb/CQRS/Events/PollingSubscription.cs @@ -21,6 +21,7 @@ namespace Squidex.Infrastructure.CQRS.Events private readonly MongoEventStore eventStore; private readonly string streamFilter; private string position; + private bool isStopped; private IDisposable subscription; private CompletionTimer timer; @@ -36,9 +37,11 @@ namespace Squidex.Infrastructure.CQRS.Events { if (disposing) { + isStopped = true; + subscription?.Dispose(); - timer?.Dispose(); + timer?.StopAsync().Forget(); } } @@ -57,22 +60,28 @@ namespace Squidex.Infrastructure.CQRS.Events { await eventStore.GetEventsAsync(async storedEvent => { - await onNext(storedEvent); + if (!isStopped) + { + await onNext(storedEvent); - position = storedEvent.EventPosition; + position = storedEvent.EventPosition; + } }, ct, streamFilter, position); } catch (Exception ex) when (!(ex is OperationCanceledException)) { - onError?.Invoke(ex); + if (!isStopped) + { + onError?.Invoke(ex); + } } }); subscription = eventNotifier.Subscribe(() => { - if (!timer.IsDisposed) + if (!isStopped) { - timer.Wakeup(); + timer.SkipCurrentDelay(); } }); diff --git a/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs b/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs index 4ab00201b..2742194f2 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs @@ -58,7 +58,7 @@ namespace Squidex.Infrastructure.CQRS.Events try { - timer?.Dispose(); + timer?.StopAsync().Wait(); } catch (Exception ex) { @@ -73,7 +73,7 @@ namespace Squidex.Infrastructure.CQRS.Events { ThrowIfDisposed(); - timer?.Wakeup(); + timer?.SkipCurrentDelay(); } public void Subscribe(IEventConsumer eventConsumer) diff --git a/src/Squidex.Infrastructure/Timers/CompletionTimer.cs b/src/Squidex.Infrastructure/Timers/CompletionTimer.cs index d19c5249f..5807bd03e 100644 --- a/src/Squidex.Infrastructure/Timers/CompletionTimer.cs +++ b/src/Squidex.Infrastructure/Timers/CompletionTimer.cs @@ -10,15 +10,19 @@ using System; using System.Threading; using System.Threading.Tasks; +// ReSharper disable RedundantJumpStatement // ReSharper disable InvertIf namespace Squidex.Infrastructure.Timers { - public sealed class CompletionTimer : DisposableObjectBase + public sealed class CompletionTimer { - private readonly CancellationTokenSource disposeToken = new CancellationTokenSource(); + private const int OneCallNotExecuted = 0; + private const int OneCallExecuted = 1; + private const int OneCallRequested = 2; + private readonly CancellationTokenSource stopToken = new CancellationTokenSource(); private readonly Task runTask; - private int requiresAtLeastOne; + private int oneCallState; private CancellationTokenSource wakeupToken; public CompletionTimer(int delayInMs, Func callback, int initialDelay = 0) @@ -29,23 +33,21 @@ namespace Squidex.Infrastructure.Timers runTask = RunInternal(delayInMs, initialDelay, callback); } - protected override void DisposeObject(bool disposing) + public Task StopAsync() { - if (disposing) - { - disposeToken.Cancel(); + stopToken.Cancel(); - runTask.Wait(); - } + return runTask; } - public void Wakeup() + public void SkipCurrentDelay() { - ThrowIfDisposed(); - - Interlocked.CompareExchange(ref requiresAtLeastOne, 2, 0); + if (!stopToken.IsCancellationRequested) + { + Interlocked.CompareExchange(ref oneCallState, OneCallRequested, OneCallNotExecuted); - wakeupToken?.Cancel(); + wakeupToken?.Cancel(); + } } private async Task RunInternal(int delay, int initialDelay, Func callback) @@ -57,26 +59,18 @@ namespace Squidex.Infrastructure.Timers await WaitAsync(initialDelay).ConfigureAwait(false); } - while (requiresAtLeastOne == 2 || !disposeToken.IsCancellationRequested) + while (oneCallState == OneCallRequested || !stopToken.IsCancellationRequested) { - try - { - await callback(disposeToken.Token).ConfigureAwait(false); - } - catch (OperationCanceledException) - { - } - finally - { - requiresAtLeastOne = 1; - } + await callback(stopToken.Token).ConfigureAwait(false); + + oneCallState = OneCallExecuted; await WaitAsync(delay).ConfigureAwait(false); } } catch { - + return; } } @@ -86,7 +80,7 @@ namespace Squidex.Infrastructure.Timers { wakeupToken = new CancellationTokenSource(); - using (var cts = CancellationTokenSource.CreateLinkedTokenSource(disposeToken.Token, wakeupToken.Token)) + using (var cts = CancellationTokenSource.CreateLinkedTokenSource(stopToken.Token, wakeupToken.Token)) { await Task.Delay(intervall, cts.Token).ConfigureAwait(false); } diff --git a/src/Squidex.Infrastructure/UsageTracking/BackgroundUsageTracker.cs b/src/Squidex.Infrastructure/UsageTracking/BackgroundUsageTracker.cs index 13b7eea11..6cd3667a5 100644 --- a/src/Squidex.Infrastructure/UsageTracking/BackgroundUsageTracker.cs +++ b/src/Squidex.Infrastructure/UsageTracking/BackgroundUsageTracker.cs @@ -59,7 +59,7 @@ namespace Squidex.Infrastructure.UsageTracking { if (disposing) { - timer.Dispose(); + timer.StopAsync().Wait(); } } @@ -67,7 +67,7 @@ namespace Squidex.Infrastructure.UsageTracking { ThrowIfDisposed(); - timer.Wakeup(); + timer.SkipCurrentDelay(); } private async Task TrackAsync() diff --git a/src/Squidex/Config/Domain/ReadModule.cs b/src/Squidex/Config/Domain/ReadModule.cs index 31308a20e..1b6944fcc 100644 --- a/src/Squidex/Config/Domain/ReadModule.cs +++ b/src/Squidex/Config/Domain/ReadModule.cs @@ -80,6 +80,10 @@ namespace Squidex.Config.Domain .As() .SingleInstance(); + builder.RegisterType() + .As() + .InstancePerDependency(); + builder.RegisterType() .As() .AsSelf() diff --git a/src/Squidex/app/shared/services/auth.service.ts b/src/Squidex/app/shared/services/auth.service.ts index 6ec12056b..d51458ff3 100644 --- a/src/Squidex/app/shared/services/auth.service.ts +++ b/src/Squidex/app/shared/services/auth.service.ts @@ -161,7 +161,12 @@ export class AuthService { }); }); - return observable.timeout(1000).retryWhen(errors => errors.filter(e => e instanceof TimeoutError)); + return observable.timeout(2000) + .retryWhen(errors => errors + .filter(e => e instanceof TimeoutError) + .delay(500) + .take(5) + .concat(Observable.throw(new Error('Retry limit exceeeded.')))); } private createProfile(user: User) { diff --git a/tests/Squidex.Infrastructure.Tests/Timers/CompletionTimerTests.cs b/tests/Squidex.Infrastructure.Tests/Timers/CompletionTimerTests.cs index e445ff387..079baf2bf 100644 --- a/tests/Squidex.Infrastructure.Tests/Timers/CompletionTimerTests.cs +++ b/tests/Squidex.Infrastructure.Tests/Timers/CompletionTimerTests.cs @@ -10,6 +10,7 @@ using System.Threading; using Squidex.Infrastructure.Tasks; using Xunit; +// ReSharper disable MethodSupportsCancellation // ReSharper disable AccessToModifiedClosure namespace Squidex.Infrastructure.Timers @@ -28,8 +29,8 @@ namespace Squidex.Infrastructure.Timers return TaskHelper.Done; }, 2000); - timer.Wakeup(); - timer.Dispose(); + timer.SkipCurrentDelay(); + timer.StopAsync().Wait(); Assert.True(called); } @@ -40,15 +41,15 @@ namespace Squidex.Infrastructure.Timers timer = new CompletionTimer(10, ct => { - timer?.Dispose(); + timer?.StopAsync().Wait(); return TaskHelper.Done; }, 10); Thread.Sleep(1000); - timer.Wakeup(); - timer.Dispose(); + timer.SkipCurrentDelay(); + timer.StopAsync().Wait(); } } }