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..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,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. 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. -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 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 . + +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..c0c26634c91 100644 --- a/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs +++ b/src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs @@ -43,8 +43,8 @@ public static void AddSiloStreaming(this IServiceCollection services) services.AddKeyedTransient(typeof(IStreamConsumerExtension), (sp, _) => { var runtime = sp.GetRequiredService(); - var grainContextAccessor = sp.GetRequiredService(); - return new StreamConsumerExtension(runtime, grainContextAccessor.GrainContext?.GrainInstance as IStreamSubscriptionObserver); + var grainContext = sp.GetRequiredService().GrainContext; + return new StreamConsumerExtension(runtime, grainContext); }); services.AddSingleton(sp => new StreamSubscriptionManagerAdmin(sp.GetRequiredService())); 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/Internal/StreamConsumerExtension.cs b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs index 7f4d45e911b..32940072fe2 100644 --- a/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs +++ b/src/Orleans.Streaming/Internal/StreamConsumerExtension.cs @@ -37,14 +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; + private readonly IGrainContext? _grainContext; - internal StreamConsumerExtension(IStreamProviderRuntime providerRt, IStreamSubscriptionObserver? streamSubscriptionObserver = null) + private IStreamSubscriptionObserver? StreamSubscriptionObserver => + _grainContext?.GrainInstance as IStreamSubscriptionObserver; + + private bool IsStatelessWorker => _grainContext is ActivationData { IsStatelessWorker: true }; + + internal StreamConsumerExtension(IStreamProviderRuntime providerRt, IGrainContext? grainContext = null) { - this.streamSubscriptionObserver = streamSubscriptionObserver; + _grainContext = grainContext; providerRuntime = providerRt; logger = providerRt.ServiceProvider.GetRequiredService>(); } @@ -58,13 +61,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; } @@ -92,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)) { @@ -124,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)) { diff --git a/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs b/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs index 14b427e7aad..cd49d1d5e90 100644 --- a/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs +++ b/src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs @@ -56,14 +56,15 @@ 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; this.batchObserver = batchObserver; this.streamImpl = streamImpl ?? throw new ArgumentNullException(nameof(streamImpl)); this.filterData = filterData; - this.isRewindable = streamImpl.IsRewindable; + this.isRewindable = streamImpl.IsRewindable && !disableHandshake; if (IsRewindable) { expectedToken = StreamHandshakeToken.CreateStartToken(token); @@ -79,7 +80,7 @@ public void Invalidate() public StreamHandshakeToken? GetSequenceToken() { - return this.expectedToken; + return expectedToken; } public override Task UnsubscribeAsync() diff --git a/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs b/src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs index 67a05bfabfe..54ada7b8627 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,12 +150,20 @@ 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; + 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) + { + 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}."); + } } - return this.grainContextAccessor.GrainContext.GetComponent()! // Grain contexts expose an extension binder. + return grainContext.GetComponent()! // Grain contexts expose an extension binder. .GetOrSetExtension(newExtensionFunc); } diff --git a/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs b/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs index 95cab53a124..7b39c333c6c 100644 --- a/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs +++ b/src/Orleans.Streaming/StreamConsumerGrainContextAction.cs @@ -22,15 +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 as IStreamSubscriptionObserver); + context.GetComponent()! // Grain contexts expose an extension binder. + .GetOrSetExtension( + () => new StreamConsumerExtension(_streamProviderRuntime, context)); } } - - private void InstallStreamConsumerExtension(IGrainContext context, IStreamSubscriptionObserver observer) - { - _streamProviderRuntime.BindExtension(() => new StreamConsumerExtension(_streamProviderRuntime, observer)); - } } } diff --git a/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs b/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs index 1379674163e..d7a19204843 100644 --- a/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs @@ -1,6 +1,19 @@ 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 + { + } + + 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..cacaa317a07 100644 --- a/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs +++ b/test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs @@ -1,33 +1,186 @@ -using Microsoft.Extensions.Logging; +using System.Collections.Concurrent; 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 _deliveriesReleased = CreateCompletionSource(); + private TaskCompletionSource _deliveryTargetReached = CreateCompletionSource(); + private bool _blockDeliveries; + private int _deliveryCount; + private int _expectedDeliveries; + private int _waitingDeliveryCount; + + public int DeliveryCount => Volatile.Read(ref _deliveryCount); + + 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(); + Interlocked.Exchange(ref _deliveryCount, 0); + _expectedDeliveries = expectedDeliveries; + Volatile.Write(ref _blockDeliveries, blockDeliveries); + _deliveriesReleased = CreateCompletionSource(); + _deliveryTargetReached = CreateCompletionSource(); + } + + public Task WaitForDeliveriesAsync(TimeSpan timeout) => _deliveryTargetReached.Task.WaitAsync(timeout); + + public Task WaitForReleasedDeliveriesAsync(TimeSpan timeout) => _deliveriesReleased.Task.WaitAsync(timeout); + + 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); + var deliveryCount = Interlocked.Increment(ref _deliveryCount); + if (!Volatile.Read(ref _blockDeliveries)) + { + if (deliveryCount >= _expectedDeliveries) + { + _deliveryTargetReached.TrySetResult(); + } + return; + } + + if (Interlocked.Increment(ref _waitingDeliveryCount) >= _expectedDeliveries) + { + _deliveryTargetReached.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 StatelessWorkerStreamConsumerState _state; - public StatelessWorkerStreamConsumerGrain(ILoggerFactory loggerFactory) + public StatelessWorkerStreamConsumerGrain(StatelessWorkerStreamConsumerState state) { - this.logger = loggerFactory.CreateLogger($"{this.GetType().Name}-{this.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) + { + _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 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..06832226f7c 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,120 @@ 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.WaitForDeliveriesAsync(Timeout); + + Assert.Equal(PartitionCount, state.WaitingDeliveryCount); + Assert.Equal(PartitionCount, state.DeliveryActivationCount); + Assert.Equal(PartitionCount, state.ObserverActivationCount); + + state.ReleaseDeliveries(); + 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(); + for (var attempt = 0; attempt < 1_000 && result.Count < PartitionCount; attempt++) + { + var streamId = Guid.NewGuid(); + var queueId = mapper.GetQueueForStream( + StreamId.Create(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId)); + result.TryAdd(queueId, streamId); + } + + Assert.Equal(PartitionCount, result.Count); + return result.Values.ToArray(); } } } \ No newline at end of file