Repository navigation
Expand file tree
/
Copy pathOrleansNodeRequestExecutor.cs
More file actions
54 lines (51 loc) · 2.63 KB
/
Copy pathOrleansNodeRequestExecutor.cs
File metadata and controls
54 lines (51 loc) · 2.63 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
using KeyLoad.Orleans;
using ManagedCode.Communication.CQRS;
using Microsoft.Extensions.Options;
using Orleans.Serialization;
namespace KeyLoad.Server;
/// <summary>Bounds one admitted native request, including its cohort check and complete CQRS stream.</summary>
internal sealed class OrleansNodeRequestExecutor(IOptions<GrainRoutingOptions> routingOptions,
ILogger<OrleansNode> logger)
{
internal async Task<GrainOperationReply> ExecuteAsync(IGrainFactory grains, IServiceProvider services,
PhysicalShardCatalogStartup? catalog, Guid requestId, string signedRequest, bool command,
CancellationToken cancellationToken)
{
var clock = services.GetRequiredService<TimeProvider>();
using var deadline = new CancellationTokenSource(routingOptions.Value.ExecutionLifetime, clock);
using var execution = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, deadline.Token);
await services.GetRequiredService<ReplicaSiloDiscoveryClient>()
.EnsureCompatibleCohortAsync(execution.Token).ConfigureAwait(false);
if (catalog is not null)
{
if (!catalog.IsReady)
{
throw Errors.Fail(ErrorCode.OwnershipLost, PhysicalShardCatalogFence.NotReady);
}
await catalog.EnsureAdmissionAsync(execution.Token).ConfigureAwait(false);
}
var reply = await DrainAsync(grains, services, requestId, signedRequest, command, clock,
execution.Token, cancellationToken).ConfigureAwait(false);
if (reply.Error is { } error)
{
throw Errors.Fail(error, reply.SafeDetail ?? OrleansNodeProtocol.ReplyRejected);
}
return reply;
}
private async Task<GrainOperationReply> DrainAsync(IGrainFactory grains, IServiceProvider services,
Guid requestId, string signedRequest, bool command, TimeProvider clock,
CancellationToken executionToken, CancellationToken callerToken)
{
try
{
return await GrainRequestStreamConsumer.DrainAsync(
createStream: token => grains.GetGrain<IRequestGrain>(requestId).ExecuteStreamAsync(signedRequest, token),
serializer: services.GetRequiredService<Serializer<CqrsStreamChunk<GrainRequestProgress, GrainOperationReply>>>(),
requestId: requestId, clock: clock, cancellationToken: executionToken, options: routingOptions).ConfigureAwait(false);
}
catch (Exception failure) when (NativeCqrsBoundaryErrors.IsNonFatal(failure))
{
throw OrleansRpcFailure.Translate(failure, command, requestId, logger, callerToken);
}
}
}