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
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ public ClientStreamingProviderRuntime(
this.clientContext = clientContext;
this.runtimeClient = serviceProvider.GetService<IRuntimeClient>()!; // 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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public SiloStreamProviderRuntime(
this.runtimeClient = runtimeClient;
this.logger = this.loggerFactory.CreateLogger<SiloProviderRuntime>();
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);
}
Expand Down
14 changes: 6 additions & 8 deletions src/Orleans.Streaming/PubSub/ImplicitStreamPubSub.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,26 +9,24 @@ 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<ISet<PubSubSubscriptionState>> RegisterProducer(QualifiedStreamId streamId, GrainId streamProducer)
{
ISet<PubSubSubscriptionState> result = new HashSet<PubSubSubscriptionState>();
if (!ImplicitStreamSubscriberTable.IsImplicitSubscribeEligibleNameSpace(streamId.GetNamespace())) return Task.FromResult(result);

IDictionary<Guid, GrainId> implicitSubscriptions = implicitTable.GetImplicitSubscribers(streamId, this.grainFactory);
IDictionary<Guid, GrainId> implicitSubscriptions = implicitTable.GetImplicitSubscribers(streamId);
foreach (var kvp in implicitSubscriptions)
{
GuidId subscriptionId = GuidId.GetGuidId(kvp.Key);
Expand Down Expand Up @@ -84,7 +82,7 @@ public Task<List<StreamSubscription>> 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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -107,14 +107,13 @@ private Cache BuildCache(MajorMinorVersion version, ImmutableDictionary<GrainTyp
}

/// <summary>
/// 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.
/// </summary>
/// <param name="streamId">A stream ID.</param>
/// <param name="grainFactory">The grain factory used to get consumer references.</param>
/// <returns>A set of GrainId that are implicitly subscribed grains. They are expected to support the streaming consumer extension.</returns>
/// <returns>A dictionary mapping subscription IDs to implicitly subscribed grain IDs. The grains are expected to support the streaming consumer extension.</returns>
/// <exception cref="System.ArgumentException">The stream ID doesn't have an associated namespace.</exception>
/// <exception cref="System.InvalidOperationException">Internal invariant violation.</exception>
internal Dictionary<Guid, GrainId> GetImplicitSubscribers(QualifiedStreamId streamId, IInternalGrainFactory grainFactory)
internal Dictionary<Guid, GrainId> GetImplicitSubscribers(QualifiedStreamId streamId)
{
var streamNamespace = streamId.GetNamespace();
if (!IsImplicitSubscribeEligibleNameSpace(streamNamespace))
Expand Down
Loading