Repository navigation
Expand file tree
/
Copy pathRecurringDueGrainService.cs
More file actions
199 lines (187 loc) · 9.1 KB
/
Copy pathRecurringDueGrainService.cs
File metadata and controls
199 lines (187 loc) · 9.1 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
using KeyLoad.Core;
using KeyLoad.Core.Features.Messaging;
using KeyLoad.Replication;
using Microsoft.Extensions.Options;
namespace KeyLoad.Orleans;
/// <summary>Per-silo leader watcher which reads canonical due hints and drains them through partition grains.</summary>
/// <param name="id">The Orleans grain-service identity for this silo instance.</param>
/// <param name="silo">The hosting Orleans silo.</param>
/// <param name="loggerFactory">The factory used by the grain-service base class.</param>
/// <param name="database">The node-local database and canonical read owner.</param>
/// <param name="consensus">The local replicated-partition leadership and read-barrier service.</param>
/// <param name="grainFactory">The Orleans factory for partition coordinator activations.</param>
/// <param name="clock">The shared UTC clock used for due discovery and bounded polling.</param>
/// <param name="diagnostics">Safe operational diagnostics for rejected pages and dispatch outcomes.</param>
/// <param name="options">The centrally validated dispatch and discovery scheduling settings.</param>
/// <param name="journalAdmission">The native scheduling admission and early shutdown token.</param>
public sealed class RecurringDueGrainService(GrainId id, Silo silo,
Microsoft.Extensions.Logging.ILoggerFactory loggerFactory,
DatabaseEngine database, ReplicaConsensus consensus, IGrainFactory grainFactory, TimeProvider clock,
Microsoft.Extensions.Logging.ILogger<RecurringDueGrainService> diagnostics,
IOptions<DueCoordinationOptions> options, RuntimeJournalAdmission journalAdmission)
: GrainService(id, silo, loggerFactory), IRecurringDueGrainService
{
private const int NoRejectedDueHints = 0;
private Task? loop;
/// <summary>Completes native per-silo grain-service initialization.</summary>
/// <param name="serviceProvider">The active silo service provider.</param>
/// <returns>The base grain-service initialization task.</returns>
public override Task Init(IServiceProvider serviceProvider) => base.Init(serviceProvider);
/// <summary>Starts the joined due-discovery loop after native grain-service startup.</summary>
/// <returns>The base grain-service startup task.</returns>
protected override Task StartInBackground()
{
var started = base.StartInBackground();
loop = RunAsync(StoppedCancellationTokenSource.Token);
return started;
}
/// <summary>Stops admission, cancels discovery and joins the active loop.</summary>
/// <returns>The joined base and due-loop shutdown task.</returns>
public override async Task Stop()
{
var stopped = base.Stop();
if (loop is { } active)
{
await Task.WhenAll(stopped, active).ConfigureAwait(true);
return;
}
await stopped.ConfigureAwait(true);
}
private async Task RunAsync(CancellationToken stopped)
{
using var lifetime = CancellationTokenSource.CreateLinkedTokenSource(stopped, journalAdmission.SchedulingToken);
var cancellationToken = lifetime.Token;
var turns = new RecurringDueTurnState();
try
{
await consensus.TransportReady.WaitAsync(cancellationToken).ConfigureAwait(true);
while (true)
{
cancellationToken.ThrowIfCancellationRequested();
var cycleStarted = clock.GetTimestamp();
var observedPosition = consensus.AppliedPosition;
await turns.RunAsync(options, RunCycleAsync, RunQueueCycleAsync, RunTransferCycleAsync,
cancellationToken).ConfigureAwait(true);
await RecurringDueWait.WaitForChangeOrFallbackAsync(consensus, observedPosition,
options.Value.PollInterval, clock, cancellationToken).ConfigureAwait(true);
await RecurringDueWait.WaitForMinimumCadenceAsync(cycleStarted, options.Value.MinimumCycleCadence,
clock, cancellationToken).ConfigureAwait(true);
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
return;
}
catch (KeyLoadException failure)
{
RecurringDueDiagnostics.ServiceFault(diagnostics, failure.Code);
}
catch (OperationCanceledException)
{
RecurringDueDiagnostics.ServiceFault(diagnostics, ErrorCode.OwnershipLost);
}
catch (Exception failure) when (GrainBoundaryErrors.Handles(failure))
{
RecurringDueDiagnostics.ServiceFault(diagnostics, ErrorCode.OwnershipLost);
}
}
private Task<RemoteTransferPendingCursor?> RunTransferCycleAsync(RemoteTransferPendingCursor? cursor, CancellationToken token)
=> RemoteTransferServiceCycle.RunAsync(database, consensus, grainFactory, clock, diagnostics, options, cursor, token);
private Task<QueueDeadlineCursor?> RunQueueCycleAsync(QueueDeadlineCursor? cursor, CancellationToken cancellationToken)
=> QueueDeadlineServiceCycle.RunAsync(database, consensus, grainFactory, clock, diagnostics, options, cursor, cancellationToken);
private async Task<DueSweepCursor?> RunCycleAsync(DueSweepCursor? cursor,
CancellationToken cancellationToken)
{
try
{
var page = await ReadAndDispatchPage(cursor, cancellationToken).ConfigureAwait(true);
cursor = page.Cursor;
if (page.Rejected.Length > NoRejectedDueHints)
{
RecurringDueDiagnostics.PageRejected(diagnostics, page.Rejected.Length);
}
}
catch (KeyLoadException failure) when (failure.Code == ErrorCode.Corruption
&& failure.Message == DueWorkProtocol.KeyExceedsBound)
{
cursor = DueWorkCursor.DeferNext(database.Store.Identity, cursor);
RecurringDueDiagnostics.PrefixRejected(diagnostics, failure.Code);
}
catch (KeyLoadException failure)
{
RecurringDueDiagnostics.ServiceFault(diagnostics, failure.Code);
}
catch (Exception failure) when (GrainBoundaryErrors.Handles(failure))
{
RecurringDueDiagnostics.ServiceFault(diagnostics, ErrorCode.OwnershipLost);
}
return cursor;
}
private async Task<DueWorkPage> ReadAndDispatchPage(DueSweepCursor? cursor,
CancellationToken cancellationToken)
{
if (!await consensus.IsLeaderAsync(cancellationToken).ConfigureAwait(true))
{
return DueWorkDiscovery.EmptyPage(cursor, database.Store.Identity);
}
await consensus.ReadBarrierAsync(cancellationToken).ConfigureAwait(true);
if (!await consensus.IsLeaderAsync(cancellationToken).ConfigureAwait(true))
{
return DueWorkDiscovery.EmptyPage(cursor, database.Store.Identity);
}
var page = DueWorkDiscovery.ReadPage(database, cursor, clock.GetUtcNow(), cancellationToken);
foreach (var hint in page.Jobs)
{
cancellationToken.ThrowIfCancellationRequested();
await Dispatch(hint, cancellationToken).ConfigureAwait(true);
}
return page;
}
private async Task Dispatch(DueWorkHint hint, CancellationToken cancellationToken)
{
var dispatchDeadline = options.Value.DispatchDeadline;
using var timeout = new CancellationTokenSource(dispatchDeadline, clock);
using var deadline = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, timeout.Token);
try
{
var partition = hint.Lane.Partition.AtomicPartitionId;
var coordinator = grainFactory.GetGrain<IRecurringDueCoordinatorGrain>(partition);
var dispatch = coordinator.ProcessDueAsync(hint, deadline.Token);
var result = await AwaitDispatch(dispatch, deadline.Token).ConfigureAwait(true);
if (result.Error is { } code)
{
RecurringDueDiagnostics.DispatchFailed(diagnostics, code);
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch (OperationCanceledException) when (deadline.IsCancellationRequested)
{
RecurringDueDiagnostics.DispatchDeadline(diagnostics);
}
catch (KeyLoadException failure) when (!cancellationToken.IsCancellationRequested)
{
RecurringDueDiagnostics.DispatchFailed(diagnostics, failure.Code);
}
catch (Exception failure) when (!cancellationToken.IsCancellationRequested
&& GrainBoundaryErrors.Handles(failure))
{
RecurringDueDiagnostics.DispatchFailed(diagnostics, ErrorCode.OwnershipLost);
}
}
private static async Task<DueDispatchResult> AwaitDispatch(Task<DueDispatchResult> dispatch,
CancellationToken cancellationToken)
{
try
{
return await dispatch.WaitAsync(cancellationToken).ConfigureAwait(true);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
await dispatch.ConfigureAwait(true);
throw;
}
}
}