Browse Source

Last errors fixed

pull/95/head
Sebastian Stehle 9 years ago
parent
commit
aca19d09f4
  1. 21
      src/Squidex.Infrastructure.MongoDb/CQRS/Events/PollingSubscription.cs
  2. 4
      src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs
  3. 50
      src/Squidex.Infrastructure/Timers/CompletionTimer.cs
  4. 4
      src/Squidex.Infrastructure/UsageTracking/BackgroundUsageTracker.cs
  5. 4
      src/Squidex/Config/Domain/ReadModule.cs
  6. 7
      src/Squidex/app/shared/services/auth.service.ts
  7. 11
      tests/Squidex.Infrastructure.Tests/Timers/CompletionTimerTests.cs

21
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();
}
});

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

50
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<CancellationToken, Task> 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<CancellationToken, Task> 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);
}

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

4
src/Squidex/Config/Domain/ReadModule.cs

@ -80,6 +80,10 @@ namespace Squidex.Config.Domain
.As<IHistoryEventsCreator>()
.SingleInstance();
builder.RegisterType<NoopAppPlanBillingManager>()
.As<IAppPlanBillingManager>()
.InstancePerDependency();
builder.RegisterType<WebhookInvoker>()
.As<IEventConsumer>()
.AsSelf()

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

11
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();
}
}
}

Loading…
Cancel
Save