From 1087c69e61a442fb5b796fa7bdf8cad735549701 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Tue, 18 Aug 2026 15:37:44 -0700 Subject: [PATCH 1/7] feat(streaming): support stateless worker consumers Enable observer-based stream consumption as local competing-consumer execution across stateless worker activations. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../docs/streaming/delivery-semantics.md | 3 +- .../streaming/streams-programming-apis.md | 10 +- .../Core/IStreamSubscriptionObserver.cs | 1 + .../StreamingServiceCollectionExtensions.cs | 6 +- .../Internal/StreamConsumerExtension.cs | 21 +- .../Internal/StreamSubscriptionHandleImpl.cs | 14 +- .../Providers/SiloStreamProviderRuntime.cs | 20 +- .../StreamConsumerGrainContextAction.cs | 7 +- .../IStatelessWorkerStreamConsumerGrain.cs | 14 ++ .../StatelessWorkerStreamConsumerGrain.cs | 186 +++++++++++++++++- .../StatelessWorkersStreamTests.cs | 131 +++++++++--- 11 files changed, 355 insertions(+), 58 deletions(-) diff --git a/docs/site/src/content/docs/streaming/delivery-semantics.md b/docs/site/src/content/docs/streaming/delivery-semantics.md index d6950b9f701..dd3d7d9e857 100644 --- a/docs/site/src/content/docs/streaming/delivery-semantics.md +++ b/docs/site/src/content/docs/streaming/delivery-semantics.md @@ -1,7 +1,7 @@ --- title: Stream delivery, ordering, replay, and recovery description: Design Orleans streaming consumers for provider-specific delivery, ordering, replay, and failure behavior. -ms.date: 08/02/2026 +ms.date: 08/18/2026 ms.topic: concept-article --- @@ -39,6 +39,7 @@ Treat ordering as scoped, not global: - Physical partitions and queues order independently. - Retries, visibility timeouts, consumer failures, and rebalancing can reorder delivery. - A grain processes its own turns serially unless its concurrency configuration says otherwise, but that doesn't impose a total order across grains or subscriptions. +- A stateless worker stream subscription assigns each delivery attempt to one selected activation. Concurrent deliveries and retries can run on different activations, so completion order across the worker pool is unspecified. Use event version numbers or domain sequence numbers when business logic requires ordering. A represents provider position; it isn't a universal business sequence and not every provider accepts application-supplied tokens. diff --git a/docs/site/src/content/docs/streaming/streams-programming-apis.md b/docs/site/src/content/docs/streaming/streams-programming-apis.md index 36ab0b3a65e..e6efcfad8cf 100644 --- a/docs/site/src/content/docs/streaming/streams-programming-apis.md +++ b/docs/site/src/content/docs/streaming/streams-programming-apis.md @@ -103,11 +103,15 @@ Clients can produce and explicitly consume streams after the provider is configu ## Stateless worker grains -Grains marked with can publish stream events. Orleans rejects stateless worker grain subscription attempts with an because a stream consumer uses a grain extension which must bind to one activation, while a stateless worker grain identity can have multiple, replaceable activations. +Grains marked with can publish and consume streams. A stateless worker stream consumer implements . Orleans installs a separate stream consumer extension on every activation and calls `OnSubscribed` when a selected activation first receives a delivery for a subscription. The grain calls `ResumeAsync` from that callback to attach the activation's observer. -Use a regular grain as the stream consumer so that the stream-to-grain binding has a stable virtual identity which can own state and the subscription lifecycle. If processing after delivery is stateless and parallelizable, have that grain call stateless worker grains and await the required work before its consumer task completes. This keeps stream acknowledgment and recovery at the regular grain boundary instead of treating a multicast subscription as a competing-consumer work queue. +The subscription belongs to the stateless worker grain identity. For persistent streams, each pulling agent delivers through that identity from its silo, and normal stateless-worker placement selects one local activation for each delivery attempt. Concurrent pulling agents can therefore process items on different activations and silos. Each delivery attempt runs on one activation, providing competing-consumer execution for stateless transformations such as decoding, validation, enrichment, filtering, and forwarding. -Support for subscribing directly from stateless worker grains is tracked by [dotnet/orleans#433](https://github.com/dotnet/orleans/issues/433). +Implicit subscriptions establish the grain-level subscription from grain metadata. Explicit `SubscribeAsync` calls establish the grain-level subscription at runtime; later activations attach their local observers through `OnSubscribed`, and `UnsubscribeAsync` removes that grain-level subscription. + +Ordering is scoped to the selected activation and provider delivery path. Concurrent deliveries can complete in any order across activations. A provider retry can select a different activation, so handlers use stateless or idempotent processing and follow the provider's delivery guarantee. A stateless worker which subscribes without implementing receives an . + +Stateless worker observers attach with a null sequence token. The pulling agent owns progress for the live subscription as deliveries move between activations. Passing a non-null sequence token to `SubscribeAsync` or `ResumeAsync` produces an . Use a regular grain consumer when application-managed rewind or checkpoint resume is required. diff --git a/src/Orleans.Streaming/Core/IStreamSubscriptionObserver.cs b/src/Orleans.Streaming/Core/IStreamSubscriptionObserver.cs index cc0d96619a3..2eea650ca2e 100644 --- a/src/Orleans.Streaming/Core/IStreamSubscriptionObserver.cs +++ b/src/Orleans.Streaming/Core/IStreamSubscriptionObserver.cs @@ -4,6 +4,7 @@ namespace Orleans.Streams.Core { /// /// When implemented by a grain, notifies the grain of any new or resuming subscriptions. + /// Stateless worker grains receive this notification on each activation so that every activation can attach its own observer. /// public interface IStreamSubscriptionObserver { diff --git a/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs b/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs index 5b5be81f671..6c3d4600f03 100644 --- a/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs +++ b/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs @@ -44,7 +44,11 @@ public static void AddSiloStreaming(this IServiceCollection services) { var runtime = sp.GetRequiredService(); var grainContextAccessor = sp.GetRequiredService(); - return new StreamConsumerExtension(runtime, grainContextAccessor.GrainContext?.GrainInstance as IStreamSubscriptionObserver); + var grainContext = grainContextAccessor.GrainContext; + return new StreamConsumerExtension( + runtime, + grainContext?.GrainInstance as IStreamSubscriptionObserver, + grainContext is ActivationData { IsStatelessWorker: true }); }); services.AddSingleton(sp => new StreamSubscriptionManagerAdmin(sp.GetRequiredService())); diff --git a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs index 7f4d45e911b..cdfb2d91aad 100644 --- a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs +++ b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs @@ -41,10 +41,16 @@ internal sealed partial class StreamConsumerExtension : IStreamConsumerExtension // then this will be not null, otherwise, it will be null [NonSerialized] private readonly IStreamSubscriptionObserver? streamSubscriptionObserver; + [NonSerialized] + private readonly bool isStatelessWorker; - internal StreamConsumerExtension(IStreamProviderRuntime providerRt, IStreamSubscriptionObserver? streamSubscriptionObserver = null) + internal StreamConsumerExtension( + IStreamProviderRuntime providerRt, + IStreamSubscriptionObserver? streamSubscriptionObserver = null, + bool isStatelessWorker = false) { this.streamSubscriptionObserver = streamSubscriptionObserver; + this.isStatelessWorker = isStatelessWorker; providerRuntime = providerRt; logger = providerRt.ServiceProvider.GetRequiredService>(); } @@ -58,13 +64,24 @@ internal StreamSubscriptionHandleImpl SetObserver( string? filterData) { if (null == stream) throw new ArgumentNullException(nameof(stream)); + if (isStatelessWorker && token is not null) + { + throw new InvalidOperationException("Stateless worker stream subscriptions use provider-managed live delivery and require a null sequence token."); + } try { LogDebugAddObserver(providerRuntime.ExecutingEntityIdentity(), stream.InternalStreamId); // Note: The caller [StreamConsumer] already handles locking for Add/Remove operations, so we don't need to repeat here. - var handle = new StreamSubscriptionHandleImpl(subscriptionId, observer, batchObserver, stream, token, filterData); + var handle = new StreamSubscriptionHandleImpl( + subscriptionId, + observer, + batchObserver, + stream, + token, + filterData, + disableHandshake: isStatelessWorker); allStreamObservers[subscriptionId] = handle; return handle; } diff --git a/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs b/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs index 14b427e7aad..6263b11f469 100644 --- a/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs +++ b/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs @@ -23,6 +23,8 @@ internal class StreamSubscriptionHandleImpl : StreamSubscriptionHandle, IS private readonly GuidId subscriptionId; [Id(3)] private readonly bool isRewindable; + [Id(4)] + private readonly bool disableHandshake; [NonSerialized] private IAsyncObserver? observer; @@ -56,7 +58,8 @@ public StreamSubscriptionHandleImpl( IAsyncBatchObserver? batchObserver, StreamImpl streamImpl, StreamSequenceToken? token, - string? filterData) + string? filterData, + bool disableHandshake = false) { this.subscriptionId = subscriptionId ?? throw new ArgumentNullException(nameof(subscriptionId)); this.observer = observer; @@ -64,7 +67,8 @@ public StreamSubscriptionHandleImpl( this.streamImpl = streamImpl ?? throw new ArgumentNullException(nameof(streamImpl)); this.filterData = filterData; this.isRewindable = streamImpl.IsRewindable; - if (IsRewindable) + this.disableHandshake = disableHandshake; + if (IsRewindable && !disableHandshake) { expectedToken = StreamHandshakeToken.CreateStartToken(token); } @@ -79,7 +83,7 @@ public void Invalidate() public StreamHandshakeToken? GetSequenceToken() { - return this.expectedToken; + return disableHandshake ? null : expectedToken; } public override Task UnsubscribeAsync() @@ -146,7 +150,7 @@ public override Task> ResumeAsync(IAsyncBatchObserve } } - if (IsRewindable) + if (IsRewindable && !disableHandshake) { this.expectedToken = StreamHandshakeToken.CreateDeliveyToken(batch.SequenceToken); } @@ -203,7 +207,7 @@ public override Task> ResumeAsync(IAsyncBatchObserve if (!this.expectedToken.Equals(handshakeToken)) return this.expectedToken; } - if (IsRewindable) + if (IsRewindable && !disableHandshake) { this.expectedToken = StreamHandshakeToken.CreateDeliveyToken(currentToken); } diff --git a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs index 67a05bfabfe..5a9fa98c0d6 100644 --- a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs +++ b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs @@ -6,6 +6,7 @@ using Microsoft.Extensions.Logging; using Orleans.Configuration; using Orleans.Streams.Filtering; +using Orleans.Streams.Core; using Orleans.Internal; namespace Orleans.Runtime.Providers @@ -149,13 +150,24 @@ public StreamDirectory GetStreamDirectory() where TExtension : class, TExtensionInterface where TExtensionInterface : class, IGrainExtension { - if (this.grainContextAccessor.GrainContext is ActivationData activationData && activationData.IsStatelessWorker) + var grainContext = this.grainContextAccessor.GrainContext; + var extensionFactory = newExtensionFunc; + if (grainContext is ActivationData { IsStatelessWorker: true } activationData) { - throw new InvalidOperationException($"The extension { typeof(TExtension) } cannot be bound to a Stateless Worker."); + if (typeof(TExtension) != typeof(StreamConsumerExtension) + || typeof(TExtensionInterface) != typeof(IStreamConsumerExtension) + || activationData.GrainInstance is not IStreamSubscriptionObserver observer) + { + throw new InvalidOperationException( + $"The extension {typeof(TExtension)} cannot be bound to stateless worker grain '{activationData.GrainId}'. " + + $"Stateless worker stream consumers must implement {typeof(IStreamSubscriptionObserver).FullName}."); + } + + extensionFactory = () => (TExtension)(object)new StreamConsumerExtension(this, observer, isStatelessWorker: true); } - return this.grainContextAccessor.GrainContext.GetComponent()! // Grain contexts expose an extension binder. - .GetOrSetExtension(newExtensionFunc); + return grainContext.GetComponent()! // Grain contexts expose an extension binder. + .GetOrSetExtension(extensionFactory); } [LoggerMessage( diff --git a/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs b/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs index 95cab53a124..140c78fba7d 100644 --- a/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs +++ b/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs @@ -24,13 +24,16 @@ public void Configure(IGrainContext context) { if (context.GrainInstance is IStreamSubscriptionObserver observer) { - InstallStreamConsumerExtension(context, observer as IStreamSubscriptionObserver); + InstallStreamConsumerExtension(context, observer); } } private void InstallStreamConsumerExtension(IGrainContext context, IStreamSubscriptionObserver observer) { - _streamProviderRuntime.BindExtension(() => new StreamConsumerExtension(_streamProviderRuntime, observer)); + var isStatelessWorker = context is ActivationData { IsStatelessWorker: true }; + context.GetComponent()! // Grain contexts expose an extension binder. + .GetOrSetExtension( + () => new StreamConsumerExtension(_streamProviderRuntime, observer, isStatelessWorker)); } } } diff --git a/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs index 1379674163e..6ed87aea42e 100644 --- a/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs @@ -1,6 +1,20 @@ namespace UnitTests.GrainInterfaces { public interface IStatelessWorkerStreamConsumerGrain : IGrainWithIntegerKey + { + Task BecomeConsumer(Guid[] streamIds, string providerToUse); + + Task BecomeConsumerFromToken(Guid streamId, string providerToUse); + + Task StopConsuming(Guid streamId, string providerToUse); + } + + public interface IImplicitStatelessWorkerStreamConsumerGrain : IGrainWithGuidKey + { + Task Ping(); + } + + public interface IUnsupportedStatelessWorkerStreamConsumerGrain : IGrainWithIntegerKey { Task BecomeConsumer(Guid streamId, string providerToUse); } diff --git a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs index b7f29817e42..9ad96f52ca7 100644 --- a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs @@ -1,33 +1,201 @@ +using System.Collections.Concurrent; +using System.Linq; using Microsoft.Extensions.Logging; using Orleans.Concurrency; +using Orleans.Providers.Streams.Common; using Orleans.Streams; +using Orleans.Streams.Core; using UnitTests.GrainInterfaces; namespace UnitTests.Grains { + public sealed class StatelessWorkerStreamConsumerState + { + private readonly ConcurrentDictionary _deliveryCounts = new(); + private readonly ConcurrentDictionary _observerCounts = new(); + private readonly SemaphoreSlim _deliverySemaphore = new(0); + private TaskCompletionSource _blockedDeliveriesReached = CreateCompletionSource(); + private TaskCompletionSource _deliveriesReleased = CreateCompletionSource(); + private TaskCompletionSource _deliveryTargetReached = CreateCompletionSource(); + private int _blockDeliveries; + private int _expectedDeliveries; + private int _waitingDeliveryCount; + + public int DeliveryCount => _deliveryCounts.Values.Sum(); + + public int DeliveryActivationCount => _deliveryCounts.Count; + + public int ObserverActivationCount => _observerCounts.Count; + + public int WaitingDeliveryCount => Volatile.Read(ref _waitingDeliveryCount); + + public void Reset(int expectedDeliveries, bool blockDeliveries = false) + { + if (WaitingDeliveryCount != 0) + { + throw new InvalidOperationException("All blocked stream deliveries must be released before resetting the test state."); + } + + while (_deliverySemaphore.Wait(0)) + { + } + + _deliveryCounts.Clear(); + _observerCounts.Clear(); + _expectedDeliveries = expectedDeliveries; + Volatile.Write(ref _blockDeliveries, blockDeliveries ? 1 : 0); + _blockedDeliveriesReached = CreateCompletionSource(); + _deliveriesReleased = CreateCompletionSource(); + _deliveryTargetReached = CreateCompletionSource(); + } + + public Task WaitForDeliveriesAsync(TimeSpan timeout) => _deliveryTargetReached.Task.WaitAsync(timeout); + + public Task WaitForBlockedDeliveriesAsync(TimeSpan timeout) => _blockedDeliveriesReached.Task.WaitAsync(timeout); + + public Task WaitForReleasedDeliveriesAsync(TimeSpan timeout) => _deliveriesReleased.Task.WaitAsync(timeout); + + public void ReleaseDeliveries(int count) => _deliverySemaphore.Release(count); + + internal void RecordObserver(Guid activationId) => _observerCounts.AddOrUpdate(activationId, 1, static (_, count) => count + 1); + + internal async Task RecordDelivery(Guid activationId) + { + _deliveryCounts.AddOrUpdate(activationId, 1, static (_, count) => count + 1); + if (_deliveryCounts.Values.Sum() >= _expectedDeliveries) + { + _deliveryTargetReached.TrySetResult(); + } + + if (Volatile.Read(ref _blockDeliveries) == 0) + { + return; + } + + if (Interlocked.Increment(ref _waitingDeliveryCount) >= _expectedDeliveries) + { + _blockedDeliveriesReached.TrySetResult(); + } + + try + { + await _deliverySemaphore.WaitAsync(); + } + finally + { + if (Interlocked.Decrement(ref _waitingDeliveryCount) == 0) + { + _deliveriesReleased.TrySetResult(); + } + } + } + + private static TaskCompletionSource CreateCompletionSource() => + new(TaskCreationOptions.RunContinuationsAsynchronously); + } + [StatelessWorker(MaxLocalWorkers)] - public class StatelessWorkerStreamConsumerGrain : Grain, IStatelessWorkerStreamConsumerGrain + public class StatelessWorkerStreamConsumerGrain + : Grain, IStatelessWorkerStreamConsumerGrain, IStreamSubscriptionObserver, IAsyncObserver { - internal const int MaxLocalWorkers = 1; - internal const string StreamNamespace = "StatelessWorkerStreamingNamespace"; + public const int MaxLocalWorkers = 4; + public const string ExplicitStreamNamespace = "StatelessWorkerStreamingNamespace"; - private readonly ILogger logger; + private readonly Guid _activationId = Guid.NewGuid(); + private readonly ILogger _logger; + private readonly StatelessWorkerStreamConsumerState _state; - public StatelessWorkerStreamConsumerGrain(ILoggerFactory loggerFactory) + public StatelessWorkerStreamConsumerGrain( + ILoggerFactory loggerFactory, + StatelessWorkerStreamConsumerState state) { - this.logger = loggerFactory.CreateLogger($"{this.GetType().Name}-{this.IdentityString}"); + _logger = loggerFactory.CreateLogger($"{GetType().Name}-{IdentityString}"); + _state = state; } public Task OnCompletedAsync() => Task.CompletedTask; public Task OnErrorAsync(Exception ex) => Task.CompletedTask; - public Task OnNextAsync(string item, StreamSequenceToken? token = null) => Task.CompletedTask; + public Task OnNextAsync(string item, StreamSequenceToken? token = null) => _state.RecordDelivery(_activationId); + public async Task BecomeConsumer(Guid[] streamIds, string providerToUse) + { + foreach (var streamId in streamIds) + { + var stream = this.GetStreamProvider(providerToUse).GetStream(ExplicitStreamNamespace, streamId); + _ = await stream.SubscribeAsync(OnNextAsync, OnErrorAsync, OnCompletedAsync); + _state.RecordObserver(_activationId); + } + } + + public async Task BecomeConsumerFromToken(Guid streamId, string providerToUse) + { + var stream = this.GetStreamProvider(providerToUse).GetStream(ExplicitStreamNamespace, streamId); + _ = await stream.SubscribeAsync(this, new EventSequenceToken(0)); + } + + public async Task StopConsuming(Guid streamId, string providerToUse) + { + var stream = this.GetStreamProvider(providerToUse).GetStream(ExplicitStreamNamespace, streamId); + var handles = await stream.GetAllSubscriptionHandles(); + foreach (var handle in handles) + { + await handle.UnsubscribeAsync(); + } + + return handles.Count; + } + + public async Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) + { + _logger.LogInformation( + "Attaching activation {ActivationId} to stream {ProviderName}/{StreamId}", + _activationId, + handleFactory.ProviderName, + handleFactory.StreamId); + _state.RecordObserver(_activationId); + await handleFactory.Create().ResumeAsync(OnNextAsync, OnErrorAsync, OnCompletedAsync); + } + } + + [StatelessWorker(MaxLocalWorkers)] + [ImplicitStreamSubscription(StreamNamespace)] + public class ImplicitStatelessWorkerStreamConsumerGrain + : Grain, IImplicitStatelessWorkerStreamConsumerGrain, IStreamSubscriptionObserver + { + public const int MaxLocalWorkers = 4; + public const string StreamNamespace = "ImplicitStatelessWorkerStreamingNamespace"; + + private readonly Guid _activationId = Guid.NewGuid(); + private readonly StatelessWorkerStreamConsumerState _state; + + public ImplicitStatelessWorkerStreamConsumerGrain(StatelessWorkerStreamConsumerState state) + { + _state = state; + } + + public Task Ping() => Task.CompletedTask; + + public async Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) + { + _state.RecordObserver(_activationId); + await handleFactory.Create().ResumeAsync( + (item, token) => _state.RecordDelivery(_activationId), + static exception => Task.CompletedTask, + static () => Task.CompletedTask); + } + } + + [StatelessWorker] + public class UnsupportedStatelessWorkerStreamConsumerGrain + : Grain, IUnsupportedStatelessWorkerStreamConsumerGrain + { public async Task BecomeConsumer(Guid streamId, string providerToUse) { - var stream = this.GetStreamProvider(providerToUse).GetStream(StreamNamespace, streamId); - _ = await stream.SubscribeAsync(OnNextAsync, OnErrorAsync, OnCompletedAsync); + var stream = this.GetStreamProvider(providerToUse) + .GetStream(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId); + _ = await stream.SubscribeAsync(static (item, token) => Task.CompletedTask); } } } diff --git a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs index 7b9c04f9092..59bc046398a 100644 --- a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs +++ b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs @@ -1,9 +1,13 @@ +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Configuration; -using Microsoft.Extensions.Logging; +using Orleans.Configuration; using Orleans.Providers; +using Orleans.Runtime; +using Orleans.Streams; using Orleans.TestingHost; using TestExtensions; using UnitTests.GrainInterfaces; +using UnitTests.Grains; using Xunit; namespace UnitTests.StreamingTests @@ -19,12 +23,11 @@ public class StatelessWorkersStreamTests : OrleansTestingBase, IClassFixture(); builder.AddClientBuilderConfigurator(); } @@ -33,8 +36,11 @@ public class SiloConfigurator : ISiloConfigurator { public void Configure(ISiloBuilder hostBuilder) { - hostBuilder.AddMemoryStreams(StreamProvider) - .AddMemoryGrainStorage("PubSubStore"); + hostBuilder.AddMemoryStreams( + StreamProvider, + streams => streams.ConfigurePartitioning(PartitionCount)) + .AddMemoryGrainStorage("PubSubStore"); + hostBuilder.Services.AddSingleton(); } } @@ -42,56 +48,119 @@ public class ClientConfiguretor : IClientBuilderConfigurator { public void Configure(IConfiguration configuration, IClientBuilder clientBuilder) { - clientBuilder.AddMemoryStreams(StreamProvider); + clientBuilder.AddMemoryStreams( + StreamProvider, + streams => streams.ConfigurePartitioning(PartitionCount)); } } } + private static readonly TimeSpan Timeout = TimeSpan.FromSeconds(30); + private const int PartitionCount = 4; private const string StreamProvider = StreamTestsConstants.MEMORY_STREAM_PROVIDER_NAME; public StatelessWorkersStreamTests(Fixture fixture) { this.fixture = fixture; - logger = this.fixture.Logger; } [Fact, TestCategory("Functional")] - public async Task SubscribeToStream_FromStatelessWorker_Fail() + public async Task ExplicitSubscription_DeliversToStatelessWorker_AndCanBeRemovedAtGrainScope() { - this.logger.LogInformation($"************************ { nameof(SubscribeToStream_FromStatelessWorker_Fail) } *********************************"); - var runner = new StatelessWorkersStreamTestsRunner(StreamProvider, this.logger, this.fixture.HostedCluster); - await Assert.ThrowsAsync( () => runner.BecomeConsumer(Guid.NewGuid())); + var state = GetConsumerState(); + state.Reset(expectedDeliveries: 1); + var streamId = Guid.NewGuid(); + var consumer = fixture.GrainFactory.GetGrain(0); + + await consumer.BecomeConsumer([streamId], StreamProvider); + await GetStream(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId).OnNextAsync("first"); + await state.WaitForDeliveriesAsync(Timeout); + + Assert.Equal(1, state.DeliveryCount); + Assert.Equal(1, await consumer.StopConsuming(streamId, StreamProvider)); + Assert.Equal(0, await consumer.StopConsuming(streamId, StreamProvider)); } - } - /// - /// Test runner class for executing stateless worker stream tests with producer and consumer functionality. - /// - public class StatelessWorkersStreamTestsRunner - { - private const string StreamNamespace = "SampleStreamNamespace"; + [Fact, TestCategory("Functional")] + public async Task ImplicitSubscription_AttachesObserverBeforeFirstDelivery() + { + var state = GetConsumerState(); + state.Reset(expectedDeliveries: 1); + var streamId = Guid.NewGuid(); + + await GetStream(ImplicitStatelessWorkerStreamConsumerGrain.StreamNamespace, streamId).OnNextAsync("first"); + await state.WaitForDeliveriesAsync(Timeout); - private readonly string streamProvider; - private readonly ILogger logger; - private readonly TestCluster cluster; + Assert.Equal(1, state.DeliveryCount); + Assert.Equal(1, state.ObserverActivationCount); + } - public StatelessWorkersStreamTestsRunner(string streamProvider, ILogger logger, TestCluster cluster) + [Fact, TestCategory("Functional")] + public async Task ConcurrentQueueDeliveries_UseMultipleLocalWorkerActivations() { - this.streamProvider = streamProvider; - this.logger = logger; - this.cluster = cluster; + var state = GetConsumerState(); + state.Reset(expectedDeliveries: PartitionCount, blockDeliveries: true); + var streamIds = CreateStreamIdsForDistinctQueues(); + var consumer = fixture.GrainFactory.GetGrain(1); + await consumer.BecomeConsumer(streamIds, StreamProvider); + + await Task.WhenAll(streamIds.Select((streamId, index) => + GetStream(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId) + .OnNextAsync($"item-{index}"))); + await state.WaitForBlockedDeliveriesAsync(Timeout); + + Assert.Equal(PartitionCount, state.WaitingDeliveryCount); + Assert.Equal(PartitionCount, state.DeliveryActivationCount); + Assert.Equal(PartitionCount, state.ObserverActivationCount); + + state.ReleaseDeliveries(PartitionCount); + await state.WaitForReleasedDeliveriesAsync(Timeout); + Assert.Equal(PartitionCount, state.DeliveryCount); } - public async Task BecomeConsumer(Guid streamId) + [Fact, TestCategory("Functional")] + public async Task SubscribeToStream_FromStatelessWorkerWithoutSubscriptionObserver_Fails() + { + var consumer = fixture.GrainFactory.GetGrain(0); + + var exception = await Assert.ThrowsAsync( + () => consumer.BecomeConsumer(Guid.NewGuid(), StreamProvider)); + + Assert.Contains(typeof(Orleans.Streams.Core.IStreamSubscriptionObserver).FullName!, exception.Message); + } + + [Fact, TestCategory("Functional")] + public async Task SubscribeToStream_FromStatelessWorkerWithSequenceToken_Fails() { - var consumer = this.cluster.GrainFactory!.GetGrain(0); - await consumer.BecomeConsumer(streamId, streamProvider); + var consumer = fixture.GrainFactory.GetGrain(2); + + var exception = await Assert.ThrowsAsync( + () => consumer.BecomeConsumerFromToken(Guid.NewGuid(), StreamProvider)); + + Assert.Contains("null sequence token", exception.Message); } - public async Task ProduceMessage(Guid streamId) + private StatelessWorkerStreamConsumerState GetConsumerState() => + fixture.HostedCluster.GetSiloServiceProvider().GetRequiredService(); + + private IAsyncStream GetStream(string streamNamespace, Guid streamId) => + fixture.Client.GetStreamProvider(StreamProvider).GetStream(streamNamespace, streamId); + + private static Guid[] CreateStreamIdsForDistinctQueues() { - var producer = this.cluster.GrainFactory!.GetGrain(0); - await producer.Produce(streamId, streamProvider, string.Empty); + var mapper = new HashRingBasedStreamQueueMapper( + new HashRingStreamQueueMapperOptions { TotalQueueCount = PartitionCount }, + StreamProvider); + var result = new Dictionary(); + while (result.Count < PartitionCount) + { + var streamId = Guid.NewGuid(); + var queueId = mapper.GetQueueForStream( + StreamId.Create(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId)); + result.TryAdd(queueId, streamId); + } + + return result.Values.ToArray(); } } } \ No newline at end of file From e053bacfb16728e3f158ce3dc26fad020d36ce68 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Wed, 19 Aug 2026 02:34:05 -0700 Subject: [PATCH 2/7] test(streaming): count deliveries atomically Track the concurrent delivery threshold with an interlocked total while retaining per-activation counts. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596 --- .../Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs | 6 ++++-- .../StreamingTests/StatelessWorkersStreamTests.cs | 1 + 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs index 9ad96f52ca7..2d7f2d7d902 100644 --- a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs @@ -18,10 +18,11 @@ public sealed class StatelessWorkerStreamConsumerState private TaskCompletionSource _deliveriesReleased = CreateCompletionSource(); private TaskCompletionSource _deliveryTargetReached = CreateCompletionSource(); private int _blockDeliveries; + private int _deliveryCount; private int _expectedDeliveries; private int _waitingDeliveryCount; - public int DeliveryCount => _deliveryCounts.Values.Sum(); + public int DeliveryCount => Volatile.Read(ref _deliveryCount); public int DeliveryActivationCount => _deliveryCounts.Count; @@ -42,6 +43,7 @@ public void Reset(int expectedDeliveries, bool blockDeliveries = false) _deliveryCounts.Clear(); _observerCounts.Clear(); + Interlocked.Exchange(ref _deliveryCount, 0); _expectedDeliveries = expectedDeliveries; Volatile.Write(ref _blockDeliveries, blockDeliveries ? 1 : 0); _blockedDeliveriesReached = CreateCompletionSource(); @@ -62,7 +64,7 @@ public void Reset(int expectedDeliveries, bool blockDeliveries = false) internal async Task RecordDelivery(Guid activationId) { _deliveryCounts.AddOrUpdate(activationId, 1, static (_, count) => count + 1); - if (_deliveryCounts.Values.Sum() >= _expectedDeliveries) + if (Interlocked.Increment(ref _deliveryCount) >= _expectedDeliveries) { _deliveryTargetReached.TrySetResult(); } diff --git a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs index 59bc046398a..a6d0b6c1eeb 100644 --- a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs +++ b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs @@ -107,6 +107,7 @@ public async Task ConcurrentQueueDeliveries_UseMultipleLocalWorkerActivations() await Task.WhenAll(streamIds.Select((streamId, index) => GetStream(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId) .OnNextAsync($"item-{index}"))); + await state.WaitForDeliveriesAsync(Timeout); await state.WaitForBlockedDeliveriesAsync(Timeout); Assert.Equal(PartitionCount, state.WaitingDeliveryCount); From d3e9467058067fd7f2a6d6ebd9841f5a711c1767 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Thu, 20 Aug 2026 14:59:42 -0700 Subject: [PATCH 3/7] refactor(streaming): simplify stateless worker consumers Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596 --- .../StreamingServiceCollectionExtensions.cs | 8 +--- .../Internal/StreamConsumerExtension.cs | 9 ++--- .../Internal/StreamSubscriptionHandleImpl.cs | 13 +++---- .../Providers/SiloStreamProviderRuntime.cs | 4 +- .../StreamConsumerGrainContextAction.cs | 14 ++----- .../IStatelessWorkerStreamConsumerGrain.cs | 1 - .../StatelessWorkerStreamConsumerGrain.cs | 38 ++++++------------- .../StatelessWorkersStreamTests.cs | 3 +- 8 files changed, 28 insertions(+), 62 deletions(-) diff --git a/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs b/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs index 6c3d4600f03..c0c26634c91 100644 --- a/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs +++ b/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs @@ -43,12 +43,8 @@ public static void AddSiloStreaming(this IServiceCollection services) services.AddKeyedTransient(typeof(IStreamConsumerExtension), (sp, _) => { var runtime = sp.GetRequiredService(); - var grainContextAccessor = sp.GetRequiredService(); - var grainContext = grainContextAccessor.GrainContext; - return new StreamConsumerExtension( - runtime, - grainContext?.GrainInstance as IStreamSubscriptionObserver, - grainContext is ActivationData { IsStatelessWorker: true }); + var grainContext = sp.GetRequiredService().GrainContext; + return new StreamConsumerExtension(runtime, grainContext); }); services.AddSingleton(sp => new StreamSubscriptionManagerAdmin(sp.GetRequiredService())); diff --git a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs index cdfb2d91aad..ba6eda719b5 100644 --- a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs +++ b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs @@ -44,13 +44,10 @@ internal sealed partial class StreamConsumerExtension : IStreamConsumerExtension [NonSerialized] private readonly bool isStatelessWorker; - internal StreamConsumerExtension( - IStreamProviderRuntime providerRt, - IStreamSubscriptionObserver? streamSubscriptionObserver = null, - bool isStatelessWorker = false) + internal StreamConsumerExtension(IStreamProviderRuntime providerRt, IGrainContext? grainContext = null) { - this.streamSubscriptionObserver = streamSubscriptionObserver; - this.isStatelessWorker = isStatelessWorker; + streamSubscriptionObserver = grainContext?.GrainInstance as IStreamSubscriptionObserver; + isStatelessWorker = grainContext is ActivationData { IsStatelessWorker: true }; providerRuntime = providerRt; logger = providerRt.ServiceProvider.GetRequiredService>(); } diff --git a/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs b/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs index 6263b11f469..cd49d1d5e90 100644 --- a/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs +++ b/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs @@ -23,8 +23,6 @@ internal class StreamSubscriptionHandleImpl : StreamSubscriptionHandle, IS private readonly GuidId subscriptionId; [Id(3)] private readonly bool isRewindable; - [Id(4)] - private readonly bool disableHandshake; [NonSerialized] private IAsyncObserver? observer; @@ -66,9 +64,8 @@ public StreamSubscriptionHandleImpl( this.batchObserver = batchObserver; this.streamImpl = streamImpl ?? throw new ArgumentNullException(nameof(streamImpl)); this.filterData = filterData; - this.isRewindable = streamImpl.IsRewindable; - this.disableHandshake = disableHandshake; - if (IsRewindable && !disableHandshake) + this.isRewindable = streamImpl.IsRewindable && !disableHandshake; + if (IsRewindable) { expectedToken = StreamHandshakeToken.CreateStartToken(token); } @@ -83,7 +80,7 @@ public void Invalidate() public StreamHandshakeToken? GetSequenceToken() { - return disableHandshake ? null : expectedToken; + return expectedToken; } public override Task UnsubscribeAsync() @@ -150,7 +147,7 @@ public override Task> ResumeAsync(IAsyncBatchObserve } } - if (IsRewindable && !disableHandshake) + if (IsRewindable) { this.expectedToken = StreamHandshakeToken.CreateDeliveyToken(batch.SequenceToken); } @@ -207,7 +204,7 @@ public override Task> ResumeAsync(IAsyncBatchObserve if (!this.expectedToken.Equals(handshakeToken)) return this.expectedToken; } - if (IsRewindable && !disableHandshake) + if (IsRewindable) { this.expectedToken = StreamHandshakeToken.CreateDeliveyToken(currentToken); } diff --git a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs index 5a9fa98c0d6..e72232b1378 100644 --- a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs +++ b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs @@ -156,14 +156,14 @@ public StreamDirectory GetStreamDirectory() { if (typeof(TExtension) != typeof(StreamConsumerExtension) || typeof(TExtensionInterface) != typeof(IStreamConsumerExtension) - || activationData.GrainInstance is not IStreamSubscriptionObserver observer) + || activationData.GrainInstance is not IStreamSubscriptionObserver) { throw new InvalidOperationException( $"The extension {typeof(TExtension)} cannot be bound to stateless worker grain '{activationData.GrainId}'. " + $"Stateless worker stream consumers must implement {typeof(IStreamSubscriptionObserver).FullName}."); } - extensionFactory = () => (TExtension)(object)new StreamConsumerExtension(this, observer, isStatelessWorker: true); + extensionFactory = () => (TExtension)(object)new StreamConsumerExtension(this, grainContext); } return grainContext.GetComponent()! // Grain contexts expose an extension binder. diff --git a/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs b/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs index 140c78fba7d..7b39c333c6c 100644 --- a/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs +++ b/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs @@ -22,18 +22,12 @@ public StreamConsumerGrainContextAction(IStreamProviderRuntime streamProviderRun /// public void Configure(IGrainContext context) { - if (context.GrainInstance is IStreamSubscriptionObserver observer) + if (context.GrainInstance is IStreamSubscriptionObserver) { - InstallStreamConsumerExtension(context, observer); + context.GetComponent()! // Grain contexts expose an extension binder. + .GetOrSetExtension( + () => new StreamConsumerExtension(_streamProviderRuntime, context)); } } - - private void InstallStreamConsumerExtension(IGrainContext context, IStreamSubscriptionObserver observer) - { - var isStatelessWorker = context is ActivationData { IsStatelessWorker: true }; - context.GetComponent()! // Grain contexts expose an extension binder. - .GetOrSetExtension( - () => new StreamConsumerExtension(_streamProviderRuntime, observer, isStatelessWorker)); - } } } diff --git a/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs index 6ed87aea42e..d7a19204843 100644 --- a/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs @@ -11,7 +11,6 @@ public interface IStatelessWorkerStreamConsumerGrain : IGrainWithIntegerKey public interface IImplicitStatelessWorkerStreamConsumerGrain : IGrainWithGuidKey { - Task Ping(); } public interface IUnsupportedStatelessWorkerStreamConsumerGrain : IGrainWithIntegerKey diff --git a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs index 2d7f2d7d902..60a4041266e 100644 --- a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs @@ -1,6 +1,5 @@ using System.Collections.Concurrent; using System.Linq; -using Microsoft.Extensions.Logging; using Orleans.Concurrency; using Orleans.Providers.Streams.Common; using Orleans.Streams; @@ -14,10 +13,9 @@ public sealed class StatelessWorkerStreamConsumerState private readonly ConcurrentDictionary _deliveryCounts = new(); private readonly ConcurrentDictionary _observerCounts = new(); private readonly SemaphoreSlim _deliverySemaphore = new(0); - private TaskCompletionSource _blockedDeliveriesReached = CreateCompletionSource(); private TaskCompletionSource _deliveriesReleased = CreateCompletionSource(); private TaskCompletionSource _deliveryTargetReached = CreateCompletionSource(); - private int _blockDeliveries; + private bool _blockDeliveries; private int _deliveryCount; private int _expectedDeliveries; private int _waitingDeliveryCount; @@ -45,38 +43,35 @@ public void Reset(int expectedDeliveries, bool blockDeliveries = false) _observerCounts.Clear(); Interlocked.Exchange(ref _deliveryCount, 0); _expectedDeliveries = expectedDeliveries; - Volatile.Write(ref _blockDeliveries, blockDeliveries ? 1 : 0); - _blockedDeliveriesReached = CreateCompletionSource(); + Volatile.Write(ref _blockDeliveries, blockDeliveries); _deliveriesReleased = CreateCompletionSource(); _deliveryTargetReached = CreateCompletionSource(); } public Task WaitForDeliveriesAsync(TimeSpan timeout) => _deliveryTargetReached.Task.WaitAsync(timeout); - public Task WaitForBlockedDeliveriesAsync(TimeSpan timeout) => _blockedDeliveriesReached.Task.WaitAsync(timeout); - public Task WaitForReleasedDeliveriesAsync(TimeSpan timeout) => _deliveriesReleased.Task.WaitAsync(timeout); - public void ReleaseDeliveries(int count) => _deliverySemaphore.Release(count); + public void ReleaseDeliveries() => _deliverySemaphore.Release(_expectedDeliveries); internal void RecordObserver(Guid activationId) => _observerCounts.AddOrUpdate(activationId, 1, static (_, count) => count + 1); internal async Task RecordDelivery(Guid activationId) { _deliveryCounts.AddOrUpdate(activationId, 1, static (_, count) => count + 1); - if (Interlocked.Increment(ref _deliveryCount) >= _expectedDeliveries) - { - _deliveryTargetReached.TrySetResult(); - } - - if (Volatile.Read(ref _blockDeliveries) == 0) + var deliveryCount = Interlocked.Increment(ref _deliveryCount); + if (!Volatile.Read(ref _blockDeliveries)) { + if (deliveryCount >= _expectedDeliveries) + { + _deliveryTargetReached.TrySetResult(); + } return; } if (Interlocked.Increment(ref _waitingDeliveryCount) >= _expectedDeliveries) { - _blockedDeliveriesReached.TrySetResult(); + _deliveryTargetReached.TrySetResult(); } try @@ -104,14 +99,10 @@ public class StatelessWorkerStreamConsumerGrain public const string ExplicitStreamNamespace = "StatelessWorkerStreamingNamespace"; private readonly Guid _activationId = Guid.NewGuid(); - private readonly ILogger _logger; private readonly StatelessWorkerStreamConsumerState _state; - public StatelessWorkerStreamConsumerGrain( - ILoggerFactory loggerFactory, - StatelessWorkerStreamConsumerState state) + public StatelessWorkerStreamConsumerGrain(StatelessWorkerStreamConsumerState state) { - _logger = loggerFactory.CreateLogger($"{GetType().Name}-{IdentityString}"); _state = state; } @@ -151,11 +142,6 @@ public async Task StopConsuming(Guid streamId, string providerToUse) public async Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) { - _logger.LogInformation( - "Attaching activation {ActivationId} to stream {ProviderName}/{StreamId}", - _activationId, - handleFactory.ProviderName, - handleFactory.StreamId); _state.RecordObserver(_activationId); await handleFactory.Create().ResumeAsync(OnNextAsync, OnErrorAsync, OnCompletedAsync); } @@ -177,8 +163,6 @@ public ImplicitStatelessWorkerStreamConsumerGrain(StatelessWorkerStreamConsumerS _state = state; } - public Task Ping() => Task.CompletedTask; - public async Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory) { _state.RecordObserver(_activationId); diff --git a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs index a6d0b6c1eeb..e5e450c1c2a 100644 --- a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs +++ b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs @@ -108,13 +108,12 @@ await Task.WhenAll(streamIds.Select((streamId, index) => GetStream(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId) .OnNextAsync($"item-{index}"))); await state.WaitForDeliveriesAsync(Timeout); - await state.WaitForBlockedDeliveriesAsync(Timeout); Assert.Equal(PartitionCount, state.WaitingDeliveryCount); Assert.Equal(PartitionCount, state.DeliveryActivationCount); Assert.Equal(PartitionCount, state.ObserverActivationCount); - state.ReleaseDeliveries(PartitionCount); + state.ReleaseDeliveries(); await state.WaitForReleasedDeliveriesAsync(Timeout); Assert.Equal(PartitionCount, state.DeliveryCount); } From ab60c282fa4b41e0cd6b72e1938bd9203acda139 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Thu, 20 Aug 2026 15:18:33 -0700 Subject: [PATCH 4/7] refactor(streaming): derive worker consumer context Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596 --- .../Internal/StreamConsumerExtension.cs | 28 +++++++++---------- 1 file changed, 14 insertions(+), 14 deletions(-) diff --git a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs index ba6eda719b5..32940072fe2 100644 --- a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs +++ b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs @@ -37,17 +37,17 @@ internal sealed partial class StreamConsumerExtension : IStreamConsumerExtension [Id(2)] private readonly ILogger logger; private const int MAXIMUM_ITEM_STRING_LOG_LENGTH = 128; - // if this extension is attached to a cosnumer grain which implements IOnSubscriptionActioner, - // then this will be not null, otherwise, it will be null [NonSerialized] - private readonly IStreamSubscriptionObserver? streamSubscriptionObserver; - [NonSerialized] - private readonly bool isStatelessWorker; + private readonly IGrainContext? _grainContext; + + private IStreamSubscriptionObserver? StreamSubscriptionObserver => + _grainContext?.GrainInstance as IStreamSubscriptionObserver; + + private bool IsStatelessWorker => _grainContext is ActivationData { IsStatelessWorker: true }; internal StreamConsumerExtension(IStreamProviderRuntime providerRt, IGrainContext? grainContext = null) { - streamSubscriptionObserver = grainContext?.GrainInstance as IStreamSubscriptionObserver; - isStatelessWorker = grainContext is ActivationData { IsStatelessWorker: true }; + _grainContext = grainContext; providerRuntime = providerRt; logger = providerRt.ServiceProvider.GetRequiredService>(); } @@ -61,7 +61,7 @@ internal StreamSubscriptionHandleImpl SetObserver( string? filterData) { if (null == stream) throw new ArgumentNullException(nameof(stream)); - if (isStatelessWorker && token is not null) + if (IsStatelessWorker && token is not null) { throw new InvalidOperationException("Stateless worker stream subscriptions use provider-managed live delivery and require a null sequence token."); } @@ -78,7 +78,7 @@ internal StreamSubscriptionHandleImpl SetObserver( stream, token, filterData, - disableHandshake: isStatelessWorker); + disableHandshake: IsStatelessWorker); allStreamObservers[subscriptionId] = handle; return handle; } @@ -106,13 +106,13 @@ public bool RemoveObserver(GuidId subscriptionId) { return await observer.DeliverItem(item, currentToken, handshakeToken); } - else if(this.streamSubscriptionObserver != null) + else if (StreamSubscriptionObserver is { } streamSubscriptionObserver) { var streamProvider = this.providerRuntime.ServiceProvider.GetKeyedService(streamId.ProviderName); - if(streamProvider != null) + if (streamProvider != null) { var subscriptionHandlerFactory = new StreamSubscriptionHandlerFactory(streamProvider, streamId, streamId.ProviderName, subscriptionId); - await this.streamSubscriptionObserver.OnSubscribed(subscriptionHandlerFactory); + await streamSubscriptionObserver.OnSubscribed(subscriptionHandlerFactory); //check if an observer were attached after handling the new subscription, deliver on it if attached if (allStreamObservers.TryGetValue(subscriptionId, out observer)) { @@ -138,13 +138,13 @@ public bool RemoveObserver(GuidId subscriptionId) { return await observer.DeliverBatch(batch, handshakeToken); } - else if(this.streamSubscriptionObserver != null) + else if (StreamSubscriptionObserver is { } streamSubscriptionObserver) { var streamProvider = this.providerRuntime.ServiceProvider.GetKeyedService(streamId.ProviderName); if (streamProvider != null) { var subscriptionHandlerFactory = new StreamSubscriptionHandlerFactory(streamProvider, streamId, streamId.ProviderName, subscriptionId); - await this.streamSubscriptionObserver.OnSubscribed(subscriptionHandlerFactory); + await streamSubscriptionObserver.OnSubscribed(subscriptionHandlerFactory); // check if an observer were attached after handling the new subscription, deliver on it if attached if (allStreamObservers.TryGetValue(subscriptionId, out observer)) { From 635e8d283ba07453dc6ccb1e973276c0bfb75179 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Thu, 20 Aug 2026 15:25:39 -0700 Subject: [PATCH 5/7] refactor(streaming): honor extension factory Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596 --- src/Orleans.Streaming/Internal/StreamConsumer.cs | 6 +++++- .../Providers/SiloStreamProviderRuntime.cs | 5 +---- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/src/Orleans.Streaming/Internal/StreamConsumer.cs b/src/Orleans.Streaming/Internal/StreamConsumer.cs index a3ac671e38e..a396367f1b3 100644 --- a/src/Orleans.Streaming/Internal/StreamConsumer.cs +++ b/src/Orleans.Streaming/Internal/StreamConsumer.cs @@ -2,6 +2,7 @@ using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Orleans.Runtime; using Orleans.Streams.Core; @@ -242,7 +243,10 @@ private async Task BindExtensionLazy() if (myExtension == null) { LogDebugBindExtensionLazy(providerRuntime); - (myExtension, myGrainReference) = providerRuntime.BindExtension(() => new StreamConsumerExtension(providerRuntime)); + (myExtension, myGrainReference) = providerRuntime.BindExtension( + () => new StreamConsumerExtension( + providerRuntime, + providerRuntime.ServiceProvider.GetRequiredService().GrainContext)); LogDebugBindExtension(myExtension, myGrainReference); } } diff --git a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs index e72232b1378..54ada7b8627 100644 --- a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs +++ b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs @@ -151,7 +151,6 @@ public StreamDirectory GetStreamDirectory() where TExtensionInterface : class, IGrainExtension { var grainContext = this.grainContextAccessor.GrainContext; - var extensionFactory = newExtensionFunc; if (grainContext is ActivationData { IsStatelessWorker: true } activationData) { if (typeof(TExtension) != typeof(StreamConsumerExtension) @@ -162,12 +161,10 @@ public StreamDirectory GetStreamDirectory() $"The extension {typeof(TExtension)} cannot be bound to stateless worker grain '{activationData.GrainId}'. " + $"Stateless worker stream consumers must implement {typeof(IStreamSubscriptionObserver).FullName}."); } - - extensionFactory = () => (TExtension)(object)new StreamConsumerExtension(this, grainContext); } return grainContext.GetComponent()! // Grain contexts expose an extension binder. - .GetOrSetExtension(extensionFactory); + .GetOrSetExtension(newExtensionFunc); } [LoggerMessage( From b7b5e9c7414c27ef43ab2073a313c6ad6fbe8978 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Fri, 21 Aug 2026 03:12:22 -0700 Subject: [PATCH 6/7] test(streaming): harden stateless worker coverage --- test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs | 1 - .../StreamingTests/StatelessWorkersStreamTests.cs | 3 ++- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs index 60a4041266e..cacaa317a07 100644 --- a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs @@ -1,5 +1,4 @@ using System.Collections.Concurrent; -using System.Linq; using Orleans.Concurrency; using Orleans.Providers.Streams.Common; using Orleans.Streams; diff --git a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs index e5e450c1c2a..06832226f7c 100644 --- a/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs +++ b/test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs @@ -152,7 +152,7 @@ private static Guid[] CreateStreamIdsForDistinctQueues() new HashRingStreamQueueMapperOptions { TotalQueueCount = PartitionCount }, StreamProvider); var result = new Dictionary(); - while (result.Count < PartitionCount) + for (var attempt = 0; attempt < 1_000 && result.Count < PartitionCount; attempt++) { var streamId = Guid.NewGuid(); var queueId = mapper.GetQueueForStream( @@ -160,6 +160,7 @@ private static Guid[] CreateStreamIdsForDistinctQueues() result.TryAdd(queueId, streamId); } + Assert.Equal(PartitionCount, result.Count); return result.Values.ToArray(); } } From 302dfab5e50d7f9d674411bc95664d4f7036ae3e Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Fri, 21 Aug 2026 03:42:27 -0700 Subject: [PATCH 7/7] docs(streaming): clarify observer attachment --- .../src/content/docs/streaming/streams-programming-apis.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/site/src/content/docs/streaming/streams-programming-apis.md b/docs/site/src/content/docs/streaming/streams-programming-apis.md index e6efcfad8cf..68af2c3516a 100644 --- a/docs/site/src/content/docs/streaming/streams-programming-apis.md +++ b/docs/site/src/content/docs/streaming/streams-programming-apis.md @@ -103,11 +103,11 @@ Clients can produce and explicitly consume streams after the provider is configu ## Stateless worker grains -Grains marked with can publish and consume streams. A stateless worker stream consumer implements . Orleans installs a separate stream consumer extension on every activation and calls `OnSubscribed` when a selected activation first receives a delivery for a subscription. The grain calls `ResumeAsync` from that callback to attach the activation's observer. +Grains marked with can publish and consume streams. A stateless worker stream consumer implements . Orleans installs a separate stream consumer extension on every activation. A delivery uses the activation's existing observer for its subscription. When no observer is attached, Orleans calls `OnSubscribed`, and the grain calls `ResumeAsync` from that callback to attach one. The subscription belongs to the stateless worker grain identity. For persistent streams, each pulling agent delivers through that identity from its silo, and normal stateless-worker placement selects one local activation for each delivery attempt. Concurrent pulling agents can therefore process items on different activations and silos. Each delivery attempt runs on one activation, providing competing-consumer execution for stateless transformations such as decoding, validation, enrichment, filtering, and forwarding. -Implicit subscriptions establish the grain-level subscription from grain metadata. Explicit `SubscribeAsync` calls establish the grain-level subscription at runtime; later activations attach their local observers through `OnSubscribed`, and `UnsubscribeAsync` removes that grain-level subscription. +Implicit subscriptions establish the grain-level subscription from grain metadata. Explicit `SubscribeAsync` calls establish the grain-level subscription at runtime and attach the calling activation's observer. A later activation attaches its local observer through `OnSubscribed` when a delivery first reaches it, and `UnsubscribeAsync` removes the grain-level subscription. Ordering is scoped to the selected activation and provider delivery path. Concurrent deliveries can complete in any order across activations. A provider retry can select a different activation, so handlers use stateless or idempotent processing and follow the provider's delivery guarantee. A stateless worker which subscribes without implementing receives an .