feat(streaming): support stateless worker consumers - #10657
Conversation
There was a problem hiding this comment.
Pull request overview
This PR updates Orleans streaming to support stateless worker grains as stream consumers (as competing consumers), gated on implementing IStreamSubscriptionObserver, with stateless-worker-specific semantics around activation-local observers and provider-managed progress.
Changes:
- Allow binding
StreamConsumerExtensionon stateless worker activations only when the grain implementsIStreamSubscriptionObserver, and enforce null sequence tokens for those subscriptions. - Disable rewind/handshake sequence token behavior for stateless worker subscriptions so the pulling agent remains authoritative for live delivery progress.
- Add/expand functional tests plus documentation describing stateless worker delivery locality, ordering, retry, lifecycle, and token rules.
Show a summary per file
| File | Description |
|---|---|
| test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs | Adds functional coverage for explicit/implicit stateless worker consumption, multi-activation concurrency, and invalid contract cases. |
| test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs | Adds observer-based stateless worker test grains plus shared test state to validate activation-local attachment and concurrency. |
| test/Grains/TestGrainInterfaces/IStatelessWorkerStreamConsumerGrain.cs | Updates/extends test grain interfaces for explicit subscriptions, unsubscribe, and negative cases. |
| src/Orleans.Streaming/StreamConsumerGrainContextAction.cs | Adjusts how stream consumer extensions are installed for grains implementing IStreamSubscriptionObserver. |
| src/Orleans.Streaming/Providers/SiloStreamProviderRuntime.cs | Relaxes the stateless-worker binding restriction specifically for StreamConsumerExtension when the observer contract is implemented, otherwise throws. |
| src/Orleans.Streaming/Internal/StreamSubscriptionHandleImpl.cs | Adds a flag to disable handshake/sequence token behavior (used for stateless worker subscriptions). |
| src/Orleans.Streaming/Internal/StreamConsumerExtension.cs | Enforces null sequence tokens for stateless worker subscriptions and passes through handshake-disable behavior. |
| src/Orleans.Streaming/Hosting/StreamingServiceCollectionExtensions.cs | Ensures activation-created StreamConsumerExtension instances know whether they’re on a stateless worker activation. |
| src/Orleans.Streaming/Core/IStreamSubscriptionObserver.cs | Clarifies that stateless worker activations receive OnSubscribed per-activation. |
| docs/site/src/content/docs/streaming/streams-programming-apis.md | Documents stateless worker consumer model, ordering/retry semantics, lifecycle, and token restrictions. |
| docs/site/src/content/docs/streaming/delivery-semantics.md | Adds ordering note for stateless worker competing-consumer processing across activations. |
Review details
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
- Files reviewed: 11/11 changed files
- Comments generated: 1
- Review effort level: Lite
0577c09 to
c2c256e
Compare
There was a problem hiding this comment.
Review details
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
test/Orleans.Streaming.Tests/StreamingTests/StatelessWorkersStreamTests.cs:162
CreateStreamIdsForDistinctQueuescan loop indefinitely ifGuid.NewGuid()keeps hashing to queues which are already present. In practice it’s unlikely, but a rare unlucky streak would hang the test run. Consider bounding the search (and failing with a clear message) so the test can’t deadlock CI.
while (result.Count < PartitionCount)
{
var streamId = Guid.NewGuid();
var queueId = mapper.GetQueueForStream(
StreamId.Create(StatelessWorkerStreamConsumerGrain.ExplicitStreamNamespace, streamId));
result.TryAdd(queueId, streamId);
}
- Files reviewed: 11/11 changed files
- Comments generated: 0 new
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs:3
using System.Linq;appears unused in this file. With code-style analyzers enabled, unused usings (IDE0005) can fail the build; please remove it (or add the missing LINQ usage if intended).
using System.Collections.Concurrent;
using System.Linq;
using Orleans.Concurrency;
- Files reviewed: 11/11 changed files
- Comments generated: 0 new
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs:3
using System.Linq;appears to be unused in this file (no LINQ APIs are referenced). With analyzers/code-style enforcement, unused usings can fail the build (IDE0005/CS8019). Remove the directive to avoid warnings/errors.
using System.Collections.Concurrent;
using System.Linq;
using Orleans.Concurrency;
- Files reviewed: 11/11 changed files
- Comments generated: 0 new
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
test/Grains/TestGrains/StatelessWorkerStreamConsumerGrain.cs:2
using System.Linq;is unused in this file, which will trigger IDE0005 (and can fail the build when code-style warnings are treated as errors). Remove the directive.
using System.Collections.Concurrent;
using System.Linq;
- Files reviewed: 12/12 changed files
- Comments generated: 0 new
- Review effort level: Lite
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>
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
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 0836f485-4052-4d24-b890-e460aac92596
d441624 to
b7b5e9c
Compare
Problem
Stateless worker grains currently reject stream consumer extension binding, forcing applications to route stream deliveries through a regular grain before parallel stateless processing.
Solution
IStreamSubscriptionObserver.OnSubscribedwhile retaining grain-level implicit and explicit subscription ownership.Rationale
Each delivery attempt targets the stateless worker grain identity and executes on one selected local activation. This removes the extra regular-grain hop while preserving provider delivery and retry behavior. Activation-local rewind handshakes are disabled because delivery can move between workers; applications which require checkpoint resume continue to use a regular grain consumer.
Closes #433
Closes #10653
Microsoft Reviewers: Open in CodeFlow