diff --git a/src/Orleans.Streaming/Providers/ClientStreamingProviderRuntime.cs b/src/Orleans.Streaming/Providers/ClientStreamingProviderRuntime.cs index d8a4903a6f7..674d56f4b0b 100644 --- a/src/Orleans.Streaming/Providers/ClientStreamingProviderRuntime.cs +++ b/src/Orleans.Streaming/Providers/ClientStreamingProviderRuntime.cs @@ -31,7 +31,7 @@ public ClientStreamingProviderRuntime( this.clientContext = clientContext; this.runtimeClient = serviceProvider.GetService()!; // Registered by DefaultClientServices. grainBasedPubSub = new GrainBasedPubSubRuntime(GrainFactory); - var tmp = new ImplicitStreamPubSub(this.grainFactory, this.implicitSubscriberTable); + var tmp = new ImplicitStreamPubSub(this.implicitSubscriberTable); implicitPubSub = tmp; combinedGrainBasedAndImplicitPubSub = new StreamPubSubImpl(grainBasedPubSub, tmp); streamDirectory = new StreamDirectory(); diff --git a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs index b8d1716c389..67a05bfabfe 100644 --- a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs +++ b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs @@ -41,7 +41,7 @@ public SiloStreamProviderRuntime( this.runtimeClient = runtimeClient; this.logger = this.loggerFactory.CreateLogger(); this.grainBasedPubSub = new GrainBasedPubSubRuntime(this.GrainFactory); - var tmp = new ImplicitStreamPubSub(this.runtimeClient.InternalGrainFactory, implicitStreamSubscriberTable); + var tmp = new ImplicitStreamPubSub(implicitStreamSubscriberTable); this.implictPubSub = tmp; this.combinedGrainBasedAndImplicitPubSub = new StreamPubSubImpl(this.grainBasedPubSub, tmp); } diff --git a/src/Orleans.Streaming/PubSub/ImplicitStreamPubSub.cs b/src/Orleans.Streaming/PubSub/ImplicitStreamPubSub.cs index c87b8126e34..f6de758456f 100644 --- a/src/Orleans.Streaming/PubSub/ImplicitStreamPubSub.cs +++ b/src/Orleans.Streaming/PubSub/ImplicitStreamPubSub.cs @@ -9,18 +9,16 @@ namespace Orleans.Streams { internal class ImplicitStreamPubSub : IStreamPubSub { - private readonly IInternalGrainFactory grainFactory; private readonly ImplicitStreamSubscriberTable implicitTable; - public ImplicitStreamPubSub(IInternalGrainFactory grainFactory, ImplicitStreamSubscriberTable implicitPubSubTable) + public ImplicitStreamPubSub(ImplicitStreamSubscriberTable implicitSubscriberTable) { - if (implicitPubSubTable == null) + if (implicitSubscriberTable == null) { - throw new ArgumentNullException(nameof(implicitPubSubTable)); + throw new ArgumentNullException(nameof(implicitSubscriberTable)); } - this.grainFactory = grainFactory; - this.implicitTable = implicitPubSubTable; + this.implicitTable = implicitSubscriberTable; } public Task> RegisterProducer(QualifiedStreamId streamId, GrainId streamProducer) @@ -28,7 +26,7 @@ public Task> RegisterProducer(QualifiedStreamId st ISet result = new HashSet(); if (!ImplicitStreamSubscriberTable.IsImplicitSubscribeEligibleNameSpace(streamId.GetNamespace())) return Task.FromResult(result); - IDictionary implicitSubscriptions = implicitTable.GetImplicitSubscribers(streamId, this.grainFactory); + IDictionary implicitSubscriptions = implicitTable.GetImplicitSubscribers(streamId); foreach (var kvp in implicitSubscriptions) { GuidId subscriptionId = GuidId.GetGuidId(kvp.Key); @@ -84,7 +82,7 @@ public Task> GetAllSubscriptions(QualifiedStreamId stre } else { - var implicitConsumers = this.implicitTable.GetImplicitSubscribers(streamId, grainFactory); + var implicitConsumers = this.implicitTable.GetImplicitSubscribers(streamId); var subscriptions = implicitConsumers.Select(consumer => { var grainId = consumer.Value; diff --git a/src/Orleans.Streaming/PubSub/ImplicitStreamSubscriberTable.cs b/src/Orleans.Streaming/PubSub/ImplicitStreamSubscriberTable.cs index e0caacc7c53..e987157ee78 100644 --- a/src/Orleans.Streaming/PubSub/ImplicitStreamSubscriberTable.cs +++ b/src/Orleans.Streaming/PubSub/ImplicitStreamSubscriberTable.cs @@ -107,14 +107,13 @@ private Cache BuildCache(MajorMinorVersion version, ImmutableDictionary - /// Retrieve a map of implicit subscriptionsIds to implicit subscribers, given a stream ID. This method throws an exception if there's no namespace associated with the stream ID. + /// Retrieves a map of implicit subscription IDs to implicit subscriber grain IDs for the specified stream. /// /// A stream ID. - /// The grain factory used to get consumer references. - /// A set of GrainId that are implicitly subscribed grains. They are expected to support the streaming consumer extension. + /// A dictionary mapping subscription IDs to implicitly subscribed grain IDs. The grains are expected to support the streaming consumer extension. /// The stream ID doesn't have an associated namespace. /// Internal invariant violation. - internal Dictionary GetImplicitSubscribers(QualifiedStreamId streamId, IInternalGrainFactory grainFactory) + internal Dictionary GetImplicitSubscribers(QualifiedStreamId streamId) { var streamNamespace = streamId.GetNamespace(); if (!IsImplicitSubscribeEligibleNameSpace(streamNamespace))