From 020c8d017369dc5002c9e26cbdc2d6dd87629199 Mon Sep 17 00:00:00 2001 From: Sebastian Stehle Date: Sat, 15 Jul 2017 13:56:39 +0200 Subject: [PATCH] A lot of fixes for ugly race conditions. --- .../GetEventStoreSubscription.cs | 11 +- .../EventStore/MongoEventStore.cs | 4 +- .../EventStore/PollingSubscription.cs | 28 +++- .../CQRS/Events/DefaultEventNotifier.cs | 4 +- .../CQRS/Events/EventReceiver.cs | 2 +- .../CQRS/Events/IEventNotifier.cs | 2 +- src/Squidex/appsettings.json | 125 +++++++++++++++++- 7 files changed, 157 insertions(+), 19 deletions(-) diff --git a/src/Squidex.Infrastructure.GetEventStore/GetEventStoreSubscription.cs b/src/Squidex.Infrastructure.GetEventStore/GetEventStoreSubscription.cs index 9bb4fa6c3..87a564e94 100644 --- a/src/Squidex.Infrastructure.GetEventStore/GetEventStoreSubscription.cs +++ b/src/Squidex.Infrastructure.GetEventStore/GetEventStoreSubscription.cs @@ -19,7 +19,7 @@ using Squidex.Infrastructure.CQRS.Events; namespace Squidex.Infrastructure.GetEventStore { - internal sealed class EventStoreSubscription : IEventSubscription + internal sealed class EventStoreSubscription : DisposableObjectBase, IEventSubscription { private static readonly ConcurrentDictionary subscriptionsCreated = new ConcurrentDictionary(); private readonly IEventStoreConnection connection; @@ -40,10 +40,13 @@ namespace Squidex.Infrastructure.GetEventStore streamName = $"by-{prefix.Simplify()}-{streamFilter.Simplify()}"; } - - public void Dispose() + + protected override void DisposeObject(bool disposing) { - internalSubscription?.Stop(); + if (disposing) + { + internalSubscription?.Stop(); + } } public async Task SubscribeAsync(Func onNext, Func onError = null) diff --git a/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs b/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs index 49fe917f3..b573c6074 100644 --- a/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs +++ b/src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs @@ -26,6 +26,7 @@ namespace Squidex.Infrastructure.MongoDb.EventStore { public class MongoEventStore : MongoRepositoryBase, IEventStore, IDisposable { + private static readonly BsonTimestamp EmptyTimestamp = new BsonTimestamp(0); private static readonly FieldDefinition TimestampField = Fields.Build(x => x.Timestamp); private static readonly FieldDefinition EventsCountField = Fields.Build(x => x.EventsCount); private static readonly FieldDefinition EventStreamOffsetField = Fields.Build(x => x.EventStreamOffset); @@ -209,7 +210,8 @@ namespace Squidex.Infrastructure.MongoDb.EventStore Events = commitEvents, EventsCount = eventsCount, EventStream = streamName, - EventStreamOffset = expectedVersion + EventStreamOffset = expectedVersion, + Timestamp = EmptyTimestamp }.ToBsonDocument(); pendingCommits.Enqueue((document, cts)); diff --git a/src/Squidex.Infrastructure.MongoDb/EventStore/PollingSubscription.cs b/src/Squidex.Infrastructure.MongoDb/EventStore/PollingSubscription.cs index b6b0162c9..921cd6187 100644 --- a/src/Squidex.Infrastructure.MongoDb/EventStore/PollingSubscription.cs +++ b/src/Squidex.Infrastructure.MongoDb/EventStore/PollingSubscription.cs @@ -12,6 +12,8 @@ using Squidex.Infrastructure.CQRS.Events; using Squidex.Infrastructure.Tasks; using Squidex.Infrastructure.Timers; +// ReSharper disable InvertIf + namespace Squidex.Infrastructure.MongoDb.EventStore { public sealed class PollingSubscription : DisposableObjectBase, IEventSubscription @@ -19,7 +21,8 @@ namespace Squidex.Infrastructure.MongoDb.EventStore private readonly IEventNotifier eventNotifier; private readonly MongoEventStore eventStore; private readonly string streamFilter; - private readonly string position; + private string position; + private IDisposable subscription; private CompletionTimer timer; public PollingSubscription(MongoEventStore eventStore, IEventNotifier eventNotifier, string streamFilter, string position) @@ -34,7 +37,9 @@ namespace Squidex.Infrastructure.MongoDb.EventStore { if (disposing) { - timer.Dispose(); + subscription?.Dispose(); + + timer?.Dispose(); } } @@ -42,7 +47,7 @@ namespace Squidex.Infrastructure.MongoDb.EventStore { Guard.NotNull(onNext, nameof(onNext)); - if (timer == null) + if (timer != null) { throw new InvalidOperationException("An handler has already been registered."); } @@ -51,15 +56,26 @@ namespace Squidex.Infrastructure.MongoDb.EventStore { try { - await eventStore.GetEventsAsync(onNext, ct, streamFilter, position); + await eventStore.GetEventsAsync(async storedEvent => + { + await onNext(storedEvent); + + position = storedEvent.EventPosition; + }, ct, streamFilter, position); } - catch (Exception ex) + catch (Exception ex) when (!(ex is OperationCanceledException)) { onError?.Invoke(ex); } }); - eventNotifier.Subscribe(timer.Wakeup); + subscription = eventNotifier.Subscribe(() => + { + if (!timer.IsDisposed) + { + timer.Wakeup(); + } + }); return TaskHelper.Done; } diff --git a/src/Squidex.Infrastructure/CQRS/Events/DefaultEventNotifier.cs b/src/Squidex.Infrastructure/CQRS/Events/DefaultEventNotifier.cs index 4d2207875..0513a3a3d 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/DefaultEventNotifier.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/DefaultEventNotifier.cs @@ -28,9 +28,9 @@ namespace Squidex.Infrastructure.CQRS.Events invalidator.Publish(ChannelName, string.Empty, true); } - public void Subscribe(Action handler) + public IDisposable Subscribe(Action handler) { - invalidator.Subscribe(ChannelName, x => handler()); + return invalidator.Subscribe(ChannelName, x => handler()); } } } diff --git a/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs b/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs index 5ff162d9c..ca1ea7947 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs @@ -105,7 +105,7 @@ namespace Squidex.Infrastructure.CQRS.Events var position = status.Position; - if (status.IsResetting || status.IsStopped) + if (status.IsResetting) { currentSubscription?.Dispose(); currentSubscription = null; diff --git a/src/Squidex.Infrastructure/CQRS/Events/IEventNotifier.cs b/src/Squidex.Infrastructure/CQRS/Events/IEventNotifier.cs index 43c586d3b..758bc94f7 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/IEventNotifier.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/IEventNotifier.cs @@ -14,6 +14,6 @@ namespace Squidex.Infrastructure.CQRS.Events { void NotifyEventsStored(); - void Subscribe(Action handler); + IDisposable Subscribe(Action handler); } } diff --git a/src/Squidex/appsettings.json b/src/Squidex/appsettings.json index bd52c76f4..5ba587f8d 100644 --- a/src/Squidex/appsettings.json +++ b/src/Squidex/appsettings.json @@ -1,70 +1,187 @@ { "urls": { + /* + * Set the base url of your application, to generate correct urls in background process. + */ "baseUrl": "http://localhost:5000" }, + "logging": { - "human": false + /* + * Setting the flag to true, enables well formatteds json logs. + */ + "human": false }, + + /* + * The pub sub mechanmism distributes messages between the nodes. + */ "pubSub": { + /* + * Define the type of the read store. + * + * Supported: InMemory (for single node only), Redis (for cluster) + */ "type": "InMemory", "redis": { + /* + * Connection string to your redis server. + * + * Read More: https://github.com/ServiceStack/ServiceStack.Redis#redis-connection-strings + */ "configuration": "localhost:6379,resolveDns=1" } }, + "assetStore": { + /* + * Define the type of the read store. + * + * Supported: Folder (local folder), GoogleCloud (hosted in Google Cloud only) + */ "type": "Folder", "folder": { + /* + * The relative or absolute path to the folder to store the assets. + */ "path": "Assets" }, "googleCloud": { + /* + * The name of the bucket in google cloud store. + */ "bucket": "squidex-assets" } }, + "eventStore": { + /* + * Define the type of the event store. + * + * Supported: MongoDb, GetEventStore + */ "type": "MongoDb", "mongoDb": { + /* + * The connection string to your Mongo Server. + * + * Read More: https://docs.mongodb.com/manual/reference/connection-string/ + */ "configuration": "mongodb://localhost", + /* + * The name of the event store database. + */ "database": "Squidex" }, "getEventStore": { + /* + * The connection string to your EventStore. + * + * Read Mode: http://docs.geteventstore.com/dotnet-api/4.0.0/connecting-to-a-server/ + */ "configuration": "ConnectTo=tcp://admin:changeit@localhost:1113; HeartBeatTimeout=500", + /* + * The host name of your EventStore where projection requests will be sent to. + */ "projectionHost": "localhost", + /* + * Prefix for all streams and projections (for multiple installations). + */ "prefix": "squidex" }, + /* + * Consume the events on this server (Ensure that it is only enabled on a single node). + */ "consume": true }, + "eventPublishers": { + /* + * Additional event publishers (advanced usage only): (Name => Config) + */ "allToRabbitMq": { + /* + * Example:: Push all events to RabbitMq. + */ "type": "RabbitMq", "configuration": "amqp://guest:guest@localhost/", "exchange": "squidex", "enabled": false, - "eventsFilter": "*" - } - }, + "eventsFilter": ".*" + } + }, + "store": { + /* + * Define the type of the read store. + * + * Supported: MongoDb + */ "type": "MongoDb", "mongoDb": { + /* + * The connection string to your Mongo Server. + * + * Read More: https://docs.mongodb.com/manual/reference/connection-string/ + */ "configuration": "mongodb://localhost", + /* + * The database for all your content collections (one collection per app). + */ "contentDatabase": "SquidexContent", + /* + * The database for all your other read collections. + */ "database": "Squidex" } }, + "identity": { + /* + * Enable password auth. + */ "allowPasswordAuth": true, + /* + * Settings for Google auth (keep empty to disable). + */ "googleClient": "1006817248705-t3lb3ge808m9am4t7upqth79hulk456l.apps.googleusercontent.com", "googleSecret": "QsEi-fHqkGw2_PjJmtNHf2wg", + /* + * Settings for Github auth (keep empty to disable). + */ "githubClient": "211ea00e726baf754c78", "githubSecret": "d0a0d0fe2c26469ae20987ac265b3a339fd73132", + /* + * Settings for Microsoft auth (keep empty to disable). + */ "microsoftClient": "b55da740-6648-4502-8746-b9003f29d5f1", "microsoftSecret": "idWbANxNYEF4cB368WXJhjN", + /* + * Lock new users automatically, the administrator must unlock them. + */ "lockAutomatically": true, + "keysStore": { + /* + * Define the type of the key store. + * + * Supported: InMemory (development only), Folder (shared or local), Redis + * + * Read More: https://docs.microsoft.com/en-us/aspnet/core/security/data-protection/implementation/key-storage-providers + */ "type": "InMemory", "redis": { + /* + * Connection string to your redis server. + * + * Read More: https://github.com/ServiceStack/ServiceStack.Redis#redis-connection-strings + */ "configuration": "localhost:6379,resolveDns=1" }, "folder": { + /* + * Relative or absolute path to your encryption keys folder. + */ "path": "keys" } }