Repository navigation
Expand file tree
/
Copy pathOrleansNodeRequestExecutor.cs
More file actions
89 lines (82 loc) · 5.5 KB
/
Copy pathOrleansNodeRequestExecutor.cs
File metadata and controls
89 lines (82 loc) · 5.5 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
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 static Task CloseAsync(IGrainFactory grains, IServiceProvider services,
Guid connectionId, CancellationToken cancellationToken)
=> grains.GetGrain<IConnectionGrain>(connectionId).CloseAsync(
services.GetRequiredService<GrainRequestCodec>().CreateConnectionClose(connectionId), cancellationToken);
internal Task<GrainOperationReply> ExecuteAsync(IGrainFactory? factory, IServiceProvider? runtime,
PhysicalShardCatalogStartup? catalog, bool databaseReady, bool requireCatalogAdmission,
Guid requestId, string signedRequest, bool command, CancellationToken cancellationToken)
=> ExecuteCoreAsync(factory, runtime, catalog, databaseReady, requireCatalogAdmission,
requestId, signedRequest, command, null, cancellationToken);
internal Task<GrainOperationReply> ExecuteAsync(NativeNodeRequestTarget target, bool requireCatalogAdmission,
Guid requestId, string signedRequest, bool command, CancellationToken cancellationToken)
=> ExecuteAsync(target.Factory, target.Runtime, target.Catalog, target.DatabaseReady,
requireCatalogAdmission, requestId, signedRequest, command, cancellationToken);
internal Task<GrainOperationReply> ExecuteWithProgressAsync(NativeNodeRequestTarget target,
Guid requestId, string signedRequest, Func<GrainRequestProgress, CancellationToken, ValueTask> progress,
CancellationToken cancellationToken)
=> ExecuteCoreAsync(target.Factory, target.Runtime, target.Catalog, target.DatabaseReady, true,
requestId, signedRequest, false, progress, cancellationToken);
private async Task<GrainOperationReply> ExecuteCoreAsync(IGrainFactory? factory, IServiceProvider? runtime,
PhysicalShardCatalogStartup? catalog, bool databaseReady, bool requireCatalogAdmission,
Guid requestId, string signedRequest, bool command,
Func<GrainRequestProgress, CancellationToken, ValueTask>? progress, CancellationToken cancellationToken)
{
if (requireCatalogAdmission && (!databaseReady || catalog is null || !catalog.IsReady))
{ throw Errors.Fail(ErrorCode.OwnershipLost, PhysicalShardCatalogFence.NotReady); }
var grains = factory ?? throw Errors.Fail(ErrorCode.OwnershipLost, OrleansNodeProtocol.RoutingUnavailable);
var services = runtime ?? throw Errors.Fail(ErrorCode.OwnershipLost, OrleansNodeProtocol.RoutingUnavailable);
var connectionId = NativeConnectionExecutionIdentity.Resolve(services);
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 (requireCatalogAdmission && 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, connectionId, requestId, signedRequest, command, clock,
progress, 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 connectionId, Guid requestId, string signedRequest, bool command, TimeProvider clock,
Func<GrainRequestProgress, CancellationToken, ValueTask>? progress,
CancellationToken executionToken, CancellationToken callerToken)
{
try
{
Func<CancellationToken, IAsyncEnumerable<CqrsStreamChunk<GrainRequestProgress, GrainOperationReply>>> source =
token => grains.GetGrain<IConnectionGrain>(connectionId).ExecuteStreamAsync(signedRequest, token);
var serializer = services.GetRequiredService<Serializer<CqrsStreamChunk<GrainRequestProgress, GrainOperationReply>>>();
var purpose = new NativeCqrsStreamPurpose(() => services.GetRequiredService<GrainRequestCodec>()
.VerifyRequest(signedRequest, requestId));
return progress is null
? await GrainRequestStreamConsumer.DrainWithPurposeAsync(source, serializer, requestId, clock,
routingOptions, purpose, executionToken).ConfigureAwait(false)
: await GrainRequestStreamConsumer.DrainWithProgressAsync(source, serializer, requestId, clock,
routingOptions, purpose, progress, executionToken).ConfigureAwait(false);
}
catch (Exception failure) when (NativeCqrsBoundaryErrors.IsNonFatal(failure))
{
throw OrleansRpcFailure.Translate(failure, command, requestId, logger, callerToken);
}
}
}