-
Notifications
You must be signed in to change notification settings - Fork 2.1k
Expand file tree
/
Copy pathAdoNetClusteringTable.cs
More file actions
346 lines (306 loc) · 13 KB
/
Copy pathAdoNetClusteringTable.cs
File metadata and controls
346 lines (306 loc) · 13 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
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
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
using System;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Orleans.Clustering.AdoNet.Storage;
using Orleans.Configuration;
#nullable disable
namespace Orleans.Runtime.MembershipService
{
public partial class AdoNetClusteringTable : IMembershipTable
{
private readonly string clusterId;
private readonly IServiceProvider serviceProvider;
private readonly ILogger logger;
private RelationalOrleansQueries orleansQueries;
private readonly AdoNetClusteringSiloOptions clusteringTableOptions;
public AdoNetClusteringTable(
IServiceProvider serviceProvider,
IOptions<ClusterOptions> clusterOptions,
IOptions<AdoNetClusteringSiloOptions> clusteringOptions,
ILogger<AdoNetClusteringTable> logger)
{
this.serviceProvider = serviceProvider;
this.logger = logger;
this.clusteringTableOptions = clusteringOptions.Value;
this.clusterId = clusterOptions.Value.ClusterId;
}
public async Task InitializeMembershipTable(bool tryInitTableVersion)
{
LogTraceInitializeMembershipTable();
//This initializes all of Orleans operational queries from the database using a well known view
//and assumes the database with appropriate definitions exists already.
orleansQueries = await RelationalOrleansQueries.CreateInstance(
clusteringTableOptions.Invariant,
clusteringTableOptions.ConnectionString);
// even if I am not the one who created the table,
// try to insert an initial table version if it is not already there,
// so we always have a first table version row, before this silo starts working.
if (tryInitTableVersion)
{
var wasCreated = await InitTableAsync();
if (wasCreated)
{
LogInfoCreatedNewTableVersionRow();
}
}
}
public async Task<MembershipTableData> ReadRow(SiloAddress key)
{
LogTraceReadRow(key);
try
{
return await orleansQueries.MembershipReadRowAsync(this.clusterId, key);
}
catch (Exception ex)
{
LogDebugReadRowFailed(ex);
throw;
}
}
public async Task<MembershipTableData> ReadAll()
{
LogTraceReadAll();
try
{
return await orleansQueries.MembershipReadAllAsync(this.clusterId);
}
catch (Exception ex)
{
LogDebugReadAllFailed(ex);
throw;
}
}
public async Task<bool> InsertRow(MembershipEntry entry, TableVersion tableVersion)
{
LogTraceInsertRow(entry, tableVersion);
//The "tableVersion" parameter should always exist when inserting a row as Init should
//have been called and membership version created and read. This is an optimization to
//not to go through all the way to database to fail a conditional check on etag (which does
//exist for the sake of robustness) as mandated by Orleans membership protocol.
//Likewise, no update can be done without membership entry.
if (entry == null)
{
LogDebugInsertRowAbortedNullEntry();
throw new ArgumentNullException(nameof(entry));
}
if (tableVersion is null)
{
LogDebugInsertRowAbortedNullTableVersion();
throw new ArgumentNullException(nameof(tableVersion));
}
try
{
ClearMetadataIfUnavailable(entry);
return await orleansQueries.InsertMembershipRowAsync(this.clusterId, entry, tableVersion.VersionEtag);
}
catch (Exception ex)
{
LogDebugInsertRowFailed(ex);
throw;
}
}
public async Task<bool> UpdateRow(MembershipEntry entry, string etag, TableVersion tableVersion)
{
LogTraceUpdateRow(entry, etag, tableVersion);
//The "tableVersion" parameter should always exist when updating a row as Init should
//have been called and membership version created and read. This is an optimization to
//not to go through all the way to database to fail a conditional check (which does
//exist for the sake of robustness) as mandated by Orleans membership protocol.
//Likewise, no update can be done without membership entry or an etag.
if (entry == null)
{
LogDebugUpdateRowAbortedNullEntry();
throw new ArgumentNullException(nameof(entry));
}
if (tableVersion is null)
{
LogDebugUpdateRowAbortedNullTableVersion();
throw new ArgumentNullException(nameof(tableVersion));
}
try
{
ClearMetadataIfUnavailable(entry);
return await orleansQueries.UpdateMembershipRowAsync(this.clusterId, entry, tableVersion.VersionEtag);
}
catch (Exception ex)
{
LogDebugUpdateRowFailed(ex);
throw;
}
}
public async Task UpdateIAmAlive(MembershipEntry entry)
{
LogTraceUpdateIAmAlive(entry);
if (entry == null)
{
LogDebugUpdateIAmAliveAbortedNullEntry();
throw new ArgumentNullException(nameof(entry));
}
try
{
await orleansQueries.UpdateIAmAliveTimeAsync(this.clusterId, entry.SiloAddress, entry.IAmAliveTime);
}
catch (Exception ex)
{
LogDebugUpdateIAmAliveFailed(ex);
throw;
}
}
public async Task DeleteMembershipTableEntries(string clusterId)
{
LogTraceDeleteMembershipTableEntries(clusterId);
try
{
await orleansQueries.DeleteMembershipTableEntriesAsync(clusterId);
}
catch (Exception ex)
{
LogDebugDeleteMembershipTableEntriesFailed(ex);
throw;
}
}
public async Task CleanupDefunctSiloEntries(DateTimeOffset beforeDate)
{
LogTraceCleanupDefunctSiloEntries(beforeDate, clusterId);
try
{
await orleansQueries.CleanupDefunctSiloEntriesAsync(beforeDate, this.clusterId);
}
catch (Exception ex)
{
LogDebugCleanupDefunctSiloEntriesFailed(ex);
throw;
}
}
private async Task<bool> InitTableAsync()
{
try
{
return await orleansQueries.InsertMembershipVersionRowAsync(this.clusterId);
}
catch (Exception ex)
{
LogTraceInsertSiloMembershipVersionFailed(ex);
throw;
}
}
private void ClearMetadataIfUnavailable(MembershipEntry entry)
{
if (!orleansQueries.SupportsMembershipMetadata)
{
entry.Metadata = null;
}
}
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(InitializeMembershipTable)} called."
)]
private partial void LogTraceInitializeMembershipTable();
[LoggerMessage(
Level = LogLevel.Information,
Message = "Created new table version row."
)]
private partial void LogInfoCreatedNewTableVersionRow();
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(ReadRow)} called with key: {{Key}}."
)]
private partial void LogTraceReadRow(SiloAddress key);
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(ReadRow)} failed"
)]
private partial void LogDebugReadRowFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(ReadAll)} called."
)]
private partial void LogTraceReadAll();
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(ReadAll)} failed"
)]
private partial void LogDebugReadAllFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(InsertRow)} called with entry {{Entry}} and tableVersion {{TableVersion}}."
)]
private partial void LogTraceInsertRow(MembershipEntry entry, TableVersion tableVersion);
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(InsertRow)} aborted due to null check. MembershipEntry is null."
)]
private partial void LogDebugInsertRowAbortedNullEntry();
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(InsertRow)} aborted due to null check. TableVersion is null "
)]
private partial void LogDebugInsertRowAbortedNullTableVersion();
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(InsertRow)} failed"
)]
private partial void LogDebugInsertRowFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(IMembershipTable)}.{nameof(UpdateRow)} called with entry {{Entry}}, etag {{ETag}} and tableVersion {{TableVersion}}."
)]
private partial void LogTraceUpdateRow(MembershipEntry entry, string etag, TableVersion tableVersion);
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(UpdateRow)} aborted due to null check. MembershipEntry is null."
)]
private partial void LogDebugUpdateRowAbortedNullEntry();
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(UpdateRow)} aborted due to null check. TableVersion is null"
)]
private partial void LogDebugUpdateRowAbortedNullTableVersion();
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(UpdateRow)} failed"
)]
private partial void LogDebugUpdateRowFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(IMembershipTable)}.{nameof(UpdateIAmAlive)} called with entry {{Entry}}."
)]
private partial void LogTraceUpdateIAmAlive(MembershipEntry entry);
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(UpdateIAmAlive)} aborted due to null check. MembershipEntry is null."
)]
private partial void LogDebugUpdateIAmAliveAbortedNullEntry();
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(UpdateIAmAlive)} failed"
)]
private partial void LogDebugUpdateIAmAliveFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(IMembershipTable)}.{nameof(DeleteMembershipTableEntries)} called with clusterId {{ClusterId}}."
)]
private partial void LogTraceDeleteMembershipTableEntries(string clusterId);
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(DeleteMembershipTableEntries)} failed"
)]
private partial void LogDebugDeleteMembershipTableEntriesFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = $"{nameof(IMembershipTable)}.{nameof(CleanupDefunctSiloEntries)} called with beforeDate {{beforeDate}} and clusterId {{ClusterId}}."
)]
private partial void LogTraceCleanupDefunctSiloEntries(DateTimeOffset beforeDate, string clusterId);
[LoggerMessage(
Level = LogLevel.Debug,
Message = $"{nameof(AdoNetClusteringTable)}.{nameof(CleanupDefunctSiloEntries)} failed"
)]
private partial void LogDebugCleanupDefunctSiloEntriesFailed(Exception exception);
[LoggerMessage(
Level = LogLevel.Trace,
Message = "Insert silo membership version failed"
)]
private partial void LogTraceInsertSiloMembershipVersionFailed(Exception exception);
}
}