Repository navigation
feat(streaming): add subscription start position - #10936
Conversation
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Review tier: Lite
Findings: 1
New issues introduced by this change (1)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — CreateStartPositionToken returns a StartPositionToken instance without initializing… |
What changed in this PR
Adds a per-subscription “start position” concept to Orleans persistent streaming so tokenless rewindable subscriptions can opt into starting at the earliest message retained in the local queue cache (or wait for the first future message when none is retained), while keeping the existing “latest” default behavior.
Changes:
- Introduces
StreamSubscriptionOptions/StreamSubscriptionStartPositionand newSubscribeWithOptionsAsyncAPIs for item and batch observables. - Extends the consumer handshake and
PersistentStreamPullingAgentto negotiate and honor “earliest available” starts. - Adds cache support (
GetCacheCursorAtPosition) for pooled/simple caches (and Event Hubs wrapper) plus comprehensive tests.
| File | Description |
|---|---|
| test/Orleans.Streaming.Tests/StreamingTests/StreamSubscriptionHandleImplTests.cs | Adds unit tests covering new options/handshake token construction and default-interface behavior. |
| test/Orleans.Streaming.Tests/StreamingTests/PersistentStreamPullingAgentTests.cs | Adds end-to-end handshake/cursor behavior tests for latest vs earliest and error paths. |
| test/Orleans.Streaming.Tests/OrleansRuntime/Streams/SimpleQueueCacheTests.cs | New tests for earliest-available cursor behavior in SimpleQueueCache. |
| test/Orleans.Streaming.Tests/OrleansRuntime/Streams/PooledQueueCacheTests.cs | Adds pooled cache tests for earliest-available positioning and purge/anchor scenarios. |
| test/Grains/TestInternalGrains/StreamLifecycleTestInternalGrains.cs | Updates internal test grain calls to pass new options parameter. |
| test/Grains/TestInternalGrains/StreamLifecycleTestGrains.cs | Updates internal test grain calls to pass new options parameter. |
| src/Orleans.Streaming/QueueAdapters/IQueueCache.cs | Adds default-interface GetCacheCursorAtPosition for start-position cursor acquisition. |
| src/Orleans.Streaming/PersistentStreams/PersistentStreamPullingAgent.cs | Propagates start-position through handshake; adds faulting behavior for unsupported caches/tokens. |
| src/Orleans.Streaming/MemoryStreams/MemoryPooledCache.cs | Implements GetCacheCursorAtPosition and fixes cursor refresh forwarding for pooled cache wrapper. |
| src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs | Extends handle construction to choose start-position token vs explicit start token. |
| src/Orleans.Streaming/Internal/StreamImpl.cs | Adds SubscribeWithOptionsAsync pass-through APIs for item and batch subscriptions. |
| src/Orleans.Streaming/Internal/StreamHandshakeToken.cs | Adds StartPositionToken and token factory for start position. |
| src/Orleans.Streaming/Internal/StreamConsumerExtension.cs | Validates options and restricts non-latest positions for stateless workers/implicit subscriptions. |
| src/Orleans.Streaming/Internal/StreamConsumer.cs | Adds SubscribeWithOptionsAsync overloads and validates option/token combinations. |
| src/Orleans.Streaming/Generator/GeneratorPooledCache.cs | Implements GetCacheCursorAtPosition and fixes cursor refresh forwarding for pooled cache wrapper. |
| src/Orleans.Streaming/Extensions/AsyncObservableExtensions.cs | Adds delegate-based SubscribeWithOptionsAsync extension overloads. |
| src/Orleans.Streaming/Extensions/AsyncBatchObservableExtensions.cs | Adds delegate-based SubscribeWithOptionsAsync extension overloads for batch observers. |
| src/Orleans.Streaming/Core/StreamSubscriptionOptions.cs | New public options + enum defining Latest vs EarliestAvailable. |
| src/Orleans.Streaming/Core/IAsyncObservable.cs | Adds default-interface SubscribeWithOptionsAsync with fallback/NotSupported behavior. |
| src/Orleans.Streaming/Core/IAsyncBatchObservable.cs | Adds default-interface SubscribeWithOptionsAsync with fallback/NotSupported behavior. |
| src/Orleans.Streaming/Common/SimpleCache/SimpleQueueCacheCursor.cs | Adds cursor state to support earliest-available waiting behavior. |
| src/Orleans.Streaming/Common/SimpleCache/SimpleQueueCache.cs | Implements earliest-available cursor positioning and refresh behavior in simple cache. |
| src/Orleans.Streaming/Common/PooledCache/PooledQueueCache.cs | Implements earliest-available cursor positioning and insertion-order tracking in pooled cache. |
| src/Orleans.Streaming/Common/PooledCache/CachedMessageBlock.cs | Adds generation tracking/metadata needed for pooled earliest-available cursor recovery. |
| src/Azure/Orleans.Streaming.EventHubs/Providers/Streams/EventHub/IEventHubQueueCache.cs | Adds default-interface GetCursorAtPosition for Event Hubs cache abstraction. |
| src/Azure/Orleans.Streaming.EventHubs/Providers/Streams/EventHub/EventHubQueueCache.cs | Implements GetCursorAtPosition forwarding to underlying pooled cache. |
| src/Azure/Orleans.Streaming.EventHubs/Providers/Streams/EventHub/EventHubAdapterReceiver.cs | Implements GetCacheCursorAtPosition in adapter receiver cursor wrapper. |
| src/api/Orleans.Streaming/Orleans.Streaming.cs | Updates public API surface for new options and cursor-position APIs. |
| src/api/Azure/Orleans.Streaming.EventHubs/Orleans.Streaming.EventHubs.cs | Updates public API surface for Event Hubs cursor-position API. |
Suppressed comments (1)
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs:32
- Method name
CreateDeliveyTokenappears to be a typo (“Delivey”). Since this helper is referenced in multiple places (including new tests), consider renaming it toCreateDeliveryTokenand keeping a temporary forwarding method if needed to avoid churn during transition.
public static StreamHandshakeToken? CreateDeliveyToken(StreamSequenceToken? token)
{
if (token == null) return default;
return new DeliveryToken {Token = token};
}
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Review tier: Lite
Findings: 1
New issues introduced by this change (1)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamConsumer.cs — ArgumentNullException is incorrect here because token is known to be non-null (the condition is… |
Pre-existing issues (1)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — CreateStartPositionToken returns a StartPositionToken instance without initializing… View comment |
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs:32
- Method name
CreateDeliveyTokencontains a spelling error ("Delivey"). Since this helper is used widely, a low-risk way to prevent further spread is to add a correctly-spelledCreateDeliveryTokenoverload and have the existing method delegate to it.
public static StreamHandshakeToken? CreateDeliveyToken(StreamSequenceToken? token)
{
if (token == null) return default;
return new DeliveryToken {Token = token};
}
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
Review tier: Lite
Findings: 1
Pre-existing issues (2)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — CreateStartPositionToken returns a StartPositionToken instance without initializing… View comment |
|
src/Orleans.Streaming/Internal/StreamConsumer.cs — ArgumentNullException is incorrect here because token is known to be non-null (the condition is… View comment |
Suppressed comments (3)
Previously missed (2) — in code that hasn't changed since the last review.
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs:32
CreateDeliveyTokenappears to be a misspelling ("Delivey" vs "Delivery"). Since this is an internal API and you’re adding new start-position-related APIs now, consider renaming it toCreateDeliveryTokenand updating the handful of call sites to avoid perpetuating the typo.
public static StreamHandshakeToken? CreateDeliveyToken(StreamSequenceToken? token)
{
if (token == null) return default;
return new DeliveryToken {Token = token};
}
src/Orleans.Streaming/Core/StreamSubscriptionStartPosition.cs:5
- The file uses block-scoped namespace syntax, but this repo generally uses file-scoped namespaces for C# source. Updating to file-scoped keeps the style consistent and avoids extra indentation in new public API files.
namespace Orleans.Streams
{
/// <summary>
src/Orleans.Streaming/Internal/StreamConsumer.cs:102
tokenis non-null here, so throwingArgumentNullExceptionis misleading. UseArgumentException(orInvalidOperationException) to reflect an invalid argument value when subscribing to a non-rewindable observable.
startPosition.Validate();
if (token != null && !IsRewindable)
throw new ArgumentNullException(nameof(token), "Passing a non-null token to a non-rewindable IAsyncObservable.");
if (startPosition != StreamSubscriptionStartPosition.Latest && !IsRewindable)
throw new InvalidOperationException("A non-latest start position requires a rewindable IAsyncObservable.");
|
CI is failing consistently in streaming resume/recovery paths introduced by this change: MemoryStreamResumeTests.ResumeAfterDeactivation*, EHStreamingResumeTests.ResumeAfterDeactivation*, StatelessWorkersStreamTests.SubscribeToStream_FromStatelessWorkerWithSequenceToken_Fails, and GeneratedImplicitSubscriptionStreamRecoveryTests.Recoverable100EventStreamsWithTransientErrorsTest fail across Redis, Event Hubs, Functional, and SlowBVT matrices on Windows, Linux, and macOS. This pattern appears specific to the new subscription start-position behavior rather than an isolated infrastructure flake. Run: https://github.com/dotnet/orleans/actions/runs/33477112547 |
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
Review tier: Lite
Findings: None
Issues resolved since last review (2)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamConsumer.cs — ArgumentNullException is incorrect here because token is known to be non-null (the condition is… View resolved comment |
|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — CreateStartPositionToken returns a StartPositionToken instance without initializing… View resolved comment |
Suppressed comments (2)
Previously missed (2) — in code that hasn't changed since the last review.
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs:32
- The helper is named
CreateDeliveyToken, which appears to be a typo. Since this is an internal API, consider adding a correctly-spelledCreateDeliveryTokenmethod and forwarding the existing misspelled method to it to avoid perpetuating the typo in new call sites.
public static StreamHandshakeToken? CreateDeliveyToken(StreamSequenceToken? token)
{
if (token == null) return default;
return new DeliveryToken {Token = token};
}
src/Orleans.Streaming/Internal/StreamConsumer.cs:133
- The local variable name
subriptionHandleis misspelled, which makes the code harder to read/search. Rename it tosubscriptionHandleand update the return accordingly.
var subriptionHandle = myExtension!.SetObserver(subscriptionId, stream, observer, batchObserver, token, startPosition, filterData);
try
{
await pubSub.RegisterConsumer(subscriptionId, stream.InternalStreamId, myGrainReference!.GetGrainId(), filterData);
return subriptionHandle;
659b064 to
8694c51
Compare
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
Review tier: Lite
Findings: None
Suppressed comments (2)
Previously missed (2) — in code that hasn't changed since the last review.
src/Orleans.Streaming/Internal/StreamConsumer.cs:133
- The local variable name
subriptionHandleis misspelled. This makes the code harder to read/search and propagates the typo into new edits in this method.
Rename it to subscriptionHandle (and update the return).
var subriptionHandle = myExtension!.SetObserver(subscriptionId, stream, observer, batchObserver, token, startPosition, filterData);
try
{
await pubSub.RegisterConsumer(subscriptionId, stream.InternalStreamId, myGrainReference!.GetGrainId(), filterData);
return subriptionHandle;
src/Orleans.Streaming/PersistentStreams/PersistentStreamPullingAgent.cs:396
DoHandshakeWithConsumertreats any non-null handshake token whoseTokenis null as an unsupported token type and faults the subscription. ForStartToken/DeliveryToken, this is a supported token type but an invalid state, so the resulting exception is misleading and makes diagnosing handshake issues harder.
Handle StartToken/DeliveryToken explicitly and throw a clear exception when the sequence token is missing (consistent with the delivery-handshake path).
else if (requestedHandshakeToken is not null)
{
forceFaultSubscription = true;
throw new InvalidOperationException($"Unsupported stream handshake token type {requestedHandshakeToken.GetType().FullName}.");
}
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Review tier: Lite
Findings: 1
New issues introduced by this change (1)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — StartPositionToken declares its startPosition field as [Id(0)], which collides with the base… |
Suppressed comments (1)
src/Orleans.Streaming/Internal/StreamConsumer.cs:134
- Local variable name
subriptionHandleis misspelled, which makes the code harder to read/search. Rename it tosubscriptionHandle(and update the return statement) to avoid propagating the typo.
var subriptionHandle = myExtension!.SetObserver(subscriptionId, stream, observer, batchObserver, token, startPosition, filterData);
try
{
await pubSub.RegisterConsumer(subscriptionId, stream.InternalStreamId, myGrainReference!.GetGrainId(), filterData);
return subriptionHandle;
}
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
StartPositionToken introduces a serializer field id collision with its base type ([Id(0)]), which is likely to break codegen/build or serialization correctness.
Review tier: Lite
Findings: 1
Pre-existing issues (1)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — StartPositionToken declares its startPosition field as [Id(0)], which collides with the base… View comment |
Suppressed comments (1)
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs:62
- StartPositionToken defines a serialized field with
[Id(0)], but the base type StreamHandshakeToken already uses[Id(0)]forToken. Orleans generated serializers require unique ids across an inheritance hierarchy; this id collision will either fail codegen/build or result in incorrect serialization/deserialization of StartPositionToken instances.
internal sealed class StartPositionToken : StreamHandshakeToken
{
[Id(0)]
private StreamSubscriptionStartPosition? startPosition;
bde3fe5 to
8bf249a
Compare
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
Duplicate Orleans serialization member ids in StartPositionToken vs its base StreamHandshakeToken can break serialization/codegen and should be fixed before merge.
Review tier: Lite
Findings: 1
Pre-existing issues (1)
| Severity | Finding |
|---|---|
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs — StartPositionToken declares its startPosition field as [Id(0)], which collides with the base… View comment |
Suppressed comments (1)
src/Orleans.Streaming/Internal/StreamHandshakeToken.cs:62
StartPositionTokendeclaresstartPositionwith[Id(0)], butStreamHandshakeTokenalready definesTokenwith[Id(0)]. Orleans serializer member ids must be unique across the inheritance hierarchy; this duplicate id can cause serialization/codegen failures or incorrect field binding at runtime.
internal sealed class StartPositionToken : StreamHandshakeToken
{
[Id(0)]
private StreamSubscriptionStartPosition? startPosition;
97ad38f to
6d5430f
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
A newly added pulling-agent test is missing the standard BVT/Streaming category attributes and will be skipped by filtered CI test runs.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Review tier: Lite
Findings: 1
New issues introduced by this change (1)
| Severity | Finding |
|---|---|
test/Orleans.Streaming.Tests/StreamingTests/PersistentStreamPullingAgentTests.cs — This test is the only one in this class without the standard BVT/Streaming categorization… |
Code coverage
Report-only conclusion: current-main baseline stale. The newest successful coverage run tested 3604a08, not current main e842631. Coverage combines every CI test matrix job, including providers, CodeGen, .NET 8/10, Linux, Windows, and macOS, using canonical physical source and branch identities. The comparison remains report-only while normal line and branch variance is calibrated. Coverage details |
6d5430f to
6cdec4e
Compare
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
It changes core persistent-stream pulling/handshake/cursoring behavior across multiple cache implementations, so a final human review is needed to validate compatibility and regression risk beyond the added tests.
Review tier: Lite
Findings: None
Issues resolved since last review (1)
| Severity | Finding |
|---|---|
test/Orleans.Streaming.Tests/StreamingTests/PersistentStreamPullingAgentTests.cs — This test is the only one in this class without the standard BVT/Streaming categorization… View resolved comment |
There was a problem hiding this comment.
Copilot review overview
🔵 Needs a closer look
There are a few concrete correctness/behavior issues to address (notably handshake suppression via disableHandshake being bypassed for explicit start positions, plus smaller API/clarity fixes) before it’s safe to approve.
Review tier: Lite
Findings: None
Suppressed comments (3)
Previously missed (3) — in code that hasn't changed since the last review.
src/Orleans.Streaming/Internal/StreamConsumer.cs:102
SubscribeAsyncImpltreats an invalidstartPositionvalue (non-latest on a non-rewindable stream) as an argument validation error, but it throwsInvalidOperationExceptioninstead of an argument exception. This is inconsistent with the adjacenttokenvalidation (which throwsArgumentException) and makes it harder for callers to diagnose parameter issues.
src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs:96disableHandshakeis intended to suppress handshake behavior (used for stateless worker subscriptions), but the constructor still initializesexpectedTokenfromstartPositioneven whendisableHandshakeis true. That re-enables handshake enforcement inDeliverBatch/DeliverItemfor cases which explicitly passLatest, potentially causing unnecessary renegotiation or delivery failures for stateless-worker subscriptions.
src/Orleans.Streaming/Internal/StreamConsumer.cs:129- Local variable
subriptionHandleis misspelled; this makes the code harder to search/read and the line was touched in this PR. Renaming tosubscriptionHandlekeeps the intent clear.



Problem
Rewindable persistent-stream subscribers can provide a sequence token, but tokenless subscriptions begin with live delivery. Applications need both per-subscription control and a provider-wide default for consuming messages which are still retained in the local queue cache.
Solution
Add
StreamSubscriptionStartPositionwithLatestandEarliestAvailable.Callers can select a position through
SubscribeAsyncoverloads for item and batch observers. Providers can configure the default for legacy tokenless subscriptions throughStreamPullingAgentOptions.InitialSubscriptionStartPosition.Latestremains the default.Start selection follows one precedence rule: a concrete sequence token, then an explicit subscription position, then the provider default, then
Latest. Both explicit and provider-selected positions use the same consumer-handshake and cache-positioning path.EarliestAvailablebegins inclusively at the oldest retained message for the targetStreamId, or waits for its first future message when none is retained. Built-in pooled, simple, and Event Hubs caches support the position. Custom caches remain compatible through an additive default-interface capability and receive deterministic subscription failure when earliest positioning is unsupported.Rationale
The shared implementation keeps per-subscription and provider-wide behavior consistent without duplicating cache algorithms. Concrete tokens remain inclusive and authoritative, and replay is scoped to the current local cache without changing receiver checkpoints or upstream retention behavior.
During rolling upgrades, use the new start-position controls after the pulling-agent silos are running the updated version.