Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 18 additions & 10 deletions src/Orleans.Runtime/Core/InsideRuntimeClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ internal sealed partial class InsideRuntimeClient : IRuntimeClient, ILifecyclePa
private IGrainCallCancellationManager _cancellationManager;
private HostedClient hostedClient;

private HostedClient HostedClient => this.hostedClient ??= this.ServiceProvider.GetRequiredService<HostedClient>();
private HostedClient HostedClient => this.hostedClient;
private readonly MessageFactory messageFactory;
private IGrainReferenceRuntime grainReferenceRuntime;
private Task callbackTimerTask;
Expand Down Expand Up @@ -113,15 +113,25 @@ public InsideRuntimeClient(

public GrainFactory ConcreteGrainFactory { get; }

private GrainLocator GrainLocator
=> this.grainLocator ?? (this.grainLocator = this.ServiceProvider.GetRequiredService<GrainLocator>());
private GrainLocator GrainLocator => this.grainLocator;

private List<IIncomingGrainCallFilter> GrainCallFilters
=> this.grainCallFilters ??= new List<IIncomingGrainCallFilter>(this.ServiceProvider.GetServices<IIncomingGrainCallFilter>());
private List<IIncomingGrainCallFilter> GrainCallFilters => this.grainCallFilters;

private MessageCenter MessageCenter => this.messageCenter ?? (this.messageCenter = this.ServiceProvider.GetRequiredService<MessageCenter>());
private MessageCenter MessageCenter => this.messageCenter;

public IGrainReferenceRuntime GrainReferenceRuntime => this.grainReferenceRuntime ?? (this.grainReferenceRuntime = this.ServiceProvider.GetRequiredService<IGrainReferenceRuntime>());
public IGrainReferenceRuntime GrainReferenceRuntime => this.grainReferenceRuntime;

internal void ConsumeServices()
{
this.grainLocator = this.ServiceProvider.GetRequiredService<GrainLocator>();
this.grainCallFilters = new List<IIncomingGrainCallFilter>(this.ServiceProvider.GetServices<IIncomingGrainCallFilter>());
this.messageCenter = this.ServiceProvider.GetRequiredService<MessageCenter>();
this.grainReferenceRuntime = this.ServiceProvider.GetRequiredService<IGrainReferenceRuntime>();
this.hostedClient = this.ServiceProvider.GetRequiredService<HostedClient>();
_cancellationManager = this.ServiceProvider.GetRequiredService<IGrainCallCancellationManager>();
sharedCallbackData.CancellationManager = _cancellationManager;
systemSharedCallbackData.CancellationManager = _cancellationManager;
}

public void SendRequest(
GrainReference target,
Expand Down Expand Up @@ -600,9 +610,7 @@ public void BreakOutstandingMessagesToSilo(SiloAddress deadSilo)

public void Participate(ISiloLifecycle lifecycle)
{
_cancellationManager = this.ServiceProvider.GetRequiredService<IGrainCallCancellationManager>();
sharedCallbackData.CancellationManager = _cancellationManager;
systemSharedCallbackData.CancellationManager = _cancellationManager;
ConsumeServices();
lifecycle.Subscribe<InsideRuntimeClient>(ServiceLifecycleStage.RuntimeInitialize, OnRuntimeInitializeStart, OnRuntimeInitializeStop);
}

Expand Down
43 changes: 43 additions & 0 deletions test/Orleans.Runtime.Tests/InsideRuntimeClientTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
using System;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Orleans.Runtime;
using Orleans.TestingHost;
using TestExtensions;
using Xunit;

namespace Tester;

[TestCategory("BVT")]
public class InsideRuntimeClientTests
{
[Fact]
public async Task LateRejectionAfterSiloDisposalDoesNotResolveServices()
{
var builder = new InProcessTestClusterBuilder(1);
builder.ConfigureHost(hostBuilder => TestDefaultConfiguration.ConfigureHostConfiguration(hostBuilder.Configuration));
await using var cluster = builder.Build();
await cluster.DeployAsync();

var silo = cluster.Silos[0];
var runtimeClient = silo.ServiceProvider.GetRequiredService<InsideRuntimeClient>();
var rejection = new Message
{
Direction = Message.Directions.Response,
Result = Message.ResponseTypes.Rejection,
TargetSilo = silo.SiloAddress,
TargetGrain = GrainId.Create("caller", Guid.NewGuid().ToString()),
SendingGrain = GrainId.Create("target", Guid.NewGuid().ToString()),
BodyObject = new RejectionResponse
{
RejectionType = Message.RejectionTypes.Unrecoverable,
RejectionInfo = "The outbound queue is stopped",
},
};

await silo.StopSiloAsync(stopGracefully: false);
await silo.DisposeAsync();

Assert.Null(Record.Exception(() => runtimeClient.ReceiveResponse(rejection)));
}
}