Repository navigation
Expand file tree
/
Copy pathOrleansSiloConfiguration.cs
More file actions
210 lines (202 loc) · 12.7 KB
/
Copy pathOrleansSiloConfiguration.cs
File metadata and controls
210 lines (202 loc) · 12.7 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
using System.Collections.Immutable;
using System.Net;
using System.Security.Cryptography;
using KeyLoad.Core;
using KeyLoad.Orleans;
using KeyLoad.Query;
using KeyLoad.Replication;
using KeyLoad.Server.Features.ClusterRouting;
using KeyLoad.Server.Features.DocumentStorage;
using ManagedCode.Communication.Orleans.Converters;
using ManagedCode.Orleans.Identity.Core.Serializations;
using Microsoft.Extensions.Options;
using Orleans.Configuration;
using Orleans.Serialization;
namespace KeyLoad.Server;
internal static class OrleansSiloConfiguration
{
internal static IHost Build(PartitionHost partition, NodeOptions options, INodeAdministration administration,
ILoggerFactory loggerFactory, NativeRequestWorkOwner requestWork, NativeConnectionOwnerIdentity connectionOwner, IPAddress address,
ServerRuntimeOptions runtimeOptions, TimeProvider clock, INativePartitionMovementCapture? movementCapture, IPartitionMovementDispatcher? movementDispatcher,
IRemoteDocumentReadRouter? remoteDocuments, IRemoteBlobReadRouter? remoteBlobs, IRemotePartitionQueryRouter? remoteQueries,
IControlledDocumentCommandRouter? controlledDocuments, IGrainPartitionMovementSealedOperationObserver? sealedObserver,
CancellationToken startupCancellation)
{
var builder = Host.CreateApplicationBuilder();
builder.Services.AddSingleton(loggerFactory);
builder.Services.AddSingleton(connectionOwner);
if (sealedObserver is not null)
{ builder.Services.AddSingleton(sealedObserver); }
if (options.MembershipAuthority.RemoteDocumentReads)
{
builder.Services.AddSingleton<IPhysicalRequestPlacement, PhysicalDocumentRequestPlacement>();
if (remoteDocuments is not null)
{ builder.Services.AddSingleton(remoteDocuments); }
if (remoteBlobs is not null)
{ builder.Services.AddSingleton(remoteBlobs); }
if (remoteQueries is not null)
{ builder.Services.AddSingleton(remoteQueries); }
}
if (controlledDocuments is not null)
{ builder.Services.AddSingleton(controlledDocuments); }
if (movementCapture is not null)
{ builder.Services.AddSingleton(movementCapture); }
if (movementCapture is PartitionMovementSourceOwner sourceOwner)
{ builder.Services.AddSingleton(sourceOwner.TransferReads); }
if (movementDispatcher is not null)
{ builder.Services.AddSingleton(movementDispatcher); }
if (movementDispatcher is IPartitionMovementParent movementParent)
{ builder.Services.AddSingleton(movementParent); }
runtimeOptions.RegisterBorrowed(builder.Services);
RegisterBorrowedServices(builder.Services, partition, administration, options, requestWork, runtimeOptions, clock, startupCancellation);
builder.UseOrleans(silo => Configure(silo, options, partition.Configuration, address, runtimeOptions.Membership.Value,
runtimeOptions.Core.RuntimeJournal, runtimeOptions.DurableJobs, runtimeOptions.GrainRouting));
return builder.Build();
}
private static void RegisterBorrowedServices(IServiceCollection services, PartitionHost partition,
INodeAdministration administration, NodeOptions options, NativeRequestWorkOwner requestWork,
ServerRuntimeOptions runtimeOptions, TimeProvider clock, CancellationToken startupCancellation)
{
var peers = runtimeOptions.Peer.Value;
peers.Validate(partition.Configuration);
var expectedOwner = new PhysicalShardRecord(options.PhysicalShardId, partition.Configuration.Incarnation,
ImmutableArray.CreateRange(partition.Configuration.VoterIds),
PhysicalShardCatalogStartupProtocol.InitialPlacementEpoch);
services.AddSingleton(expectedOwner);
services.AddSingleton(partition.Database);
services.AddSingleton<ICacheMemoryBudget>(partition.CacheMemory);
services.AddSingleton<INativeAnnMaintenance>(partition.AnnMaintenance);
services.AddSingleton<INativeTextMaintenance>(partition.TextMaintenance);
services.AddSingleton<ICommitCoordinator>(partition.Coordinator);
services.AddSingleton<IReplicaEndpoint>(partition.Consensus);
services.AddSingleton(partition.Consensus);
services.AddSingleton(clock);
services.AddSingleton(requestWork);
services.AddSingleton(administration);
services.AddSingleton<QueryEngine>();
services.AddSingleton(_ => new SearchEngine(partition.Database, runtimeOptions.Core.QueryExecution, partition.TextProjection, partition.AnnMaintenance));
RegisterRequestCodec(services, partition, options);
services.AddSerializer(serialization => serialization
.AddAssembly(typeof(GrainRequestProgress).Assembly)
.AddAssembly(typeof(CqrsStreamChunkSurrogateConverter<GrainRequestProgress, GrainOperationReply>).Assembly)
.AddAssembly(typeof(ClaimsPrincipalSurrogateConverter).Assembly));
services.AddSingleton(provider => new ReplicaSiloDiscoveryState(
provider.GetRequiredService<IOptions<ReplicaConfiguration>>(),
provider.GetRequiredService<IOptions<ReplicaPeerOptions>>(),
provider.GetRequiredService<ILocalSiloDetails>(),
RuntimeJournalStorePreparation.CurrentCapabilityEvidence(partition)));
services.AddSingleton(provider => new ReplicaEnvelopeAuthenticator(provider.GetRequiredService<IOptions<ReplicaConfiguration>>(),
provider.GetRequiredService<IOptions<ReplicaPeerOptions>>(),
provider.GetRequiredService<ReplicaSiloDiscoveryState>(), provider.GetRequiredService<TimeProvider>(),
provider.GetRequiredService<IOptions<ReplicaTransportOptions>>(), provider.GetRequiredService<IOptions<ReplicaReplayLimits>>(),
logger: provider.GetService<ILogger<ReplicaEnvelopeAuthenticator>>(), canonicalDatabase: partition.Database));
services.AddSingleton<ReplicaSiloDiscoveryClient>(provider => new ReplicaSiloDiscoveryClient(
provider.GetRequiredService<IOptions<ReplicaConfiguration>>(), provider.GetRequiredService<IOptions<ReplicaPeerOptions>>(),
provider.GetRequiredService<ReplicaSiloDiscoveryState>(),
provider.GetRequiredService<ReplicaEnvelopeAuthenticator>(), provider.GetRequiredService<TimeProvider>(),
provider.GetRequiredService<IOptions<PeerDiscoveryOptions>>()));
services.AddSingleton<ReplicaGrainServiceClient>();
services.AddSingleton<ILifecycleParticipant<ISiloLifecycle>, ReplicaTransportLifecycle>();
RegisterMembershipTable(services, partition, options, startupCancellation);
}
private static void RegisterMembershipTable(IServiceCollection services, PartitionHost partition,
NodeOptions options, CancellationToken startupCancellation)
{
if (options.MembershipAuthority.Mode == MembershipAuthoritySettingsProtocol.Proxy)
{
services.AddSingleton<IMembershipTable>(provider => CreateMembershipProxy(provider, options));
return;
}
services.AddSingleton<IMembershipTable>(provider => new ReplicaMembershipTable(partition.Database, partition.Coordinator,
partition.Consensus, options.ClusterId, ClusterPrincipalPolicy.InternalPrincipalId, provider.GetRequiredService<TimeProvider>(),
options.MembershipAuthority.Mode == MembershipAuthoritySettingsProtocol.Authority
? provider.GetRequiredService<IOptions<OrleansMembershipOptions>>().Value.MaximumRows
: ReplicaMembershipProtocol.UnboundedRows,
provider.GetRequiredService<IOptions<OrleansMembershipOptions>>(),
provider.GetRequiredService<IOptions<ReplicaExecutionOptions>>(), startupCancellation));
}
private static ReplicaMembershipAuthorityClientTable CreateMembershipProxy(IServiceProvider provider, NodeOptions options)
{
var authority = options.MembershipAuthority;
var callerSecret = Convert.FromBase64String(options.PeerSecret);
var authoritySecret = Convert.FromBase64String(authority.AuthorityPeerSecret!);
try
{
var local = provider.GetRequiredService<ILocalSiloDetails>();
var settings = new ReplicaMembershipAuthorityExchangeOptions(options.ClusterId,
authority.AuthorityPhysicalShardId, authority.AuthorityIncarnation, options.PhysicalShardId,
options.Incarnation, options.PublicEndpoint, local.SiloAddress.ToParsableString(),
authority.AuthorityEndpoints.Select(endpoint => new Uri(endpoint)).ToArray(), callerSecret,
authoritySecret, provider.GetRequiredService<TimeProvider>());
return new(settings, provider.GetRequiredService<IOptions<OrleansMembershipOptions>>(),
provider.GetRequiredService<IOptions<ReplicaExecutionOptions>>());
}
finally
{
CryptographicOperations.ZeroMemory(callerSecret);
CryptographicOperations.ZeroMemory(authoritySecret);
}
}
private static void RegisterRequestCodec(IServiceCollection services, PartitionHost partition, NodeOptions options)
{
if (options.RequestCqrsProbe.Enabled)
{
services.AddSingleton<IReplicaTransportObservation>(provider =>
{
_ = provider.GetRequiredService<RequestCqrsProbeObserver>();
return partition.ApplyProbe;
});
services.AddSingleton(provider => AttachApplyObserver(partition, RequestCqrsProbeObserverFactory.Create(
options.RequestCqrsProbe, provider.GetRequiredService<IOptions<ReplicaConfiguration>>(), options.AllowPrivateNetworkHttp,
provider.GetRequiredService<ILocalSiloDetails>(), provider.GetRequiredService<IHostApplicationLifetime>(),
provider.GetRequiredService<IOptions<RequestProbeExecutionOptions>>(), provider.GetRequiredService<TimeProvider>())
?? throw new InvalidOperationException(RequestCqrsProbeProtocol.InvalidOptions)));
services.AddSingleton<IGrainRequestPhaseObserver>(provider => provider.GetRequiredService<RequestCqrsProbeObserver>());
services.AddSingleton<IGrainActivationMigrationObserver>(provider => provider.GetRequiredService<RequestCqrsProbeObserver>()
.CreateMigration(provider.GetRequiredService<ReplicaSiloDiscoveryClient>()));
}
services.AddSingleton(provider => new GrainRequestCodec(partition.Database, provider.GetRequiredService<TimeProvider>(),
provider.GetRequiredService<IOptions<GrainRoutingOptions>>())
{
PhaseObserver = provider.GetService<IGrainRequestPhaseObserver>()
});
}
private static RequestCqrsProbeObserver AttachApplyObserver(PartitionHost partition, RequestCqrsProbeObserver observer)
{
partition.ApplyProbe.Attach(observer);
return observer;
}
private static void Configure(ISiloBuilder silo, NodeOptions options, ReplicaConfiguration replica, IPAddress address,
OrleansMembershipOptions membershipOptions, IOptions<RuntimeJournalOptions> journal,
IOptions<NativeDurableJobOptions> jobs, IOptions<GrainRoutingOptions> routing)
{
silo.Configure<ClusterOptions>(cluster =>
{
cluster.ClusterId = options.ClusterId;
cluster.ServiceId = OrleansNodeProtocol.ServiceId;
});
silo.ConfigureEndpoints(address, options.SiloPort, OrleansNodeProtocol.GatewayPort);
silo.Configure<ClusterMembershipOptions>(membership =>
{
membership.IAmAliveTablePublishTimeout = membershipOptions.MembershipRefresh;
membership.TableRefreshTimeout = membershipOptions.MembershipRefresh;
});
silo.Configure<SiloMessagingOptions>(messaging => messaging.MaxMessageBodySize = Math.Max(
checked(replica.MaxAppendBytes + ReplicaTransportProtocol.MaximumMetadataBytes
+ ReplicaTransportProtocol.MaximumEnvelopeOverheadBytes),
checked(routing.Value.MaximumCompletedBytes + ReplicaTransportProtocol.MaximumEnvelopeOverheadBytes)));
silo.AddGrainService<PartitionReplicaGrainService>();
silo.AddGrainService<RecurringDueGrainService>();
silo.AddGrainService<SampleChunkGrainService>();
silo.AddActivityPropagation();
// ADR-036: owner explicitly requires these two native experimental services.
#pragma warning disable ORLEANSEXP003
silo.AddDistributedGrainDirectory();
#pragma warning restore ORLEANSEXP003
#pragma warning disable ORLEANSEXP001
silo.AddActivationRepartitioner();
#pragma warning restore ORLEANSEXP001
ConnectionGrainGraphRegistration.Register(silo, options.RequestCqrsProbe.Enabled);
NativeRuntimeJournalRegistration.Register(silo, journal, jobs);
}
}