Browse Source

Tests improved.

pull/206/head
Sebastian Stehle 9 years ago
parent
commit
0ec9c7a66b
  1. 10
      src/Squidex.Domain.Apps.Entities.MongoDb/Apps/MongoAppEntity.cs
  2. 19
      src/Squidex.Domain.Apps.Entities.MongoDb/Apps/MongoAppRepository.cs
  3. 4
      src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository.cs
  4. 8
      src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs
  5. 15
      src/Squidex.Domain.Apps.Entities.MongoDb/Rules/MongoRuleEntity.cs
  6. 16
      src/Squidex.Domain.Apps.Entities.MongoDb/Rules/MongoRuleRepository.cs
  7. 13
      src/Squidex.Domain.Apps.Entities.MongoDb/Schemas/MongoSchemaEntity.cs
  8. 16
      src/Squidex.Domain.Apps.Entities.MongoDb/Schemas/MongoSchemaRepository.cs
  9. 68
      src/Squidex.Domain.Apps.Entities/AppProvider.cs
  10. 9
      src/Squidex.Domain.Apps.Entities/EntityMapper.cs
  11. 13
      src/Squidex.Infrastructure.MongoDb/MongoDb/MongoExtensions.cs
  12. 2
      src/Squidex.Infrastructure.MongoDb/States/MongoSnapshotStore.cs
  13. 2
      src/Squidex.Infrastructure/Commands/AggregateHandler.cs
  14. 34
      src/Squidex.Infrastructure/Commands/DomainObjectBase.cs
  15. 5
      src/Squidex.Infrastructure/States/IPersistence.cs
  16. 36
      src/Squidex.Infrastructure/States/Persistence.cs
  17. 16
      src/Squidex/Config/Domain/StoreServices.cs
  18. 49
      tests/Squidex.Infrastructure.Tests/Commands/AggregateHandlerTests.cs
  19. 93
      tests/Squidex.Infrastructure.Tests/Commands/DomainObjectBaseTests.cs
  20. 14
      tests/Squidex.Infrastructure.Tests/Commands/TestHelpers/MyDomainObject.cs
  21. 8
      tests/Squidex.Infrastructure.Tests/DispatchingTests.cs
  22. 22
      tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs
  23. 10
      tests/Squidex.Infrastructure.Tests/States/StateEventSourcingTests.cs
  24. 4
      tests/Squidex.Infrastructure.Tests/Tasks/SingleThreadedDispatcherTests.cs

10
src/Squidex.Domain.Apps.Entities.MongoDb/Apps/MongoAppEntity.cs

@ -9,6 +9,7 @@
using MongoDB.Bson; using MongoDB.Bson;
using MongoDB.Bson.Serialization.Attributes; using MongoDB.Bson.Serialization.Attributes;
using Squidex.Domain.Apps.Entities.Apps.State; using Squidex.Domain.Apps.Entities.Apps.State;
using Squidex.Infrastructure.MongoDb;
namespace Squidex.Domain.Apps.Entities.MongoDb.Apps namespace Squidex.Domain.Apps.Entities.MongoDb.Apps
{ {
@ -21,14 +22,19 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Apps
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public AppState State { get; set; } public int Version { get; set; }
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public int Version { get; set; } public string Name { get; set; }
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public string[] UserIds { get; set; } public string[] UserIds { get; set; }
[BsonJson]
[BsonElement]
[BsonRequired]
public AppState State { get; set; }
} }
} }

19
src/Squidex.Domain.Apps.Entities.MongoDb/Apps/MongoAppRepository.cs

@ -25,18 +25,24 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Apps
{ {
} }
protected override Task SetupCollectionAsync(IMongoCollection<MongoAppEntity> collection) protected override string CollectionName()
{ {
return collection.Indexes.CreateOneAsync(Index.Ascending(x => x.UserIds)); return "Snapshots_Apps";
}
protected override async Task SetupCollectionAsync(IMongoCollection<MongoAppEntity> collection)
{
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.UserIds));
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.Name));
} }
public async Task<Guid> FindAppIdByNameAsync(string name) public async Task<Guid> FindAppIdByNameAsync(string name)
{ {
var appEntity = var appEntity =
await Collection.Find(x => x.State.Name == name).Only(x => x.Id) await Collection.Find(x => x.Name == name).Only(x => x.Id)
.FirstOrDefaultAsync(); .FirstOrDefaultAsync();
return appEntity != null ? Guid.Parse(appEntity.Id) : Guid.Empty; return appEntity != null ? Guid.Parse(appEntity["_id"].AsString) : Guid.Empty;
} }
public async Task<IReadOnlyList<Guid>> QueryUserAppIdsAsync(string userId) public async Task<IReadOnlyList<Guid>> QueryUserAppIdsAsync(string userId)
@ -45,7 +51,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Apps
await Collection.Find(x => x.UserIds.Contains(userId)).Only(x => x.Id) await Collection.Find(x => x.UserIds.Contains(userId)).Only(x => x.Id)
.ToListAsync(); .ToListAsync();
return appEntities.Select(x => Guid.Parse(x.Id)).ToList(); return appEntities.Select(x => Guid.Parse(x["_id"].AsString)).ToList();
} }
public async Task<IReadOnlyList<string>> QueryUserAppNamesAsync(string userId) public async Task<IReadOnlyList<string>> QueryUserAppNamesAsync(string userId)
@ -75,9 +81,12 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Apps
{ {
try try
{ {
value.Version = newVersion;
await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion, await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion,
Update Update
.Set(x => x.UserIds, value.Contributors.Keys.ToArray()) .Set(x => x.UserIds, value.Contributors.Keys.ToArray())
.Set(x => x.Name, value.Name)
.Set(x => x.State, value) .Set(x => x.State, value)
.Set(x => x.Version, newVersion), .Set(x => x.Version, newVersion),
Upsert); Upsert);

4
src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository.cs

@ -116,6 +116,8 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Assets
{ {
try try
{ {
value.Version = newVersion;
await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion, await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion,
Update Update
.Set(x => x.State, value) .Set(x => x.State, value)
@ -132,7 +134,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Assets
if (existingVersion != null) if (existingVersion != null)
{ {
throw new InconsistentStateException(existingVersion.Version, oldVersion, ex); throw new InconsistentStateException(existingVersion["Version"].AsInt64, oldVersion, ex);
} }
} }
else else

8
src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs

@ -45,7 +45,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents
protected override string CollectionName() protected override string CollectionName()
{ {
return "Snapshots_Assets"; return "Snapshots_Contents";
} }
protected override async Task SetupCollectionAsync(IMongoCollection<MongoContentEntity> collection) protected override async Task SetupCollectionAsync(IMongoCollection<MongoContentEntity> collection)
@ -89,6 +89,8 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents
try try
{ {
value.Version = newVersion;
await Collection.InsertOneAsync(document); await Collection.InsertOneAsync(document);
} }
catch (MongoWriteException ex) catch (MongoWriteException ex)
@ -101,7 +103,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents
if (existingVersion != null) if (existingVersion != null)
{ {
throw new InconsistentStateException(existingVersion.Version, oldVersion, ex); throw new InconsistentStateException(existingVersion["Version"].AsInt64, oldVersion, ex);
} }
} }
else else
@ -206,7 +208,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents
await Collection.Find(x => contentIds.Contains(x.Id) && x.AppId == appId).Only(x => x.Id) await Collection.Find(x => contentIds.Contains(x.Id) && x.AppId == appId).Only(x => x.Id)
.ToListAsync(); .ToListAsync();
return contentIds.Except(contentEntities.Select(x => x.Id)).ToList(); return contentIds.Except(contentEntities.Select(x => Guid.Parse(x["_id"].AsString))).ToList();
} }
public async Task<IContentEntity> FindContentAsync(IAppEntity app, ISchemaEntity schema, Guid id, long version) public async Task<IContentEntity> FindContentAsync(IAppEntity app, ISchemaEntity schema, Guid id, long version)

15
src/Squidex.Domain.Apps.Entities.MongoDb/Rules/MongoRuleEntity.cs

@ -6,9 +6,11 @@
// All rights reserved. // All rights reserved.
// ========================================================================== // ==========================================================================
using System;
using MongoDB.Bson; using MongoDB.Bson;
using MongoDB.Bson.Serialization.Attributes; using MongoDB.Bson.Serialization.Attributes;
using Squidex.Domain.Apps.Entities.Rules.State; using Squidex.Domain.Apps.Entities.Rules.State;
using Squidex.Infrastructure.MongoDb;
namespace Squidex.Domain.Apps.Entities.MongoDb.Rules namespace Squidex.Domain.Apps.Entities.MongoDb.Rules
{ {
@ -21,10 +23,19 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Rules
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public RuleState State { get; set; } public int Version { get; set; }
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public int Version { get; set; } public Guid AppId { get; set; }
[BsonElement]
[BsonRequired]
public bool IsDeleted { get; set; }
[BsonJson]
[BsonElement]
[BsonRequired]
public RuleState State { get; set; }
} }
} }

16
src/Squidex.Domain.Apps.Entities.MongoDb/Rules/MongoRuleRepository.cs

@ -27,13 +27,13 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Rules
protected override string CollectionName() protected override string CollectionName()
{ {
return "States_Rules"; return "Snapshots_Rules";
} }
protected override async Task SetupCollectionAsync(IMongoCollection<MongoRuleEntity> collection) protected override async Task SetupCollectionAsync(IMongoCollection<MongoRuleEntity> collection)
{ {
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.State.AppId)); await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.AppId));
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.State.IsDeleted)); await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.IsDeleted));
} }
public async Task<(RuleState Value, long Version)> ReadAsync(string key) public async Task<(RuleState Value, long Version)> ReadAsync(string key)
@ -53,19 +53,23 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Rules
public async Task<IReadOnlyList<Guid>> QueryRuleIdsAsync(Guid appId) public async Task<IReadOnlyList<Guid>> QueryRuleIdsAsync(Guid appId)
{ {
var ruleEntities = var ruleEntities =
await Collection.Find(x => x.State.AppId == appId && !x.State.IsDeleted).Only(x => x.Id) await Collection.Find(x => x.AppId == appId && !x.IsDeleted).Only(x => x.Id)
.ToListAsync(); .ToListAsync();
return ruleEntities.Select(x => Guid.Parse(x.Id)).ToList(); return ruleEntities.Select(x => Guid.Parse(x["_id"].AsString)).ToList();
} }
public async Task WriteAsync(string key, RuleState value, long oldVersion, long newVersion) public async Task WriteAsync(string key, RuleState value, long oldVersion, long newVersion)
{ {
try try
{ {
value.Version = newVersion;
await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion, await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion,
Update Update
.Set(x => x.State, value) .Set(x => x.State, value)
.Set(x => x.AppId, value.AppId)
.Set(x => x.IsDeleted, value.IsDeleted)
.Set(x => x.Version, newVersion), .Set(x => x.Version, newVersion),
Upsert); Upsert);
} }
@ -79,7 +83,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Rules
if (existingVersion != null) if (existingVersion != null)
{ {
throw new InconsistentStateException(existingVersion.Version, oldVersion, ex); throw new InconsistentStateException(existingVersion["Version"].AsInt64, oldVersion, ex);
} }
} }
else else

13
src/Squidex.Domain.Apps.Entities.MongoDb/Schemas/MongoSchemaEntity.cs

@ -6,9 +6,11 @@
// All rights reserved. // All rights reserved.
// ========================================================================== // ==========================================================================
using System;
using MongoDB.Bson; using MongoDB.Bson;
using MongoDB.Bson.Serialization.Attributes; using MongoDB.Bson.Serialization.Attributes;
using Squidex.Domain.Apps.Entities.Schemas.State; using Squidex.Domain.Apps.Entities.Schemas.State;
using Squidex.Infrastructure.MongoDb;
namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas
{ {
@ -21,10 +23,19 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public SchemaState State { get; set; } public string Name { get; set; }
[BsonElement] [BsonElement]
[BsonRequired] [BsonRequired]
public int Version { get; set; } public int Version { get; set; }
[BsonElement]
[BsonRequired]
public Guid AppId { get; set; }
[BsonJson]
[BsonElement]
[BsonRequired]
public SchemaState State { get; set; }
} }
} }

16
src/Squidex.Domain.Apps.Entities.MongoDb/Schemas/MongoSchemaRepository.cs

@ -32,8 +32,8 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas
protected override async Task SetupCollectionAsync(IMongoCollection<MongoSchemaEntity> collection) protected override async Task SetupCollectionAsync(IMongoCollection<MongoSchemaEntity> collection)
{ {
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.State.AppId)); await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.AppId));
await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.State.Name)); await collection.Indexes.CreateOneAsync(Index.Ascending(x => x.Name));
} }
public async Task<(SchemaState Value, long Version)> ReadAsync(string key) public async Task<(SchemaState Value, long Version)> ReadAsync(string key)
@ -53,10 +53,10 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas
public async Task<Guid> FindSchemaIdAsync(Guid appId, string name) public async Task<Guid> FindSchemaIdAsync(Guid appId, string name)
{ {
var schemaEntity = var schemaEntity =
await Collection.Find(x => x.State.Name == name).Only(x => x.Id) await Collection.Find(x => x.Name == name).Only(x => x.Id)
.FirstOrDefaultAsync(); .FirstOrDefaultAsync();
return schemaEntity != null ? Guid.Parse(schemaEntity.Id) : Guid.Empty; return schemaEntity != null ? Guid.Parse(schemaEntity["_id"].AsString) : Guid.Empty;
} }
public async Task<IReadOnlyList<Guid>> QuerySchemaIdsAsync(Guid appId) public async Task<IReadOnlyList<Guid>> QuerySchemaIdsAsync(Guid appId)
@ -65,16 +65,20 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas
await Collection.Find(x => x.State.AppId == appId).Only(x => x.Id) await Collection.Find(x => x.State.AppId == appId).Only(x => x.Id)
.ToListAsync(); .ToListAsync();
return schemaEntities.Select(x => Guid.Parse(x.Id)).ToList(); return schemaEntities.Select(x => Guid.Parse(x["_id"].AsString)).ToList();
} }
public async Task WriteAsync(string key, SchemaState value, long oldVersion, long newVersion) public async Task WriteAsync(string key, SchemaState value, long oldVersion, long newVersion)
{ {
try try
{ {
value.Version = newVersion;
await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion, await Collection.UpdateOneAsync(x => x.Id == key && x.Version == oldVersion,
Update Update
.Set(x => x.State, value) .Set(x => x.State, value)
.Set(x => x.AppId, value.AppId)
.Set(x => x.Name, value.Name)
.Set(x => x.Version, newVersion), .Set(x => x.Version, newVersion),
Upsert); Upsert);
} }
@ -88,7 +92,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas
if (existingVersion != null) if (existingVersion != null)
{ {
throw new InconsistentStateException(existingVersion.Version, oldVersion, ex); throw new InconsistentStateException(existingVersion["Version"].AsInt64, oldVersion, ex);
} }
} }
else else

68
src/Squidex.Domain.Apps.Entities/AppProvider.cs

@ -24,8 +24,6 @@ namespace Squidex.Domain.Apps.Entities
{ {
public sealed class AppProvider : IAppProvider public sealed class AppProvider : IAppProvider
{ {
private readonly ConcurrentDictionary<string, Guid> appIds = new ConcurrentDictionary<string, Guid>();
private readonly ConcurrentDictionary<Tuple<Guid, string>, Guid> schemaIds = new ConcurrentDictionary<Tuple<Guid, string>, Guid>();
private readonly IAppRepository appRepository; private readonly IAppRepository appRepository;
private readonly IRuleRepository ruleRepository; private readonly IRuleRepository ruleRepository;
private readonly ISchemaRepository schemaRepository; private readonly ISchemaRepository schemaRepository;
@ -52,33 +50,23 @@ namespace Squidex.Domain.Apps.Entities
{ {
var app = await stateFactory.GetSingleAsync<AppDomainObject>(appId.ToString()); var app = await stateFactory.GetSingleAsync<AppDomainObject>(appId.ToString());
if (app.Version < 0) if (IsNotFound(app))
{ {
throw new DomainObjectNotFoundException(appId.ToString(), typeof(SchemaDomainObject)); return (null, null);
} }
var schema = await stateFactory.GetSingleAsync<SchemaDomainObject>(id.ToString()); var schema = await stateFactory.GetSingleAsync<SchemaDomainObject>(id.ToString());
if (schema.Version < 0 || schema.State.IsDeleted) return IsNotFound(false, schema) ? (null, null) : (app.State, schema.State);
{
throw new DomainObjectNotFoundException(id.ToString(), typeof(SchemaDomainObject));
}
return (app.State, schema.State);
} }
public async Task<IAppEntity> GetAppAsync(string appName) public async Task<IAppEntity> GetAppAsync(string appName)
{ {
var appId = await GetAppIdAsync(appName); var appId = await GetAppIdAsync(appName);
var app = await stateFactory.GetSingleAsync<AppDomainObject>(appName); var app = await stateFactory.GetSingleAsync<AppDomainObject>(appId.ToString());
if (app.Version < 0)
{
throw new DomainObjectNotFoundException(appName, typeof(SchemaDomainObject));
}
return app.State; return IsNotFound(app) ? null : app.State;
} }
public async Task<ISchemaEntity> GetSchemaAsync(Guid appId, string name, bool provideDeleted = false) public async Task<ISchemaEntity> GetSchemaAsync(Guid appId, string name, bool provideDeleted = false)
@ -87,24 +75,14 @@ namespace Squidex.Domain.Apps.Entities
var schema = await stateFactory.GetSingleAsync<SchemaDomainObject>(schemaId.ToString()); var schema = await stateFactory.GetSingleAsync<SchemaDomainObject>(schemaId.ToString());
if (schema.Version < 0 || (schema.State.IsDeleted && !provideDeleted)) return IsNotFound(provideDeleted, schema) ? null : schema.State;
{
throw new DomainObjectNotFoundException(schemaId.ToString(), typeof(SchemaDomainObject));
}
return schema.State;
} }
public async Task<ISchemaEntity> GetSchemaAsync(Guid appId, Guid id, bool provideDeleted = false) public async Task<ISchemaEntity> GetSchemaAsync(Guid appId, Guid id, bool provideDeleted = false)
{ {
var schema = await stateFactory.GetSingleAsync<SchemaDomainObject>(id.ToString()); var schema = await stateFactory.GetSingleAsync<SchemaDomainObject>(id.ToString());
if (schema.Version < 0 || (schema.State.IsDeleted && !provideDeleted)) return IsNotFound(provideDeleted, schema) ? null : schema.State;
{
throw new DomainObjectNotFoundException(id.ToString(), typeof(SchemaDomainObject));
}
return schema.State;
} }
public async Task<List<ISchemaEntity>> GetSchemasAsync(Guid appId) public async Task<List<ISchemaEntity>> GetSchemasAsync(Guid appId)
@ -140,32 +118,24 @@ namespace Squidex.Domain.Apps.Entities
return apps.Select(a => (IAppEntity)a.State).ToList(); return apps.Select(a => (IAppEntity)a.State).ToList();
} }
private async Task<Guid> GetAppIdAsync(string name) private Task<Guid> GetAppIdAsync(string name)
{ {
var key = name; return appRepository.FindAppIdByNameAsync(name);
if (!appIds.TryGetValue(key, out var id))
{
id = await appRepository.FindAppIdByNameAsync(name);
appIds[key] = id;
}
return id;
} }
private async Task<Guid> GetSchemaIdAsync(Guid appId, string name) private Task<Guid> GetSchemaIdAsync(Guid appId, string name)
{ {
var key = Tuple.Create(appId, name); return schemaRepository.FindSchemaIdAsync(appId, name);
}
if (!schemaIds.TryGetValue(key, out var id))
{
id = await schemaRepository.FindSchemaIdAsync(appId, name);
schemaIds[key] = id; private static bool IsNotFound(AppDomainObject app)
} {
return app.Version < 0;
}
return id; private static bool IsNotFound(bool provideDeleted, SchemaDomainObject schema)
{
return schema.Version < 0 || (schema.State.IsDeleted && !provideDeleted);
} }
} }
} }

9
src/Squidex.Domain.Apps.Entities/EntityMapper.cs

@ -23,7 +23,6 @@ namespace Squidex.Domain.Apps.Entities
SetCreatedBy(entity, @event); SetCreatedBy(entity, @event);
SetLastModified(entity, headers); SetLastModified(entity, headers);
SetLastModifiedBy(entity, @event); SetLastModifiedBy(entity, @event);
SetVersion(entity, headers);
updater?.Invoke(entity); updater?.Invoke(entity);
@ -38,14 +37,6 @@ namespace Squidex.Domain.Apps.Entities
} }
} }
private static void SetVersion(IEntity entity, EnvelopeHeaders headers)
{
if (entity is IUpdateableEntityWithVersion withVersion)
{
withVersion.Version = headers.EventStreamNumber();
}
}
private static void SetCreated(IEntity entity, EnvelopeHeaders headers) private static void SetCreated(IEntity entity, EnvelopeHeaders headers)
{ {
if (entity is IUpdateableEntity updateable && updateable.Created == default(Instant)) if (entity is IUpdateableEntity updateable && updateable.Created == default(Instant))

13
src/Squidex.Infrastructure.MongoDb/MongoDb/MongoExtensions.cs

@ -9,6 +9,7 @@
using System; using System;
using System.Linq.Expressions; using System.Linq.Expressions;
using System.Threading.Tasks; using System.Threading.Tasks;
using MongoDB.Bson;
using MongoDB.Driver; using MongoDB.Driver;
namespace Squidex.Infrastructure.MongoDb namespace Squidex.Infrastructure.MongoDb
@ -34,25 +35,25 @@ namespace Squidex.Infrastructure.MongoDb
return true; return true;
} }
public static IFindFluent<TDocument, TDocument> Only<TDocument>(this IFindFluent<TDocument, TDocument> find, public static IFindFluent<TDocument, BsonDocument> Only<TDocument>(this IFindFluent<TDocument, TDocument> find,
Expression<Func<TDocument, object>> include) Expression<Func<TDocument, object>> include)
{ {
return find.Project<TDocument>(Builders<TDocument>.Projection.Include(include)); return find.Project<BsonDocument>(Builders<TDocument>.Projection.Include(include));
} }
public static IFindFluent<TDocument, TDocument> Only<TDocument>(this IFindFluent<TDocument, TDocument> find, public static IFindFluent<TDocument, BsonDocument> Only<TDocument>(this IFindFluent<TDocument, TDocument> find,
Expression<Func<TDocument, object>> include1, Expression<Func<TDocument, object>> include1,
Expression<Func<TDocument, object>> include2) Expression<Func<TDocument, object>> include2)
{ {
return find.Project<TDocument>(Builders<TDocument>.Projection.Include(include1).Include(include2)); return find.Project<BsonDocument>(Builders<TDocument>.Projection.Include(include1).Include(include2));
} }
public static IFindFluent<TDocument, TDocument> Only<TDocument>(this IFindFluent<TDocument, TDocument> find, public static IFindFluent<TDocument, BsonDocument> Only<TDocument>(this IFindFluent<TDocument, TDocument> find,
Expression<Func<TDocument, object>> include1, Expression<Func<TDocument, object>> include1,
Expression<Func<TDocument, object>> include2, Expression<Func<TDocument, object>> include2,
Expression<Func<TDocument, object>> include3) Expression<Func<TDocument, object>> include3)
{ {
return find.Project<TDocument>(Builders<TDocument>.Projection.Include(include1).Include(include2).Include(include3)); return find.Project<BsonDocument>(Builders<TDocument>.Projection.Include(include1).Include(include2).Include(include3));
} }
} }
} }

2
src/Squidex.Infrastructure.MongoDb/States/MongoSnapshotStore.cs

@ -64,7 +64,7 @@ namespace Squidex.Infrastructure.States
if (existingVersion != null) if (existingVersion != null)
{ {
throw new InconsistentStateException(existingVersion.Version, oldVersion, ex); throw new InconsistentStateException(existingVersion["Version"].AsInt64, oldVersion, ex);
} }
} }
else else

2
src/Squidex.Infrastructure/Commands/AggregateHandler.cs

@ -69,6 +69,8 @@ namespace Squidex.Infrastructure.Commands
var domainObjectId = domainObjectCommand.AggregateId; var domainObjectId = domainObjectCommand.AggregateId;
var domainObject = await stateFactory.CreateAsync<T>(domainObjectId.ToString()); var domainObject = await stateFactory.CreateAsync<T>(domainObjectId.ToString());
await handler(domainObject);
await domainObject.WriteAsync(log); await domainObject.WriteAsync(log);
if (!context.IsCompleted) if (!context.IsCompleted)

34
src/Squidex.Infrastructure/Commands/DomainObjectBase.cs

@ -18,12 +18,13 @@ namespace Squidex.Infrastructure.Commands
public abstract class DomainObjectBase<TBase, TState> : IDomainObject where TState : new() public abstract class DomainObjectBase<TBase, TState> : IDomainObject where TState : new()
{ {
private readonly List<Envelope<IEvent>> uncomittedEvents = new List<Envelope<IEvent>>(); private readonly List<Envelope<IEvent>> uncomittedEvents = new List<Envelope<IEvent>>();
private Guid id;
private TState state = new TState(); private TState state = new TState();
private IPersistence<TState> persistence; private IPersistence<TState> persistence;
public long Version public long Version
{ {
get { return persistence.Version; } get { return persistence?.Version ?? -1; }
} }
public TState State public TState State
@ -43,6 +44,8 @@ namespace Squidex.Infrastructure.Commands
public Task ActivateAsync(string key, IStore store) public Task ActivateAsync(string key, IStore store)
{ {
id = Guid.Parse(key);
persistence = store.WithSnapshots<TBase, TState>(key, s => state = s); persistence = store.WithSnapshots<TBase, TState>(key, s => state = s);
return persistence.ReadAsync(); return persistence.ReadAsync();
@ -57,6 +60,8 @@ namespace Squidex.Infrastructure.Commands
{ {
Guard.NotNull(@event, nameof(@event)); Guard.NotNull(@event, nameof(@event));
@event.SetAggregateId(id);
OnRaised(@event.To<IEvent>()); OnRaised(@event.To<IEvent>());
uncomittedEvents.Add(@event.To<IEvent>()); uncomittedEvents.Add(@event.To<IEvent>());
@ -73,19 +78,24 @@ namespace Squidex.Infrastructure.Commands
public async Task WriteAsync(ISemanticLog log) public async Task WriteAsync(ISemanticLog log)
{ {
await persistence.WriteSnapshotAsync(state); var newVersion = Version + uncomittedEvents.Count;
try if (newVersion != Version)
{
await persistence.WriteEventsAsync(uncomittedEvents.ToArray());
}
catch (Exception ex)
{
log.LogFatal(ex, w => w.WriteProperty("action", "writeEvents"));
}
finally
{ {
uncomittedEvents.Clear(); await persistence.WriteSnapshotAsync(state, newVersion);
try
{
await persistence.WriteEventsAsync(uncomittedEvents.ToArray());
}
catch (Exception ex)
{
log.LogFatal(ex, w => w.WriteProperty("action", "writeEvents"));
}
finally
{
uncomittedEvents.Clear();
}
} }
} }
} }

5
src/Squidex.Infrastructure/States/IPersistence.cs

@ -6,6 +6,7 @@
// All rights reserved. // All rights reserved.
// ========================================================================== // ==========================================================================
using System.Collections.Generic;
using System.Threading.Tasks; using System.Threading.Tasks;
using Squidex.Infrastructure.EventSourcing; using Squidex.Infrastructure.EventSourcing;
@ -15,9 +16,9 @@ namespace Squidex.Infrastructure.States
{ {
long Version { get; } long Version { get; }
Task WriteEventsAsync(params Envelope<IEvent>[] @events); Task WriteEventsAsync(IEnumerable<Envelope<IEvent>> @events);
Task WriteSnapshotAsync(TState state); Task WriteSnapshotAsync(TState state, long newVersion = -1);
Task ReadAsync(long? expectedVersion = null); Task ReadAsync(long? expectedVersion = null);
} }

36
src/Squidex.Infrastructure/States/Persistence.cs

@ -7,6 +7,7 @@
// ========================================================================== // ==========================================================================
using System; using System;
using System.Collections.Generic;
using System.Linq; using System.Linq;
using System.Threading.Tasks; using System.Threading.Tasks;
using Squidex.Infrastructure.EventSourcing; using Squidex.Infrastructure.EventSourcing;
@ -55,7 +56,7 @@ namespace Squidex.Infrastructure.States
positionSnapshot = -1; positionSnapshot = -1;
positionEvent = -1; positionEvent = -1;
if (snapshotStore != null) if (applyState != null)
{ {
var (state, position) = await snapshotStore.ReadAsync(ownerKey); var (state, position) = await snapshotStore.ReadAsync(ownerKey);
@ -68,7 +69,7 @@ namespace Squidex.Infrastructure.States
} }
} }
if (eventStore != null && streamNameResolver != null) if (applyEvent != null && streamNameResolver != null)
{ {
var events = await eventStore.GetEventsAsync(GetStreamName(), positionEvent + 1); var events = await eventStore.GetEventsAsync(GetStreamName(), positionEvent + 1);
@ -105,51 +106,56 @@ namespace Squidex.Infrastructure.States
} }
} }
public async Task WriteSnapshotAsync(TState state) public async Task WriteSnapshotAsync(TState state, long newVersion = -1)
{ {
var newPosition = if (newVersion < 0)
eventStore != null ? {
positionEvent : newVersion =
positionSnapshot + 1; applyEvent != null ?
positionEvent :
positionSnapshot + 1;
}
if (newPosition != positionSnapshot) if (newVersion != positionSnapshot)
{ {
try try
{ {
await snapshotStore.WriteAsync(ownerKey, state, positionSnapshot, newPosition); await snapshotStore.WriteAsync(ownerKey, state, positionSnapshot, newVersion);
} }
catch (InconsistentStateException ex) catch (InconsistentStateException ex)
{ {
throw new DomainObjectVersionException(ownerKey, typeof(TOwner), ex.CurrentVersion, ex.ExpectedVersion); throw new DomainObjectVersionException(ownerKey, typeof(TOwner), ex.CurrentVersion, ex.ExpectedVersion);
} }
positionSnapshot = newPosition; positionSnapshot = newVersion;
} }
invalidate?.Invoke(); invalidate?.Invoke();
} }
public async Task WriteEventsAsync(params Envelope<IEvent>[] @events) public async Task WriteEventsAsync(IEnumerable<Envelope<IEvent>> events)
{ {
Guard.NotNull(events, nameof(@events)); Guard.NotNull(events, nameof(@events));
if (@events.Length > 0) var eventArray = events.ToArray();
if (eventArray.Length > 0)
{ {
var commitId = Guid.NewGuid(); var commitId = Guid.NewGuid();
var eventStream = GetStreamName(); var eventStream = GetStreamName();
var eventData = GetEventData(events, commitId); var eventData = GetEventData(eventArray, commitId);
try try
{ {
await eventStore.AppendEventsAsync(commitId, GetStreamName(), positionEvent, eventData); await eventStore.AppendEventsAsync(commitId, GetStreamName(), Version, eventData);
} }
catch (WrongEventVersionException ex) catch (WrongEventVersionException ex)
{ {
throw new DomainObjectVersionException(ownerKey, typeof(TOwner), ex.CurrentVersion, ex.ExpectedVersion); throw new DomainObjectVersionException(ownerKey, typeof(TOwner), ex.CurrentVersion, ex.ExpectedVersion);
} }
positionEvent += events.Length; positionEvent += eventArray.Length;
} }
invalidate?.Invoke(); invalidate?.Invoke();

16
src/Squidex/Config/Domain/StoreServices.cs

@ -31,6 +31,7 @@ using Squidex.Domain.Apps.Entities.MongoDb.Schemas;
using Squidex.Domain.Apps.Entities.Rules.Repositories; using Squidex.Domain.Apps.Entities.Rules.Repositories;
using Squidex.Domain.Apps.Entities.Rules.State; using Squidex.Domain.Apps.Entities.Rules.State;
using Squidex.Domain.Apps.Entities.Schemas.Repositories; using Squidex.Domain.Apps.Entities.Schemas.Repositories;
using Squidex.Domain.Apps.Entities.Schemas.State;
using Squidex.Domain.Users; using Squidex.Domain.Users;
using Squidex.Domain.Users.MongoDb; using Squidex.Domain.Users.MongoDb;
using Squidex.Domain.Users.MongoDb.Infrastructure; using Squidex.Domain.Users.MongoDb.Infrastructure;
@ -100,19 +101,20 @@ namespace Squidex.Config.Domain
.As<ISnapshotStore<AssetState>>() .As<ISnapshotStore<AssetState>>()
.As<IExternalSystem>(); .As<IExternalSystem>();
services.AddSingletonAs(c => new MongoContentRepository(mongoContentDatabase, c.GetService<IAppProvider>()))
.As<IContentRepository>()
.As<ISnapshotStore<ContentState>>()
.As<IEventConsumer>();
services.AddSingletonAs(c => new MongoRuleRepository(mongoContentDatabase)) services.AddSingletonAs(c => new MongoRuleRepository(mongoContentDatabase))
.As<IRuleRepository>() .As<IRuleRepository>()
.As<ISnapshotStore<RuleState>>() .As<ISnapshotStore<RuleState>>()
.As<IEventConsumer>(); .As<IExternalSystem>();
services.AddSingletonAs(c => new MongoSchemaRepository(mongoDatabase)) services.AddSingletonAs(c => new MongoSchemaRepository(mongoDatabase))
.As<ISchemaRepository>() .As<ISchemaRepository>()
.As<ISnapshotStore<AppState>>() .As<ISnapshotStore<SchemaState>>()
.As<IExternalSystem>();
services.AddSingletonAs(c => new MongoContentRepository(mongoContentDatabase, c.GetService<IAppProvider>()))
.As<IContentRepository>()
.As<ISnapshotStore<ContentState>>()
.As<IEventConsumer>()
.As<IExternalSystem>(); .As<IExternalSystem>();
services.AddSingletonAs(c => new MongoHistoryEventRepository(mongoDatabase, c.GetServices<IHistoryEventsCreator>())) services.AddSingletonAs(c => new MongoHistoryEventRepository(mongoDatabase, c.GetServices<IHistoryEventsCreator>()))

49
tests/Squidex.Infrastructure.Tests/Commands/AggregateHandlerTests.cs

@ -7,6 +7,7 @@
// ========================================================================== // ==========================================================================
using System; using System;
using System.Collections.Generic;
using System.Threading.Tasks; using System.Threading.Tasks;
using FakeItEasy; using FakeItEasy;
using Squidex.Infrastructure.Commands.TestHelpers; using Squidex.Infrastructure.Commands.TestHelpers;
@ -32,32 +33,40 @@ namespace Squidex.Infrastructure.Commands
private readonly MyDomainObject domainObject = new MyDomainObject(); private readonly MyDomainObject domainObject = new MyDomainObject();
private readonly AggregateHandler sut; private readonly AggregateHandler sut;
public sealed class MyEvent : IEvent
{
}
public AggregateHandlerTests() public AggregateHandlerTests()
{ {
context = new CommandContext(new MyCommand { AggregateId = domainObjectId }); context = new CommandContext(new MyCommand { AggregateId = domainObjectId });
A.CallTo(() => store.WithEventSourcing<MyDomainObject>(domainObjectId.ToString(), A<Func<Envelope<IEvent>, Task>>.Ignored)) A.CallTo(() => store.WithSnapshots<MyDomainObject, object>(domainObjectId.ToString(), A<Func<object, Task>>.Ignored))
.Returns(persistence); .Returns(persistence);
A.CallTo(() => stateFactory.CreateAsync<MyDomainObject>(domainObjectId.ToString())) A.CallTo(() => stateFactory.CreateAsync<MyDomainObject>(domainObjectId.ToString()))
.Returns(Task.FromResult(domainObject)); .Returns(Task.FromResult(domainObject));
sut = new AggregateHandler(stateFactory, serviceProvider, log); sut = new AggregateHandler(stateFactory, serviceProvider, log);
domainObject.ActivateAsync(domainObjectId.ToString(), store).Wait();
} }
[Fact] [Fact]
public Task Create_async_should_throw_exception_if_not_aggregate_command() public Task Create_with_task_should_throw_exception_if_not_aggregate_command()
{ {
return Assert.ThrowsAnyAsync<ArgumentException>(() => sut.CreateAsync<MyDomainObject>(new CommandContext(A.Dummy<ICommand>()), x => TaskHelper.False)); return Assert.ThrowsAnyAsync<ArgumentException>(() => sut.CreateAsync<MyDomainObject>(new CommandContext(A.Dummy<ICommand>()), x => TaskHelper.False));
} }
[Fact] [Fact]
public async Task Create_async_should_create_domain_object_and_save() public async Task Create_with_task_should_create_domain_object_and_save()
{ {
MyDomainObject passedDomainObject = null; MyDomainObject passedDomainObject = null;
await sut.CreateAsync<MyDomainObject>(context, async x => await sut.CreateAsync<MyDomainObject>(context, async x =>
{ {
x.RaiseEvent(new MyEvent());
await Task.Yield(); await Task.Yield();
passedDomainObject = x; passedDomainObject = x;
@ -66,46 +75,44 @@ namespace Squidex.Infrastructure.Commands
Assert.Equal(domainObject, passedDomainObject); Assert.Equal(domainObject, passedDomainObject);
Assert.NotNull(context.Result<EntityCreatedResult<Guid>>()); Assert.NotNull(context.Result<EntityCreatedResult<Guid>>());
A.CallTo(() => persistence.ReadAsync(-1)) A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.Ignored))
.MustHaveHappened();
A.CallTo(() => persistence.WriteEventsAsync(A<Envelope<IEvent>[]>.Ignored))
.MustHaveHappened(); .MustHaveHappened();
} }
[Fact] [Fact]
public async Task Create_sync_should_create_domain_object_and_save() public async Task Create_should_create_domain_object_and_save()
{ {
MyDomainObject passedDomainObject = null; MyDomainObject passedDomainObject = null;
await sut.CreateAsync<MyDomainObject>(context, x => await sut.CreateAsync<MyDomainObject>(context, x =>
{ {
x.RaiseEvent(new MyEvent());
passedDomainObject = x; passedDomainObject = x;
}); });
Assert.Equal(domainObject, passedDomainObject); Assert.Equal(domainObject, passedDomainObject);
Assert.NotNull(context.Result<EntityCreatedResult<Guid>>()); Assert.NotNull(context.Result<EntityCreatedResult<Guid>>());
A.CallTo(() => persistence.ReadAsync(-1)) A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.Ignored))
.MustHaveHappened();
A.CallTo(() => persistence.WriteEventsAsync(A<Envelope<IEvent>[]>.Ignored))
.MustHaveHappened(); .MustHaveHappened();
} }
[Fact] [Fact]
public Task Update_async_should_throw_exception_if_not_aggregate_command() public Task Update_with_task_should_throw_exception_if_not_aggregate_command()
{ {
return Assert.ThrowsAnyAsync<ArgumentException>(() => sut.UpdateAsync<MyDomainObject>(new CommandContext(A.Dummy<ICommand>()), x => TaskHelper.False)); return Assert.ThrowsAnyAsync<ArgumentException>(() => sut.UpdateAsync<MyDomainObject>(new CommandContext(A.Dummy<ICommand>()), x => TaskHelper.False));
} }
[Fact] [Fact]
public async Task Update_async_should_create_domain_object_and_save() public async Task Update_with_task_should_create_domain_object_and_save()
{ {
MyDomainObject passedDomainObject = null; MyDomainObject passedDomainObject = null;
await sut.UpdateAsync<MyDomainObject>(context, async x => await sut.UpdateAsync<MyDomainObject>(context, async x =>
{ {
x.RaiseEvent(new MyEvent());
await Task.Yield(); await Task.Yield();
passedDomainObject = x; passedDomainObject = x;
@ -114,30 +121,26 @@ namespace Squidex.Infrastructure.Commands
Assert.Equal(domainObject, passedDomainObject); Assert.Equal(domainObject, passedDomainObject);
Assert.NotNull(context.Result<EntitySavedResult>()); Assert.NotNull(context.Result<EntitySavedResult>());
A.CallTo(() => persistence.ReadAsync(null)) A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.Ignored))
.MustHaveHappened();
A.CallTo(() => persistence.WriteEventsAsync(A<Envelope<IEvent>[]>.Ignored))
.MustHaveHappened(); .MustHaveHappened();
} }
[Fact] [Fact]
public async Task Update_sync_should_create_domain_object_and_save() public async Task Update_should_create_domain_object_and_save()
{ {
MyDomainObject passedDomainObject = null; MyDomainObject passedDomainObject = null;
await sut.UpdateAsync<MyDomainObject>(context, x => await sut.UpdateAsync<MyDomainObject>(context, x =>
{ {
x.RaiseEvent(new MyEvent());
passedDomainObject = x; passedDomainObject = x;
}); });
Assert.Equal(domainObject, passedDomainObject); Assert.Equal(domainObject, passedDomainObject);
Assert.NotNull(context.Result<EntitySavedResult>()); Assert.NotNull(context.Result<EntitySavedResult>());
A.CallTo(() => persistence.ReadAsync(null)) A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.Ignored))
.MustHaveHappened();
A.CallTo(() => persistence.WriteEventsAsync(A<Envelope<IEvent>[]>.Ignored))
.MustHaveHappened(); .MustHaveHappened();
} }
} }

93
tests/Squidex.Infrastructure.Tests/Commands/DomainObjectBaseTests.cs

@ -7,60 +7,115 @@
// ========================================================================== // ==========================================================================
using System; using System;
using System.Collections.Generic;
using System.Linq; using System.Linq;
using System.Threading.Tasks;
using FakeItEasy;
using Squidex.Infrastructure.Commands.TestHelpers; using Squidex.Infrastructure.Commands.TestHelpers;
using Squidex.Infrastructure.EventSourcing; using Squidex.Infrastructure.EventSourcing;
using Squidex.Infrastructure.Log;
using Squidex.Infrastructure.States;
using Xunit; using Xunit;
namespace Squidex.Infrastructure.Commands namespace Squidex.Infrastructure.Commands
{ {
public class DomainObjectBaseTests public class DomainObjectBaseTests
{ {
private readonly IStore store = A.Fake<IStore>();
private readonly IPersistence<object> persistence = A.Fake<IPersistence<object>>();
private readonly Guid id = Guid.NewGuid();
private readonly MyDomainObject sut = new MyDomainObject();
public DomainObjectBaseTests()
{
A.CallTo(() => store.WithSnapshots<MyDomainObject, object>(id.ToString(), A<Func<object, Task>>.Ignored))
.Returns(persistence);
}
[Fact] [Fact]
public void Should_instantiate() public void Should_instantiate()
{ {
var domainObjectId = Guid.NewGuid(); Assert.Equal(-1, sut.Version);
var domainObjectVersion = 123; }
var sut = new MyDomainObject(); [Fact]
public void Should_add_event_to_uncommitted_events_and_not_increase_version_when_raised()
{
var event1 = new MyEvent();
var event2 = new MyEvent();
sut.RaiseEvent(event1);
sut.RaiseEvent(event2);
Assert.Equal(-1, sut.Version);
Assert.Equal(new IEvent[] { event1, event2 }, sut.GetUncomittedEvents().Select(x => x.Payload).ToArray());
Assert.Equal(domainObjectId, sut.Id); sut.ClearUncommittedEvents();
Assert.Equal(domainObjectVersion, sut.Version);
Assert.Equal(0, sut.GetUncomittedEvents().Count);
} }
[Fact] [Fact]
public void Should_add_event_to_uncommitted_events_and_increase_version_when_raised() public async Task Should_write_state_and_events_when_saved()
{ {
A.CallTo(() => persistence.Version)
.Returns(100);
await sut.ActivateAsync(id.ToString(), store);
Assert.Equal(100, sut.Version);
var event1 = new MyEvent(); var event1 = new MyEvent();
var event2 = new MyEvent(); var event2 = new MyEvent();
var sut = new MyDomainObject(); sut.RaiseEvent(event1);
sut.RaiseEvent(event2);
sut.RaiseNewEvent(event1); var newState = "STATE";
sut.RaiseNewEvent(event2);
Assert.Equal(12, sut.Version); sut.UpdateState(newState);
Assert.Equal(new IEvent[] { event1, event2 }, sut.GetUncomittedEvents().Select(x => x.Payload).ToArray()); await sut.WriteAsync(A.Fake<ISemanticLog>());
sut.ClearUncommittedEvents(); A.CallTo(() => persistence.WriteSnapshotAsync(newState, 102))
.MustHaveHappened();
A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.That.Matches(x => x.Count() == 2)))
.MustHaveHappened();
Assert.Equal(0, sut.GetUncomittedEvents().Count); Assert.Empty(sut.GetUncomittedEvents());
} }
[Fact] [Fact]
public void Should_not_add_event_to_uncommitted_events_and_increase_version_when_raised() public async Task Should_ignore_exception_when_saving()
{ {
A.CallTo(() => persistence.Version)
.Returns(100);
A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.Ignored))
.Throws(new InvalidOperationException());
await sut.ActivateAsync(id.ToString(), store);
Assert.Equal(100, sut.Version);
var event1 = new MyEvent(); var event1 = new MyEvent();
var event2 = new MyEvent(); var event2 = new MyEvent();
var sut = new MyDomainObject(); sut.RaiseEvent(event1);
sut.RaiseEvent(event2);
sut.RaiseEvent(new Envelope<IEvent>(event1)); var newState = "STATE";
sut.RaiseEvent(new Envelope<IEvent>(event2));
Assert.Equal(12, sut.Version); sut.UpdateState(newState);
Assert.Equal(0, sut.GetUncomittedEvents().Count);
await sut.WriteAsync(A.Fake<ISemanticLog>());
A.CallTo(() => persistence.WriteSnapshotAsync(newState, 102))
.MustHaveHappened();
A.CallTo(() => persistence.WriteEventsAsync(A<IEnumerable<Envelope<IEvent>>>.That.Matches(x => x.Count() == 2)))
.MustHaveHappened();
Assert.Empty(sut.GetUncomittedEvents());
} }
} }
} }

14
tests/Squidex.Infrastructure.Tests/Commands/TestHelpers/MyDomainObject.cs

@ -6,25 +6,11 @@
// All rights reserved. // All rights reserved.
// ========================================================================== // ==========================================================================
using System;
using Squidex.Infrastructure.EventSourcing; using Squidex.Infrastructure.EventSourcing;
namespace Squidex.Infrastructure.Commands.TestHelpers namespace Squidex.Infrastructure.Commands.TestHelpers
{ {
internal sealed class MyDomainObject : DomainObjectBase<MyDomainObject, object> internal sealed class MyDomainObject : DomainObjectBase<MyDomainObject, object>
{ {
public MyDomainObject RaiseNewEvent(IEvent @event)
{
RaiseEvent(@event);
return this;
}
public MyDomainObject RaiseNewEvent(Envelope<IEvent> @event)
{
RaiseEvent(@event);
return this;
}
} }
} }

8
tests/Squidex.Infrastructure.Tests/DispatchingTests.cs

@ -208,7 +208,7 @@ namespace Squidex.Infrastructure
} }
[Fact] [Fact]
public async Task Should_invoke_correct_event_asynchronously() public async Task Should_invoke_correct_event_with_taskhronously()
{ {
var consumer = new MyAsyncConsumer(); var consumer = new MyAsyncConsumer();
@ -222,7 +222,7 @@ namespace Squidex.Infrastructure
} }
[Fact] [Fact]
public async Task Should_invoke_correct_event_with_context_asynchronously() public async Task Should_invoke_correct_event_with_context_with_taskhronously()
{ {
var consumer = new MyAsyncConsumer(); var consumer = new MyAsyncConsumer();
@ -264,7 +264,7 @@ namespace Squidex.Infrastructure
} }
[Fact] [Fact]
public async Task Should_invoke_correct_event_and_return_synchronously() public async Task Should_invoke_correct_event_and_returnhronously()
{ {
var consumer = new MyAsyncFuncConsumer(); var consumer = new MyAsyncFuncConsumer();
@ -278,7 +278,7 @@ namespace Squidex.Infrastructure
} }
[Fact] [Fact]
public async Task Should_invoke_correct_event_with_context_and_return_synchronously() public async Task Should_invoke_correct_event_with_context_and_returnhronously()
{ {
var consumer = new MyAsyncFuncConsumer(); var consumer = new MyAsyncFuncConsumer();

22
tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs

@ -70,8 +70,8 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
A.CallTo(() => persistence.ReadAsync(null)) A.CallTo(() => persistence.ReadAsync(null))
.Invokes(new Action<long?>(s => apply(state))); .Invokes(new Action<long?>(s => apply(state)));
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.Invokes(new Action<EventConsumerState>(s => state = s)); .Invokes(new Action<EventConsumerState, long>((s, v) => state = s));
A.CallTo(() => formatter.Parse(eventData, true)).Returns(envelope); A.CallTo(() => formatter.Parse(eventData, true)).Returns(envelope);
@ -132,7 +132,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = true, Position = initialPosition, Error = null }); state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = true, Position = initialPosition, Error = null });
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Once); .MustHaveHappened(Repeated.Exactly.Once);
A.CallTo(() => eventSubscription.StopAsync()) A.CallTo(() => eventSubscription.StopAsync())
@ -150,7 +150,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = null, Error = null }); state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = null, Error = null });
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Twice); .MustHaveHappened(Repeated.Exactly.Twice);
A.CallTo(() => eventConsumer.ClearAsync()) A.CallTo(() => eventConsumer.ClearAsync())
@ -180,7 +180,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = @event.EventPosition, Error = null }); state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = @event.EventPosition, Error = null });
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Once); .MustHaveHappened(Repeated.Exactly.Once);
A.CallTo(() => eventConsumer.On(envelope)) A.CallTo(() => eventConsumer.On(envelope))
@ -204,7 +204,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = @event.EventPosition, Error = null }); state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = @event.EventPosition, Error = null });
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Once); .MustHaveHappened(Repeated.Exactly.Once);
A.CallTo(() => eventConsumer.On(envelope)) A.CallTo(() => eventConsumer.On(envelope))
@ -243,7 +243,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = initialPosition, Error = null }); state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = false, Position = initialPosition, Error = null });
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustNotHaveHappened(); .MustNotHaveHappened();
} }
@ -263,7 +263,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = true, Position = initialPosition, Error = ex.ToString() }); state.ShouldBeEquivalentTo(new EventConsumerState { IsStopped = true, Position = initialPosition, Error = ex.ToString() });
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Once); .MustHaveHappened(Repeated.Exactly.Once);
A.CallTo(() => eventSubscription.StopAsync()) A.CallTo(() => eventSubscription.StopAsync())
@ -292,7 +292,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
A.CallTo(() => eventConsumer.On(envelope)) A.CallTo(() => eventConsumer.On(envelope))
.MustHaveHappened(); .MustHaveHappened();
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Once); .MustHaveHappened(Repeated.Exactly.Once);
A.CallTo(() => eventSubscription.StopAsync()) A.CallTo(() => eventSubscription.StopAsync())
@ -323,7 +323,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
A.CallTo(() => eventConsumer.On(envelope)) A.CallTo(() => eventConsumer.On(envelope))
.MustNotHaveHappened(); .MustNotHaveHappened();
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Once); .MustHaveHappened(Repeated.Exactly.Once);
A.CallTo(() => eventSubscription.StopAsync()) A.CallTo(() => eventSubscription.StopAsync())
@ -354,7 +354,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains
A.CallTo(() => eventConsumer.On(envelope)) A.CallTo(() => eventConsumer.On(envelope))
.MustHaveHappened(); .MustHaveHappened();
A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored)) A.CallTo(() => persistence.WriteSnapshotAsync(A<EventConsumerState>.Ignored, -1))
.MustHaveHappened(Repeated.Exactly.Twice); .MustHaveHappened(Repeated.Exactly.Twice);
A.CallTo(() => eventSubscription.StopAsync()) A.CallTo(() => eventSubscription.StopAsync())

10
tests/Squidex.Infrastructure.Tests/States/StateEventSourcingTests.cs

@ -52,13 +52,13 @@ namespace Squidex.Infrastructure.States
private class MyStatefulObjectWithSnapshot : IStatefulObject private class MyStatefulObjectWithSnapshot : IStatefulObject
{ {
private IPersistence<int> persistence; private IPersistence<object> persistence;
public long? ExpectedVersion { get; set; } public long? ExpectedVersion { get; set; }
public Task ActivateAsync(string key, IStore store) public Task ActivateAsync(string key, IStore store)
{ {
persistence = store.WithSnapshotsAndEventSourcing<MyStatefulObject, int>(key, s => TaskHelper.Done, s => TaskHelper.Done); persistence = store.WithSnapshotsAndEventSourcing<MyStatefulObject, object>(key, s => TaskHelper.Done, s => TaskHelper.Done);
return persistence.ReadAsync(ExpectedVersion); return persistence.ReadAsync(ExpectedVersion);
} }
@ -72,7 +72,7 @@ namespace Squidex.Infrastructure.States
private readonly IMemoryCache cache = new MemoryCache(Options.Create(new MemoryCacheOptions())); private readonly IMemoryCache cache = new MemoryCache(Options.Create(new MemoryCacheOptions()));
private readonly IPubSub pubSub = new InMemoryPubSub(true); private readonly IPubSub pubSub = new InMemoryPubSub(true);
private readonly IServiceProvider services = A.Fake<IServiceProvider>(); private readonly IServiceProvider services = A.Fake<IServiceProvider>();
private readonly ISnapshotStore<int> snapshotStore = A.Fake<ISnapshotStore<int>>(); private readonly ISnapshotStore<object> snapshotStore = A.Fake<ISnapshotStore<object>>();
private readonly IStreamNameResolver streamNameResolver = A.Fake<IStreamNameResolver>(); private readonly IStreamNameResolver streamNameResolver = A.Fake<IStreamNameResolver>();
private readonly StateFactory sut; private readonly StateFactory sut;
@ -82,7 +82,7 @@ namespace Squidex.Infrastructure.States
.Returns(statefulObject); .Returns(statefulObject);
A.CallTo(() => services.GetService(typeof(MyStatefulObjectWithSnapshot))) A.CallTo(() => services.GetService(typeof(MyStatefulObjectWithSnapshot)))
.Returns(statefulObjectWithSnapShot); .Returns(statefulObjectWithSnapShot);
A.CallTo(() => services.GetService(typeof(ISnapshotStore<int>))) A.CallTo(() => services.GetService(typeof(ISnapshotStore<object>)))
.Returns(snapshotStore); .Returns(snapshotStore);
A.CallTo(() => streamNameResolver.GetStreamName(typeof(MyStatefulObject), key)) A.CallTo(() => streamNameResolver.GetStreamName(typeof(MyStatefulObject), key))
@ -278,7 +278,7 @@ namespace Squidex.Infrastructure.States
statefulObject.ExpectedVersion = null; statefulObject.ExpectedVersion = null;
A.CallTo(() => snapshotStore.ReadAsync(key)) A.CallTo(() => snapshotStore.ReadAsync(key))
.ReturnsLazily(() => Task.Delay(1).ContinueWith(x => (1, 1L))); .ReturnsLazily(() => Task.Delay(1).ContinueWith(x => ((object)1, 1L)));
var tasks = new List<Task<MyStatefulObject>>(); var tasks = new List<Task<MyStatefulObject>>();

4
tests/Squidex.Infrastructure.Tests/Tasks/SingleThreadedDispatcherTests.cs

@ -18,7 +18,7 @@ namespace Squidex.Infrastructure.Tasks
private readonly SingleThreadedDispatcher sut = new SingleThreadedDispatcher(); private readonly SingleThreadedDispatcher sut = new SingleThreadedDispatcher();
[Fact] [Fact]
public async Task Should_handle_async_messages_sequentially() public async Task Should_handle_with_task_messages_sequentially()
{ {
var source = Enumerable.Range(1, 100); var source = Enumerable.Range(1, 100);
var target = new List<int>(); var target = new List<int>();
@ -39,7 +39,7 @@ namespace Squidex.Infrastructure.Tasks
} }
[Fact] [Fact]
public async Task Should_handle_sync_messages_sequentially() public async Task Should_handle_messages_sequentially()
{ {
var source = Enumerable.Range(1, 100); var source = Enumerable.Range(1, 100);
var target = new List<int>(); var target = new List<int>();

Loading…
Cancel
Save