Skip to content

Commit c9ae01e

Browse files
authored
Merge branch 'develop' into hangfire-update
2 parents c7c8280 + 8173d09 commit c9ae01e

16 files changed

Lines changed: 180 additions & 75 deletions

File tree

RELEASE_NOTES.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,8 @@
77
`Hangfire`
88
* New: Added an overload to `IDomainEventPublisher.PublishAsync` that isn't
99
generic and doesn't require an aggregate ID
10+
* New: Added `IReadModelPopulator.DeleteAsync` that allows deletion of single
11+
read models
1012
* Obsolete: `IDomainEventPublisher.PublishAsync<,>` (generic) in favor of the
1113
new less restrictive non-generic overload
1214

Source/EventFlow.Elasticsearch.Tests/IntegrationTests/ElasticsearchReadModelStoreTests.cs

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -23,8 +23,6 @@
2323

2424
using System;
2525
using System.Net;
26-
using System.Threading;
27-
using System.Threading.Tasks;
2826
using EventFlow.Configuration;
2927
using EventFlow.Elasticsearch.Extensions;
3028
using EventFlow.Elasticsearch.ReadStores;
@@ -45,6 +43,8 @@ namespace EventFlow.Elasticsearch.Tests.IntegrationTests
4543
[Category(Categories.Integration)]
4644
public class ElasticsearchReadModelStoreTests : TestSuiteForReadModelStore
4745
{
46+
protected override Type ReadModelType { get; } = typeof(ElasticsearchThingyReadModel);
47+
4848
private IElasticClient _elasticClient;
4949
private ElasticsearchRunner.ElasticsearchInstance _elasticsearchInstance;
5050
private string _indexName;
@@ -122,16 +122,6 @@ protected override IRootResolver CreateRootResolver(IEventFlowOptions eventFlowO
122122
}
123123
}
124124

125-
protected override Task PurgeTestAggregateReadModelAsync()
126-
{
127-
return ReadModelPopulator.PurgeAsync<ElasticsearchThingyReadModel>(CancellationToken.None);
128-
}
129-
130-
protected override Task PopulateTestAggregateReadModelAsync()
131-
{
132-
return ReadModelPopulator.PopulateAsync<ElasticsearchThingyReadModel>(CancellationToken.None);
133-
}
134-
135125
[TearDown]
136126
public void TearDown()
137127
{

Source/EventFlow.Elasticsearch/ReadStores/ElasticsearchReadModelStore.cs

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -79,6 +79,22 @@ public async Task<ReadModelEnvelope<TReadModel>> GetAsync(
7979
return ReadModelEnvelope<TReadModel>.With(id, getResponse.Source, getResponse.Version);
8080
}
8181

82+
public async Task DeleteAsync(
83+
string id,
84+
CancellationToken cancellationToken)
85+
{
86+
var readModelDescription = _readModelDescriptionProvider.GetReadModelDescription<TReadModel>();
87+
88+
await _elasticClient.DeleteAsync(
89+
new DocumentPath<TReadModel>(id),
90+
d => d
91+
.Index(readModelDescription.IndexName.Value)
92+
.RequestConfiguration(c => c
93+
.AllowedStatusCodes((int) HttpStatusCode.NotFound)),
94+
cancellationToken)
95+
.ConfigureAwait(false);
96+
}
97+
8298
public async Task DeleteAllAsync(
8399
CancellationToken cancellationToken)
84100
{

Source/EventFlow.MsSql.Tests/IntegrationTests/ReadStores/MsSqlReadModelStoreTests.cs

Lines changed: 3 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -21,8 +21,7 @@
2121
// IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN
2222
// CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
2323

24-
using System.Threading;
25-
using System.Threading.Tasks;
24+
using System;
2625
using EventFlow.Configuration;
2726
using EventFlow.Extensions;
2827
using EventFlow.MsSql.EventStores;
@@ -40,6 +39,8 @@ namespace EventFlow.MsSql.Tests.IntegrationTests.ReadStores
4039
[Category(Categories.Integration)]
4140
public class MsSqlReadModelStoreTests : TestSuiteForReadModelStore
4241
{
42+
protected override Type ReadModelType { get; } = typeof(MsSqlThingyReadModel);
43+
4344
private IMsSqlDatabase _testDatabase;
4445

4546
protected override IRootResolver CreateRootResolver(IEventFlowOptions eventFlowOptions)
@@ -64,16 +65,6 @@ protected override IRootResolver CreateRootResolver(IEventFlowOptions eventFlowO
6465
return resolver;
6566
}
6667

67-
protected override Task PurgeTestAggregateReadModelAsync()
68-
{
69-
return ReadModelPopulator.PurgeAsync<MsSqlThingyReadModel>(CancellationToken.None);
70-
}
71-
72-
protected override Task PopulateTestAggregateReadModelAsync()
73-
{
74-
return ReadModelPopulator.PopulateAsync<MsSqlThingyReadModel>(CancellationToken.None);
75-
}
76-
7768
[TearDown]
7869
public void TearDown()
7970
{

Source/EventFlow.MsSql/ReadStores/MssqlReadModelStore.cs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,26 @@ public override async Task<ReadModelEnvelope<TReadModel>> GetAsync(string id, Ca
172172
: ReadModelEnvelope<TReadModel>.With(id, readModel);
173173
}
174174

175+
public override async Task DeleteAsync(
176+
string id,
177+
CancellationToken cancellationToken)
178+
{
179+
var sql = _readModelSqlGenerator.CreateDeleteSql<TReadModel>();
180+
var readModelName = typeof(TReadModel).Name;
181+
182+
var rowsAffected = await _connection.ExecuteAsync(
183+
Label.Named("mssql-delete-read-model", readModelName),
184+
cancellationToken,
185+
sql,
186+
new { EventFlowReadModelId = id })
187+
.ConfigureAwait(false);
188+
189+
if (rowsAffected != 0)
190+
{
191+
Log.Verbose($"Deleted read model '{id}' of type '{readModelName}'");
192+
}
193+
}
194+
175195
public override async Task DeleteAllAsync(CancellationToken cancellationToken)
176196
{
177197
var sql = _readModelSqlGenerator.CreatePurgeSql<TReadModel>();

Source/EventFlow.SQLite.Tests/IntegrationTests/ReadStores/SQLiteReadStoreTests.cs

Lines changed: 2 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@
2424
using System;
2525
using System.IO;
2626
using System.Threading;
27-
using System.Threading.Tasks;
2827
using EventFlow.Configuration;
2928
using EventFlow.Core;
3029
using EventFlow.Extensions;
@@ -42,6 +41,8 @@ namespace EventFlow.SQLite.Tests.IntegrationTests.ReadStores
4241
[Category(Categories.Integration)]
4342
public class SQLiteReadStoreTests : TestSuiteForReadModelStore
4443
{
44+
protected override Type ReadModelType { get; } = typeof(SQLiteThingyReadModel);
45+
4546
private string _databasePath;
4647

4748
protected override IRootResolver CreateRootResolver(IEventFlowOptions eventFlowOptions)
@@ -84,16 +85,6 @@ [Message] [nvarchar](512) NOT NULL
8485
return resolver;
8586
}
8687

87-
protected override Task PurgeTestAggregateReadModelAsync()
88-
{
89-
return ReadModelPopulator.PurgeAsync<SQLiteThingyReadModel>(CancellationToken.None);
90-
}
91-
92-
protected override Task PopulateTestAggregateReadModelAsync()
93-
{
94-
return ReadModelPopulator.PopulateAsync<SQLiteThingyReadModel>(CancellationToken.None);
95-
}
96-
9788
[TearDown]
9889
public void TearDown()
9990
{

Source/EventFlow.Sql/ReadModels/IReadModelSqlGenerator.cs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ string CreateInsertSql<TReadModel>()
3333
string CreateSelectSql<TReadModel>()
3434
where TReadModel : IReadModel;
3535

36+
string CreateDeleteSql<TReadModel>()
37+
where TReadModel : IReadModel;
38+
3639
string CreateUpdateSql<TReadModel>()
3740
where TReadModel : IReadModel;
3841

Source/EventFlow.Sql/ReadModels/ReadModelSqlGenerator.cs

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ public class ReadModelSqlGenerator : IReadModelSqlGenerator
4040
private static readonly ConcurrentDictionary<Type, string> IdentityColumns = new ConcurrentDictionary<Type, string>();
4141
private readonly Dictionary<Type, string> _insertSqls = new Dictionary<Type, string>();
4242
private readonly Dictionary<Type, string> _purgeSqls = new Dictionary<Type, string>();
43+
private readonly Dictionary<Type, string> _deleteSqls = new Dictionary<Type, string>();
4344
private readonly Dictionary<Type, string> _selectSqls = new Dictionary<Type, string>();
4445
private readonly Dictionary<Type, string> _updateSqls = new Dictionary<Type, string>();
4546

@@ -82,6 +83,21 @@ public string CreateSelectSql<TReadModel>()
8283
return sql;
8384
}
8485

86+
public string CreateDeleteSql<TReadModel>()
87+
where TReadModel : IReadModel
88+
{
89+
var readModelType = typeof(TReadModel);
90+
if (_deleteSqls.TryGetValue(readModelType, out var sql))
91+
{
92+
return sql;
93+
}
94+
95+
sql = $"DELETE FROM {GetTableName<TReadModel>()} WHERE {GetIdentityColumn<TReadModel>()} = @EventFlowReadModelId";
96+
_deleteSqls[readModelType] = sql;
97+
98+
return sql;
99+
}
100+
85101
public string CreateUpdateSql<TReadModel>()
86102
where TReadModel : IReadModel
87103
{

Source/EventFlow.Sql/ReadModels/SqlReadModelStore.cs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,26 @@ public override async Task<ReadModelEnvelope<TReadModel>> GetAsync(string id, Ca
165165
: ReadModelEnvelope<TReadModel>.With(id, readModel);
166166
}
167167

168+
public override async Task DeleteAsync(
169+
string id,
170+
CancellationToken cancellationToken)
171+
{
172+
var sql = _readModelSqlGenerator.CreateDeleteSql<TReadModel>();
173+
var readModelName = typeof(TReadModel).Name;
174+
175+
var rowsAffected = await _connection.ExecuteAsync(
176+
Label.Named("mssql-delete-read-model", readModelName),
177+
cancellationToken,
178+
sql,
179+
new { EventFlowReadModelId = id })
180+
.ConfigureAwait(false);
181+
182+
if (rowsAffected != 0)
183+
{
184+
Log.Verbose($"Deleted read model '{id}' of type '{readModelName}'");
185+
}
186+
}
187+
168188
public override async Task DeleteAllAsync(CancellationToken cancellationToken)
169189
{
170190
var sql = _readModelSqlGenerator.CreatePurgeSql<TReadModel>();

Source/EventFlow.TestHelpers/Suites/TestSuiteForReadModelStore.cs

Lines changed: 35 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -21,8 +21,10 @@
2121
// IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN
2222
// CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
2323

24+
using System;
2425
using System.Collections.Generic;
2526
using System.Linq;
27+
using System.Threading;
2628
using System.Threading.Tasks;
2729
using EventFlow.TestHelpers.Aggregates;
2830
using EventFlow.TestHelpers.Aggregates.Commands;
@@ -141,13 +143,40 @@ public async Task PurgeRemovesReadModels()
141143
await PublishPingCommandAsync(id).ConfigureAwait(false);
142144

143145
// Act
144-
await PurgeTestAggregateReadModelAsync().ConfigureAwait(false);
146+
await ReadModelPopulator.PurgeAsync(ReadModelType, CancellationToken.None).ConfigureAwait(false);
145147
var readModel = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id)).ConfigureAwait(false);
146148

147149
// Assert
148150
readModel.Should().BeNull();
149151
}
150152

153+
[Test]
154+
public async Task DeleteRemovesSpecificReadModel()
155+
{
156+
// Arrange
157+
var id1 = ThingyId.New;
158+
var id2 = ThingyId.New;
159+
await PublishPingCommandAsync(id1).ConfigureAwait(false);
160+
await PublishPingCommandAsync(id2).ConfigureAwait(false);
161+
var readModel1 = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id1)).ConfigureAwait(false);
162+
var readModel2 = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id2)).ConfigureAwait(false);
163+
readModel1.Should().NotBeNull();
164+
readModel2.Should().NotBeNull();
165+
166+
// Act
167+
await ReadModelPopulator.DeleteAsync(
168+
id1.Value,
169+
ReadModelType,
170+
CancellationToken.None)
171+
.ConfigureAwait(false);
172+
173+
// Assert
174+
readModel1 = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id1)).ConfigureAwait(false);
175+
readModel2 = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id2)).ConfigureAwait(false);
176+
readModel1.Should().BeNull();
177+
readModel2.Should().NotBeNull();
178+
}
179+
151180
[Test]
152181
public async Task RePopulateHandlesManyAggregates()
153182
{
@@ -158,8 +187,8 @@ public async Task RePopulateHandlesManyAggregates()
158187
await PublishPingCommandsAsync(id2, 5).ConfigureAwait(false);
159188

160189
// Act
161-
await PurgeTestAggregateReadModelAsync().ConfigureAwait(false);
162-
await PopulateTestAggregateReadModelAsync().ConfigureAwait(false);
190+
await ReadModelPopulator.PurgeAsync(ReadModelType, CancellationToken.None).ConfigureAwait(false);
191+
await ReadModelPopulator.PopulateAsync(ReadModelType, CancellationToken.None).ConfigureAwait(false);
163192

164193
// Assert
165194
var readModel1 = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id1)).ConfigureAwait(false);
@@ -175,10 +204,10 @@ public async Task PopulateCreatesReadModels()
175204
// Arrange
176205
var id = ThingyId.New;
177206
await PublishPingCommandsAsync(id, 2).ConfigureAwait(false);
178-
await PurgeTestAggregateReadModelAsync().ConfigureAwait(false);
207+
await ReadModelPopulator.PurgeAsync(ReadModelType, CancellationToken.None).ConfigureAwait(false);
179208

180209
// Act
181-
await PopulateTestAggregateReadModelAsync().ConfigureAwait(false);
210+
await ReadModelPopulator.PopulateAsync(ReadModelType, CancellationToken.None).ConfigureAwait(false);
182211
var readModel = await QueryProcessor.ProcessAsync(new ThingyGetQuery(id)).ConfigureAwait(false);
183212

184213
// Assert
@@ -193,8 +222,6 @@ private async Task<IReadOnlyCollection<ThingyMessage>> CreateAndPublishThingyMes
193222
return thingyMessages;
194223
}
195224

196-
protected abstract Task PurgeTestAggregateReadModelAsync();
197-
198-
protected abstract Task PopulateTestAggregateReadModelAsync();
225+
protected abstract Type ReadModelType { get; }
199226
}
200227
}

0 commit comments

Comments
 (0)