-
Notifications
You must be signed in to change notification settings - Fork 1.7k
/
SqlStreamAggregateStore.cs
167 lines (137 loc) · 5.38 KB
/
SqlStreamAggregateStore.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using CompanyName.MyMeetings.BuildingBlocks.Application.Data;
using CompanyName.MyMeetings.BuildingBlocks.Domain;
using CompanyName.MyMeetings.BuildingBlocks.Infrastructure.Serialization;
using CompanyName.MyMeetings.Modules.Payments.Domain.SeedWork;
using CompanyName.MyMeetings.Modules.Payments.Infrastructure.Configuration;
using Newtonsoft.Json;
using SqlStreamStore;
using SqlStreamStore.Streams;
namespace CompanyName.MyMeetings.Modules.Payments.Infrastructure.AggregateStore
{
public class SqlStreamAggregateStore : IAggregateStore
{
private readonly IStreamStore _streamStore;
private readonly List<IDomainEvent> _appendedChanges;
private readonly List<AggregateToSave> _aggregatesToSave;
public SqlStreamAggregateStore(
ISqlConnectionFactory sqlConnectionFactory)
{
_appendedChanges = new List<IDomainEvent>();
_streamStore =
new MsSqlStreamStore(
new MsSqlStreamStoreSettings(sqlConnectionFactory.GetConnectionString())
{
Schema = DatabaseSchema.Name
});
_aggregatesToSave = new List<AggregateToSave>();
}
public async Task Save()
{
foreach (var aggregateToSave in _aggregatesToSave)
{
await _streamStore.AppendToStream(
GetStreamId(aggregateToSave.Aggregate),
aggregateToSave.Aggregate.Version,
aggregateToSave.Messages.ToArray());
}
_aggregatesToSave.Clear();
}
public async Task<T> Load<T>(AggregateId<T> aggregateId)
where T : AggregateRoot
{
var streamId = GetStreamId(aggregateId);
IList<IDomainEvent> domainEvents = new List<IDomainEvent>();
ReadStreamPage readStreamPage;
int position = StreamVersion.Start;
int take = 100;
do
{
readStreamPage = await _streamStore.ReadStreamForwards(streamId, position, take);
var messages = readStreamPage.Messages;
foreach (var streamMessage in messages)
{
Type type = DomainEventTypeMappings.Dictionary[streamMessage.Type];
var jsonData = await streamMessage.GetJsonData();
var domainEvent = JsonConvert.DeserializeObject(jsonData, type) as IDomainEvent;
domainEvents.Add(domainEvent);
}
position += take;
}
while (!readStreamPage.IsEnd);
if (!domainEvents.Any())
{
return null;
}
var aggregate = (T)Activator.CreateInstance(typeof(T), true);
aggregate.Load(domainEvents);
return aggregate;
}
public List<IDomainEvent> GetChanges()
{
return _appendedChanges;
}
public void AppendChanges<T>(T aggregate)
where T : AggregateRoot
{
_aggregatesToSave.Add(new AggregateToSave(aggregate, CreateStreamMessages(aggregate).ToList()));
}
public void ClearChanges()
{
_appendedChanges.Clear();
}
private class AggregateToSave
{
public AggregateToSave(AggregateRoot aggregate, List<NewStreamMessage> messages)
{
Aggregate = aggregate;
Messages = messages;
}
public AggregateRoot Aggregate { get; }
public List<NewStreamMessage> Messages { get; }
}
private NewStreamMessage[] CreateStreamMessages<T>(
T aggregate)
where T : AggregateRoot
{
List<NewStreamMessage> newStreamMessages = new List<NewStreamMessage>();
var domainEvents = aggregate.GetDomainEvents();
foreach (var domainEvent in domainEvents)
{
string jsonData = JsonConvert.SerializeObject(domainEvent, new JsonSerializerSettings
{
ContractResolver = new AllPropertiesContractResolver()
});
var message = new NewStreamMessage(
domainEvent.Id,
MapDomainEventToType(domainEvent),
jsonData);
newStreamMessages.Add(message);
_appendedChanges.Add(domainEvent);
}
return newStreamMessages.ToArray();
}
private string MapDomainEventToType(IDomainEvent domainEvent)
{
foreach (var key in DomainEventTypeMappings.Dictionary.Keys)
{
if (DomainEventTypeMappings.Dictionary[key] == domainEvent.GetType())
{
return key;
}
}
throw new ArgumentException("Invalid Domain Event type", nameof(domainEvent));
}
private static string GetStreamId<T>(T aggregate)
where T : AggregateRoot
{
return $"{aggregate.GetType().Name}-{aggregate.Id:N}";
}
private static string GetStreamId<T>(AggregateId<T> aggregateId)
where T : AggregateRoot
=> $"{typeof(T).Name}-{aggregateId.Value:N}";
}
}