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
@@ -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
---

Expand Down Expand Up @@ -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 <xref:Orleans.Streams.StreamSequenceToken> represents provider position; it isn't a universal business sequence and not every provider accepts application-supplied tokens.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,11 +103,15 @@ Clients can produce and explicitly consume streams after the provider is configu

## Stateless worker grains

Grains marked with <xref:Orleans.Concurrency.StatelessWorkerAttribute> can publish stream events. Orleans rejects stateless worker grain subscription attempts with an <xref:System.InvalidOperationException> 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 <xref:Orleans.Concurrency.StatelessWorkerAttribute> can publish and consume streams. A stateless worker stream consumer implements <xref:Orleans.Streams.Core.IStreamSubscriptionObserver>. 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 <xref:Orleans.Streams.Core.IStreamSubscriptionObserver> receives an <xref:System.InvalidOperationException>.

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 <xref:System.InvalidOperationException>. Use a regular grain consumer when application-managed rewind or checkpoint resume is required.

<a id="stream-order-and-sequence-tokens"></a>
<a id="rewindable-streams"></a>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ namespace Orleans.Streams.Core
{
/// <summary>
/// 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.
/// </summary>
public interface IStreamSubscriptionObserver
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,8 +43,8 @@ public static void AddSiloStreaming(this IServiceCollection services)
services.AddKeyedTransient<IGrainExtension>(typeof(IStreamConsumerExtension), (sp, _) =>
{
var runtime = sp.GetRequiredService<IStreamProviderRuntime>();
var grainContextAccessor = sp.GetRequiredService<IGrainContextAccessor>();
return new StreamConsumerExtension(runtime, grainContextAccessor.GrainContext?.GrainInstance as IStreamSubscriptionObserver);
var grainContext = sp.GetRequiredService<IGrainContextAccessor>().GrainContext;
return new StreamConsumerExtension(runtime, grainContext);
});
services.AddSingleton<IStreamSubscriptionManagerAdmin>(sp =>
new StreamSubscriptionManagerAdmin(sp.GetRequiredService<IStreamProviderRuntime>()));
Expand Down
6 changes: 5 additions & 1 deletion src/Orleans.Streaming/Internal/StreamConsumer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -242,7 +243,10 @@ private async Task BindExtensionLazy()
if (myExtension == null)
{
LogDebugBindExtensionLazy(providerRuntime);
(myExtension, myGrainReference) = providerRuntime.BindExtension<StreamConsumerExtension, IStreamConsumerExtension>(() => new StreamConsumerExtension(providerRuntime));
(myExtension, myGrainReference) = providerRuntime.BindExtension<StreamConsumerExtension, IStreamConsumerExtension>(
() => new StreamConsumerExtension(
providerRuntime,
providerRuntime.ServiceProvider.GetRequiredService<IGrainContextAccessor>().GrainContext));
LogDebugBindExtension(myExtension, myGrainReference);
}
}
Expand Down
36 changes: 25 additions & 11 deletions src/Orleans.Streaming/Internal/StreamConsumerExtension.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<ILogger<StreamConsumerExtension>>();
}
Expand All @@ -58,13 +61,24 @@ internal StreamSubscriptionHandleImpl<T> SetObserver<T>(
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<T>(subscriptionId, observer, batchObserver, stream, token, filterData);
var handle = new StreamSubscriptionHandleImpl<T>(
subscriptionId,
observer,
batchObserver,
stream,
token,
filterData,
disableHandshake: IsStatelessWorker);
allStreamObservers[subscriptionId] = handle;
return handle;
}
Expand Down Expand Up @@ -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<IStreamProvider>(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))
{
Expand All @@ -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<IStreamProvider>(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))
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,14 +56,15 @@ public StreamSubscriptionHandleImpl(
IAsyncBatchObserver<T>? batchObserver,
StreamImpl<T> 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);
Expand All @@ -79,7 +80,7 @@ public void Invalidate()

public StreamHandshakeToken? GetSequenceToken()
{
return this.expectedToken;
return expectedToken;
}

public override Task UnsubscribeAsync()
Expand Down
15 changes: 12 additions & 3 deletions src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<IGrainExtensionBinder>()! // Grain contexts expose an extension binder.
return grainContext.GetComponent<IGrainExtensionBinder>()! // Grain contexts expose an extension binder.
.GetOrSetExtension<TExtension, TExtensionInterface>(newExtensionFunc);
}

Expand Down
11 changes: 4 additions & 7 deletions src/Orleans.Streaming/StreamConsumerGrainContextAction.cs
Original file line number Diff line number Diff line change
Expand Up @@ -22,15 +22,12 @@ public StreamConsumerGrainContextAction(IStreamProviderRuntime streamProviderRun
/// <inheritdoc/>
public void Configure(IGrainContext context)
{
if (context.GrainInstance is IStreamSubscriptionObserver observer)
if (context.GrainInstance is IStreamSubscriptionObserver)
{
InstallStreamConsumerExtension(context, observer as IStreamSubscriptionObserver);
context.GetComponent<IGrainExtensionBinder>()! // Grain contexts expose an extension binder.
.GetOrSetExtension<StreamConsumerExtension, IStreamConsumerExtension>(
() => new StreamConsumerExtension(_streamProviderRuntime, context));
}
}

private void InstallStreamConsumerExtension(IGrainContext context, IStreamSubscriptionObserver observer)
{
_streamProviderRuntime.BindExtension<StreamConsumerExtension, IStreamConsumerExtension>(() => new StreamConsumerExtension(_streamProviderRuntime, observer));
}
}
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,19 @@
namespace UnitTests.GrainInterfaces
{
public interface IStatelessWorkerStreamConsumerGrain : IGrainWithIntegerKey
{
Task BecomeConsumer(Guid[] streamIds, string providerToUse);

Task BecomeConsumerFromToken(Guid streamId, string providerToUse);

Task<int> StopConsuming(Guid streamId, string providerToUse);
}

public interface IImplicitStatelessWorkerStreamConsumerGrain : IGrainWithGuidKey
{
}

public interface IUnsupportedStatelessWorkerStreamConsumerGrain : IGrainWithIntegerKey
{
Task BecomeConsumer(Guid streamId, string providerToUse);
}
Expand Down
Loading