Browse Source

A lot of fixes for ugly race conditions.

pull/130/head
Sebastian Stehle 9 years ago
parent
commit
020c8d0173
  1. 11
      src/Squidex.Infrastructure.GetEventStore/GetEventStoreSubscription.cs
  2. 4
      src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs
  3. 28
      src/Squidex.Infrastructure.MongoDb/EventStore/PollingSubscription.cs
  4. 4
      src/Squidex.Infrastructure/CQRS/Events/DefaultEventNotifier.cs
  5. 2
      src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs
  6. 2
      src/Squidex.Infrastructure/CQRS/Events/IEventNotifier.cs
  7. 125
      src/Squidex/appsettings.json

11
src/Squidex.Infrastructure.GetEventStore/GetEventStoreSubscription.cs

@ -19,7 +19,7 @@ using Squidex.Infrastructure.CQRS.Events;
namespace Squidex.Infrastructure.GetEventStore namespace Squidex.Infrastructure.GetEventStore
{ {
internal sealed class EventStoreSubscription : IEventSubscription internal sealed class EventStoreSubscription : DisposableObjectBase, IEventSubscription
{ {
private static readonly ConcurrentDictionary<string, bool> subscriptionsCreated = new ConcurrentDictionary<string, bool>(); private static readonly ConcurrentDictionary<string, bool> subscriptionsCreated = new ConcurrentDictionary<string, bool>();
private readonly IEventStoreConnection connection; private readonly IEventStoreConnection connection;
@ -40,10 +40,13 @@ namespace Squidex.Infrastructure.GetEventStore
streamName = $"by-{prefix.Simplify()}-{streamFilter.Simplify()}"; 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<StoredEvent, Task> onNext, Func<Exception, Task> onError = null) public async Task SubscribeAsync(Func<StoredEvent, Task> onNext, Func<Exception, Task> onError = null)

4
src/Squidex.Infrastructure.MongoDb/EventStore/MongoEventStore.cs

@ -26,6 +26,7 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
{ {
public class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IEventStore, IDisposable public class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IEventStore, IDisposable
{ {
private static readonly BsonTimestamp EmptyTimestamp = new BsonTimestamp(0);
private static readonly FieldDefinition<MongoEventCommit, BsonTimestamp> TimestampField = Fields.Build(x => x.Timestamp); private static readonly FieldDefinition<MongoEventCommit, BsonTimestamp> TimestampField = Fields.Build(x => x.Timestamp);
private static readonly FieldDefinition<MongoEventCommit, long> EventsCountField = Fields.Build(x => x.EventsCount); private static readonly FieldDefinition<MongoEventCommit, long> EventsCountField = Fields.Build(x => x.EventsCount);
private static readonly FieldDefinition<MongoEventCommit, long> EventStreamOffsetField = Fields.Build(x => x.EventStreamOffset); private static readonly FieldDefinition<MongoEventCommit, long> EventStreamOffsetField = Fields.Build(x => x.EventStreamOffset);
@ -209,7 +210,8 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
Events = commitEvents, Events = commitEvents,
EventsCount = eventsCount, EventsCount = eventsCount,
EventStream = streamName, EventStream = streamName,
EventStreamOffset = expectedVersion EventStreamOffset = expectedVersion,
Timestamp = EmptyTimestamp
}.ToBsonDocument(); }.ToBsonDocument();
pendingCommits.Enqueue((document, cts)); pendingCommits.Enqueue((document, cts));

28
src/Squidex.Infrastructure.MongoDb/EventStore/PollingSubscription.cs

@ -12,6 +12,8 @@ using Squidex.Infrastructure.CQRS.Events;
using Squidex.Infrastructure.Tasks; using Squidex.Infrastructure.Tasks;
using Squidex.Infrastructure.Timers; using Squidex.Infrastructure.Timers;
// ReSharper disable InvertIf
namespace Squidex.Infrastructure.MongoDb.EventStore namespace Squidex.Infrastructure.MongoDb.EventStore
{ {
public sealed class PollingSubscription : DisposableObjectBase, IEventSubscription public sealed class PollingSubscription : DisposableObjectBase, IEventSubscription
@ -19,7 +21,8 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
private readonly IEventNotifier eventNotifier; private readonly IEventNotifier eventNotifier;
private readonly MongoEventStore eventStore; private readonly MongoEventStore eventStore;
private readonly string streamFilter; private readonly string streamFilter;
private readonly string position; private string position;
private IDisposable subscription;
private CompletionTimer timer; private CompletionTimer timer;
public PollingSubscription(MongoEventStore eventStore, IEventNotifier eventNotifier, string streamFilter, string position) public PollingSubscription(MongoEventStore eventStore, IEventNotifier eventNotifier, string streamFilter, string position)
@ -34,7 +37,9 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
{ {
if (disposing) if (disposing)
{ {
timer.Dispose(); subscription?.Dispose();
timer?.Dispose();
} }
} }
@ -42,7 +47,7 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
{ {
Guard.NotNull(onNext, nameof(onNext)); Guard.NotNull(onNext, nameof(onNext));
if (timer == null) if (timer != null)
{ {
throw new InvalidOperationException("An handler has already been registered."); throw new InvalidOperationException("An handler has already been registered.");
} }
@ -51,15 +56,26 @@ namespace Squidex.Infrastructure.MongoDb.EventStore
{ {
try 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); onError?.Invoke(ex);
} }
}); });
eventNotifier.Subscribe(timer.Wakeup); subscription = eventNotifier.Subscribe(() =>
{
if (!timer.IsDisposed)
{
timer.Wakeup();
}
});
return TaskHelper.Done; return TaskHelper.Done;
} }

4
src/Squidex.Infrastructure/CQRS/Events/DefaultEventNotifier.cs

@ -28,9 +28,9 @@ namespace Squidex.Infrastructure.CQRS.Events
invalidator.Publish(ChannelName, string.Empty, true); 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());
} }
} }
} }

2
src/Squidex.Infrastructure/CQRS/Events/EventReceiver.cs

@ -105,7 +105,7 @@ namespace Squidex.Infrastructure.CQRS.Events
var position = status.Position; var position = status.Position;
if (status.IsResetting || status.IsStopped) if (status.IsResetting)
{ {
currentSubscription?.Dispose(); currentSubscription?.Dispose();
currentSubscription = null; currentSubscription = null;

2
src/Squidex.Infrastructure/CQRS/Events/IEventNotifier.cs

@ -14,6 +14,6 @@ namespace Squidex.Infrastructure.CQRS.Events
{ {
void NotifyEventsStored(); void NotifyEventsStored();
void Subscribe(Action handler); IDisposable Subscribe(Action handler);
} }
} }

125
src/Squidex/appsettings.json

@ -1,70 +1,187 @@
{ {
"urls": { "urls": {
/*
* Set the base url of your application, to generate correct urls in background process.
*/
"baseUrl": "http://localhost:5000" "baseUrl": "http://localhost:5000"
}, },
"logging": { "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": { "pubSub": {
/*
* Define the type of the read store.
*
* Supported: InMemory (for single node only), Redis (for cluster)
*/
"type": "InMemory", "type": "InMemory",
"redis": { "redis": {
/*
* Connection string to your redis server.
*
* Read More: https://github.com/ServiceStack/ServiceStack.Redis#redis-connection-strings
*/
"configuration": "localhost:6379,resolveDns=1" "configuration": "localhost:6379,resolveDns=1"
} }
}, },
"assetStore": { "assetStore": {
/*
* Define the type of the read store.
*
* Supported: Folder (local folder), GoogleCloud (hosted in Google Cloud only)
*/
"type": "Folder", "type": "Folder",
"folder": { "folder": {
/*
* The relative or absolute path to the folder to store the assets.
*/
"path": "Assets" "path": "Assets"
}, },
"googleCloud": { "googleCloud": {
/*
* The name of the bucket in google cloud store.
*/
"bucket": "squidex-assets" "bucket": "squidex-assets"
} }
}, },
"eventStore": { "eventStore": {
/*
* Define the type of the event store.
*
* Supported: MongoDb, GetEventStore
*/
"type": "MongoDb", "type": "MongoDb",
"mongoDb": { "mongoDb": {
/*
* The connection string to your Mongo Server.
*
* Read More: https://docs.mongodb.com/manual/reference/connection-string/
*/
"configuration": "mongodb://localhost", "configuration": "mongodb://localhost",
/*
* The name of the event store database.
*/
"database": "Squidex" "database": "Squidex"
}, },
"getEventStore": { "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", "configuration": "ConnectTo=tcp://admin:changeit@localhost:1113; HeartBeatTimeout=500",
/*
* The host name of your EventStore where projection requests will be sent to.
*/
"projectionHost": "localhost", "projectionHost": "localhost",
/*
* Prefix for all streams and projections (for multiple installations).
*/
"prefix": "squidex" "prefix": "squidex"
}, },
/*
* Consume the events on this server (Ensure that it is only enabled on a single node).
*/
"consume": true "consume": true
}, },
"eventPublishers": { "eventPublishers": {
/*
* Additional event publishers (advanced usage only): (Name => Config)
*/
"allToRabbitMq": { "allToRabbitMq": {
/*
* Example:: Push all events to RabbitMq.
*/
"type": "RabbitMq", "type": "RabbitMq",
"configuration": "amqp://guest:guest@localhost/", "configuration": "amqp://guest:guest@localhost/",
"exchange": "squidex", "exchange": "squidex",
"enabled": false, "enabled": false,
"eventsFilter": "*" "eventsFilter": ".*"
} }
}, },
"store": { "store": {
/*
* Define the type of the read store.
*
* Supported: MongoDb
*/
"type": "MongoDb", "type": "MongoDb",
"mongoDb": { "mongoDb": {
/*
* The connection string to your Mongo Server.
*
* Read More: https://docs.mongodb.com/manual/reference/connection-string/
*/
"configuration": "mongodb://localhost", "configuration": "mongodb://localhost",
/*
* The database for all your content collections (one collection per app).
*/
"contentDatabase": "SquidexContent", "contentDatabase": "SquidexContent",
/*
* The database for all your other read collections.
*/
"database": "Squidex" "database": "Squidex"
} }
}, },
"identity": { "identity": {
/*
* Enable password auth.
*/
"allowPasswordAuth": true, "allowPasswordAuth": true,
/*
* Settings for Google auth (keep empty to disable).
*/
"googleClient": "1006817248705-t3lb3ge808m9am4t7upqth79hulk456l.apps.googleusercontent.com", "googleClient": "1006817248705-t3lb3ge808m9am4t7upqth79hulk456l.apps.googleusercontent.com",
"googleSecret": "QsEi-fHqkGw2_PjJmtNHf2wg", "googleSecret": "QsEi-fHqkGw2_PjJmtNHf2wg",
/*
* Settings for Github auth (keep empty to disable).
*/
"githubClient": "211ea00e726baf754c78", "githubClient": "211ea00e726baf754c78",
"githubSecret": "d0a0d0fe2c26469ae20987ac265b3a339fd73132", "githubSecret": "d0a0d0fe2c26469ae20987ac265b3a339fd73132",
/*
* Settings for Microsoft auth (keep empty to disable).
*/
"microsoftClient": "b55da740-6648-4502-8746-b9003f29d5f1", "microsoftClient": "b55da740-6648-4502-8746-b9003f29d5f1",
"microsoftSecret": "idWbANxNYEF4cB368WXJhjN", "microsoftSecret": "idWbANxNYEF4cB368WXJhjN",
/*
* Lock new users automatically, the administrator must unlock them.
*/
"lockAutomatically": true, "lockAutomatically": true,
"keysStore": { "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", "type": "InMemory",
"redis": { "redis": {
/*
* Connection string to your redis server.
*
* Read More: https://github.com/ServiceStack/ServiceStack.Redis#redis-connection-strings
*/
"configuration": "localhost:6379,resolveDns=1" "configuration": "localhost:6379,resolveDns=1"
}, },
"folder": { "folder": {
/*
* Relative or absolute path to your encryption keys folder.
*/
"path": "keys" "path": "keys"
} }
} }

Loading…
Cancel
Save