Skip to content

Commit e05fc11

Browse files
committed
perf(runtime): reduce invokable pooling overhead
1 parent 33f976b commit e05fc11

19 files changed

Lines changed: 392 additions & 753 deletions

src/Orleans.CodeGenerator/InvokableGenerator.cs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -897,6 +897,7 @@ private List<InvokerFieldDescription> GetFieldDescriptions(
897897
constructor.MethodKind == MethodKind.Constructor
898898
&& constructor.HasAttribute(LibraryTypes.GeneratedActivatorConstructorAttribute));
899899
if (method.MethodTypeParameters.Count == 0
900+
&& method.Method.Parameters.Length >= 2
900901
&& method.CustomInitializerMethods.Count == 0
901902
&& !requiresDependencyInjection
902903
&& IsPoolableBaseType(baseClassType))

src/Orleans.Core/Messaging/ClientMessageCenter.cs

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -191,8 +191,8 @@ public void SendMessage(Message msg)
191191
return;
192192
}
193193

194-
connection.Send(msg);
195194
LogSendingMessage(msg, connection.RemoteEndPoint);
195+
connection.Send(msg);
196196
}
197197
else
198198
{
@@ -211,9 +211,8 @@ async Task SendAsync(ValueTask<Connection?> task, Message message)
211211
return;
212212
}
213213

214-
connection.Send(message);
215-
216214
LogSendingMessage(message, connection.RemoteEndPoint);
215+
connection.Send(message);
217216
}
218217
catch (Exception exception)
219218
{

src/Orleans.Core/Messaging/MessageFactory.cs

Lines changed: 3 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -32,37 +32,22 @@ public MessageFactory(DeepCopier deepCopier, ILogger<MessageFactory> logger, Mes
3232

3333
public Message CreateMessage(object? body, InvokeMethodOptions options)
3434
{
35-
var oneWay = (options & InvokeMethodOptions.OneWay) != 0;
36-
var messageBody = body;
37-
var ownsBodyObject = false;
38-
if (body is RequestBase)
39-
{
40-
messageBody = CopyBodyObject(body);
41-
ownsBodyObject = true;
42-
}
43-
else if (body is IInvokable)
44-
{
45-
ownsBodyObject = true;
46-
}
47-
4835
var message = new Message
4936
{
50-
Direction = oneWay ? Message.Directions.OneWay : Message.Directions.Request,
37+
Direction = (options & InvokeMethodOptions.OneWay) != 0 ? Message.Directions.OneWay : Message.Directions.Request,
5138
Id = GetNextCorrelationId(),
5239
IsReadOnly = (options & InvokeMethodOptions.ReadOnly) != 0,
5340
IsUnordered = (options & InvokeMethodOptions.Unordered) != 0,
5441
IsAlwaysInterleave = (options & InvokeMethodOptions.AlwaysInterleave) != 0,
55-
BodyObject = messageBody,
56-
OwnsBodyObject = ownsBodyObject,
42+
BodyObject = body,
43+
OwnsBodyObject = body is IInvokable,
5744
RequestContextData = RequestContextExtensions.Export(_deepCopier),
5845
};
5946

6047
_messagingTrace.OnCreateMessage(message);
6148
return message;
6249
}
6350

64-
private object CopyBodyObject(object body) => _deepCopier.Copy(body)!;
65-
6651
private CorrelationId GetNextCorrelationId()
6752
{
6853
var id = _seed ^ Interlocked.Increment(ref _nextId);

src/Orleans.Core/Runtime/CallbackData.cs

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -210,6 +210,20 @@ public void OnHostShutdown()
210210
this.context.Complete(Response.FromException(exception));
211211
}
212212

213+
public void OnSendFailure(Exception exception)
214+
{
215+
if (!TryComplete())
216+
{
217+
return;
218+
}
219+
220+
this.stopwatch.Stop();
221+
this.shared.Unregister(this.Message);
222+
DisposeCancellationRegistration();
223+
_applicationRequestInstruments.OnAppRequestsEnd((long)this.stopwatch.Elapsed.TotalMilliseconds);
224+
this.context.Complete(Response.FromException(exception));
225+
}
226+
213227
public void DoCallback(Message response)
214228
{
215229
if (!TryComplete())

src/Orleans.Core/Runtime/GrainReferenceRuntime.cs

Lines changed: 38 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66
using Orleans.CodeGeneration;
77
using Orleans.GrainReferences;
88
using Orleans.Metadata;
9+
using Orleans.Serialization;
910
using Orleans.Serialization.Invocation;
1011

1112
namespace Orleans.Runtime
@@ -17,20 +18,23 @@ internal class GrainReferenceRuntime : IGrainReferenceRuntime
1718
private readonly IGrainCancellationTokenRuntime cancellationTokenRuntime;
1819
private readonly IOutgoingGrainCallFilter[] filters;
1920
private readonly Action<GrainReference, IResponseCompletionSource, IInvokable, InvokeMethodOptions> sendRequest;
21+
private readonly DeepCopier requestCopier;
2022

2123
public GrainReferenceRuntime(
2224
IRuntimeClient runtimeClient,
2325
IGrainCancellationTokenRuntime cancellationTokenRuntime,
2426
IEnumerable<IOutgoingGrainCallFilter> outgoingCallFilters,
2527
GrainReferenceActivator referenceActivator,
26-
GrainInterfaceTypeResolver interfaceTypeResolver)
28+
GrainInterfaceTypeResolver interfaceTypeResolver,
29+
DeepCopier requestCopier)
2730
{
2831
this.RuntimeClient = runtimeClient;
2932
this.cancellationTokenRuntime = cancellationTokenRuntime;
3033
this.referenceActivator = referenceActivator;
3134
this.interfaceTypeResolver = interfaceTypeResolver;
3235
this.filters = outgoingCallFilters.ToArray();
33-
this.sendRequest = (GrainReference reference, IResponseCompletionSource callback, IInvokable body, InvokeMethodOptions options) => RuntimeClient.SendRequest(reference, body, callback, options);
36+
this.requestCopier = requestCopier;
37+
this.sendRequest = SendFilteredRequest;
3438
}
3539

3640
public IRuntimeClient RuntimeClient { get; private set; }
@@ -64,25 +68,26 @@ public ValueTask InvokeMethodAsync(GrainReference reference, IInvokable request,
6468
public void InvokeMethod(GrainReference reference, IInvokable request, InvokeMethodOptions options)
6569
{
6670
Debug.Assert((options & InvokeMethodOptions.OneWay) != 0);
67-
var disposeRequest = true;
71+
var requestTransferred = false;
6872

6973
try
7074
{
7175
// TODO: Remove expensive interface type check
7276
if (filters.Length == 0 && request is not IOutgoingGrainCallFilter)
7377
{
7478
SetGrainCancellationTokensTarget(reference, request);
79+
requestTransferred = true;
7580
this.RuntimeClient.SendRequest(reference, request, context: null, options);
7681
}
7782
else
7883
{
7984
InvokeMethodWithFiltersAsync(reference, request, options).AsTask().Ignore();
80-
disposeRequest = false;
85+
requestTransferred = true;
8186
}
8287
}
8388
finally
8489
{
85-
if (disposeRequest)
90+
if (!requestTransferred)
8691
{
8792
DisposeRequest(request);
8893
}
@@ -121,60 +126,48 @@ private async ValueTask InvokeMethodWithFiltersAsync(GrainReference reference, I
121126
private ValueTask<TResult?> InvokeMethodAsyncCore<TResult>(GrainReference reference, IInvokable request, InvokeMethodOptions options)
122127
{
123128
ResponseCompletionSource<TResult>? responseCompletionSource = null;
129+
var requestTransferred = false;
124130
try
125131
{
126132
SetGrainCancellationTokensTarget(reference, request);
127133
responseCompletionSource = ResponseCompletionSourcePool.Get<TResult>();
134+
requestTransferred = true;
128135
this.RuntimeClient.SendRequest(reference, request, responseCompletionSource, options);
129-
return CompleteInvokeAsync(responseCompletionSource, request, options);
136+
return responseCompletionSource.AsValueTask();
130137
}
131138
catch
132139
{
133140
responseCompletionSource?.Reset();
134-
DisposeRequest(request);
135-
throw;
136-
}
137-
}
141+
if (!requestTransferred)
142+
{
143+
DisposeRequest(request);
144+
}
138145

139-
private static async ValueTask<TResult?> CompleteInvokeAsync<TResult>(ResponseCompletionSource<TResult> responseCompletionSource, IInvokable request, InvokeMethodOptions options)
140-
{
141-
try
142-
{
143-
return await responseCompletionSource.AsValueTask();
144-
}
145-
finally
146-
{
147-
DisposeRequest(request);
146+
throw;
148147
}
149148
}
150149

151150
private ValueTask InvokeMethodAsyncCore(GrainReference reference, IInvokable request, InvokeMethodOptions options)
152151
{
153152
ResponseCompletionSource? responseCompletionSource = null;
153+
var requestTransferred = false;
154154
try
155155
{
156156
SetGrainCancellationTokensTarget(reference, request);
157157
responseCompletionSource = ResponseCompletionSourcePool.Get();
158+
requestTransferred = true;
158159
this.RuntimeClient.SendRequest(reference, request, responseCompletionSource, options);
159-
return CompleteInvokeAsync(responseCompletionSource, request, options);
160+
return responseCompletionSource.AsVoidValueTask();
160161
}
161162
catch
162163
{
163164
responseCompletionSource?.Reset();
164-
DisposeRequest(request);
165-
throw;
166-
}
167-
}
165+
if (!requestTransferred)
166+
{
167+
DisposeRequest(request);
168+
}
168169

169-
private static async ValueTask CompleteInvokeAsync(ResponseCompletionSource responseCompletionSource, IInvokable request, InvokeMethodOptions options)
170-
{
171-
try
172-
{
173-
await responseCompletionSource.AsVoidValueTask();
174-
}
175-
finally
176-
{
177-
DisposeRequest(request);
170+
throw;
178171
}
179172
}
180173

@@ -186,6 +179,18 @@ private static void DisposeRequest(IInvokable request)
186179
}
187180
}
188181

182+
private void SendFilteredRequest(
183+
GrainReference reference,
184+
IResponseCompletionSource callback,
185+
IInvokable request,
186+
InvokeMethodOptions options)
187+
{
188+
var messageRequest = request is RequestBase
189+
? (IInvokable)this.requestCopier.Copy(request)!
190+
: request;
191+
RuntimeClient.SendRequest(reference, messageRequest, callback, options);
192+
}
193+
189194
public object Cast(IAddressable grain, Type grainInterface)
190195
{
191196
var grainId = grain.GetGrainId();

src/Orleans.Core/Runtime/OutgoingCallInvoker.cs

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -113,8 +113,22 @@ public async Task Invoke()
113113
// Finally call the root-level invoker.
114114
stage++;
115115
var responseCompletionSource = ResponseCompletionSourcePool.Get();
116-
this.sendRequest(this.grainReference, responseCompletionSource, this.request, this.options);
117-
this.Response = await responseCompletionSource.AsValueTask().ConfigureAwait(false);
116+
var requestSent = false;
117+
try
118+
{
119+
this.sendRequest(this.grainReference, responseCompletionSource, this.request, this.options);
120+
requestSent = true;
121+
this.Response = await responseCompletionSource.AsValueTask().ConfigureAwait(false);
122+
}
123+
catch
124+
{
125+
if (!requestSent)
126+
{
127+
responseCompletionSource.Reset();
128+
}
129+
130+
throw;
131+
}
118132

119133
return;
120134
}

0 commit comments

Comments
 (0)