-
Notifications
You must be signed in to change notification settings - Fork 78
Expand file tree
/
Copy pathSagaEventStoreRepository.cs
More file actions
136 lines (115 loc) · 3.62 KB
/
SagaEventStoreRepository.cs
File metadata and controls
136 lines (115 loc) · 3.62 KB
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
namespace CommonDomain.Persistence.EventStore
{
using System;
using System.Collections.Generic;
using System.Linq;
using global::EventStore;
using global::EventStore.Persistence;
public class SagaEventStoreRepository : ISagaRepository, IDisposable
{
private const string SagaTypeHeader = "SagaType";
private const string UndispatchedMessageHeader = "UndispatchedMessage.";
private readonly IDictionary<Guid, IEventStream> streams = new Dictionary<Guid, IEventStream>();
private readonly IStoreEvents eventStore;
private readonly IConstructSagas sagaFactory;
public SagaEventStoreRepository(IStoreEvents eventStore, IConstructSagas sagaFactory)
{
this.eventStore = eventStore;
this.sagaFactory = sagaFactory;
}
public void Dispose()
{
this.Dispose(true);
GC.SuppressFinalize(this);
}
protected virtual void Dispose(bool disposing)
{
if (!disposing)
return;
lock (this.streams)
{
foreach (var stream in this.streams)
stream.Value.Dispose();
this.streams.Clear();
}
}
public TSaga GetById<TSaga>(Guid sagaId) where TSaga : class, ISaga
{
return BuildSaga<TSaga>(sagaId, this.OpenStream(sagaId));
}
private IEventStream OpenStream(Guid sagaId)
{
IEventStream stream;
if (this.streams.TryGetValue(sagaId, out stream))
return stream;
try
{
stream = this.eventStore.OpenStream(sagaId, 0, int.MaxValue);
}
catch (StreamNotFoundException)
{
stream = this.eventStore.CreateStream(sagaId);
}
return this.streams[sagaId] = stream;
}
private TSaga BuildSaga<TSaga>(Guid sagaId, IEventStream stream) where TSaga : class, ISaga
{
var saga = sagaFactory.Build<TSaga>(sagaId);
foreach (var @event in stream.CommittedEvents.Select(x => x.Body))
saga.Transition(@event);
saga.ClearUncommittedEvents();
saga.ClearUndispatchedMessages();
return saga;
}
public void Save(ISaga saga, Guid commitId, Action<IDictionary<string, object>> updateHeaders)
{
if (saga == null)
throw new ArgumentNullException("saga", ExceptionMessages.NullArgument);
var headers = PrepareHeaders(saga, updateHeaders);
var stream = this.PrepareStream(saga, headers);
Persist(stream, commitId);
saga.ClearUncommittedEvents();
saga.ClearUndispatchedMessages();
}
private static Dictionary<string, object> PrepareHeaders(ISaga saga, Action<IDictionary<string, object>> updateHeaders)
{
var headers = new Dictionary<string, object>();
headers[SagaTypeHeader] = saga.GetType().FullName;
if (updateHeaders != null)
updateHeaders(headers);
var i = 0;
foreach (var command in saga.GetUndispatchedMessages())
headers[UndispatchedMessageHeader + i++] = command;
return headers;
}
private IEventStream PrepareStream(ISaga saga, Dictionary<string, object> headers)
{
IEventStream stream;
if (!this.streams.TryGetValue(saga.Id, out stream))
this.streams[saga.Id] = stream = this.eventStore.CreateStream(saga.Id);
foreach (var item in headers)
stream.UncommittedHeaders[item.Key] = item.Value;
saga.GetUncommittedEvents()
.Cast<object>()
.Select(x => new EventMessage { Body = x })
.ToList()
.ForEach(stream.Add);
return stream;
}
private static void Persist(IEventStream stream, Guid commitId)
{
try
{
stream.CommitChanges(commitId);
}
catch (DuplicateCommitException)
{
stream.ClearChanges();
}
catch (StorageException e)
{
throw new PersistenceException(e.Message, e);
}
}
}
}