Skip to content

feat(durable-messaging): compose opaque inbox and outbox runtime - #10693

Open
ReubenBond wants to merge 224 commits into
dotnet:mainfrom
ReubenBond:rb-laughing-system
Open

ReubenBond wants to merge 224 commits into
dotnet:mainfrom
ReubenBond:rb-laughing-system

Conversation

@ReubenBond

@ReubenBond ReubenBond commented Aug 20, 2026 •

Copy link
Copy Markdown
Member

Problem

Applications need durable delivery which commits handler effects and outgoing intent with inbox completion, while keeping application protocols and byte ownership clear.

Solution

Compose the journaled inbox and outbox through AddDurableMessaging, IDurableMessagingGrain, and the existing durable grain model. The transport carries only message, sender, and receiver IDs plus an opaque ArcBuffer payload. Each inbox has one non-generic handler, its context exposes the received envelope and Complete(), and outgoing staging uses the injected outbox's synchronous Send operation.

Applications perform asynchronous local preparation, then execute shared updates, outgoing staging, and completion synchronously through handler return. The runtime owns persistence acknowledgement and dispatch. Independent Arc lifetimes cover acceptance, handler execution, journaling, serialization, retries, deletion, and activation teardown.

Encoding, dispatch, reply destinations, and business-operation deduplication live in application code. The new self-contained samples/DurableMessaging stock-reservation application submits fresh transport IDs for one hierarchical business operation and verifies the original outcome, two journal-acknowledged replies, and one stock decrement.

The conceptual, idempotency, practical-recipe, and operations guides describe the simplified surfaces and actual persistence/ownership boundaries. Compiled examples and the production-path sequential benchmark use the same opaque protocol. Provider-cutover controls cover persisted owner handles, recovery, and cleanup in the original provider.

Layering

The updated dependency chain is:

Review only this composition layer. Each layer includes its exact published parent. The dependent fork PRs remain targeted at main.

Copilot AI lite review requested due to automatic review settings August 20, 2026 00:09

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Adds a new durable, grain-scoped inbox/outbox messaging primitive (Microsoft.Orleans.DurableMessaging) built on Orleans Journaling + Durable Jobs, alongside supporting journaling participant/observer capabilities and Durable Jobs execution semantics needed for reliability, routing, and recovery.

Changes:

  • Introduces Microsoft.Orleans.DurableMessaging (envelopes, routing handlers, inbox/outbox contracts, DI/hosting integration, diagnostics, and instrumentation).
  • Extends Journaling with composable grain participants and write/recovery observers, plus rollback fencing and codec fixes for snapshot reference-scope replay.
  • Extends Durable Jobs with feature handler registry, explicit RetryAt disposition/reset-rescheduling support, and improved execution deduplication.
Show a summary per file
File Description
test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs Adds regression tests ensuring pre-change snapshot payloads replay with independent reference scopes.
test/Orleans.Journaling.Tests/JournaledGrainParticipantTests.cs Adds integration coverage for composed journaling participants initialization and failure propagation.
test/Orleans.Journaling.Tests/DurableListDirectWriteTests.cs Updates test double to implement new revert/rollback surface.
test/Orleans.Journaling.Tests/DurableCollectionDirectWriteTests.cs Updates test double to implement new revert/rollback surface.
test/Orleans.DurableMessaging.Tests/Support/SnapshotProbe.cs Adds snapshot-based test synchronization utility.
test/Orleans.DurableMessaging.Tests/Support/HandlerProbe.cs Adds barrier utility to block/release handlers deterministically in tests.
test/Orleans.DurableMessaging.Tests/Support/DurableMessagingTestGrains.cs Adds durable-messaging test grain, message/effect models, and handler behavior for scenarios.
test/Orleans.DurableMessaging.Tests/Support/DurableMessagingClusterFixture.cs Adds reusable in-process cluster fixture with Durable Jobs + Journaling + Durable Messaging wiring.
test/Orleans.DurableMessaging.Tests/Support/ControlledJournalStorageProvider.cs Adds controllable journal storage provider for injecting write blocking/failures in tests.
test/Orleans.DurableMessaging.Tests/Orleans.DurableMessaging.Tests.csproj Introduces new Durable Messaging test project.
test/Orleans.DurableMessaging.Tests/Hosting/PublicDurableMessagingRegistrationTests.cs Validates DI registrations, options behavior, and absence of friend-access requirements.
test/Orleans.DurableMessaging.Tests/Functional/MultiSiloDurableMessagingFailoverTests.cs Adds multi-silo failover recovery test for stable inbox job ownership and exactly-once effects.
test/Orleans.DurableMessaging.Tests/Functional/InboxCapacityBehaviorTests.cs Adds backpressure + recovery behavior tests under capacity constraints.
test/Orleans.DurableMessaging.Tests/Functional/DedupeExpiryBehaviorTests.cs Adds dedupe expiry/compaction behavior test allowing later reprocessing.
test/Orleans.DurableMessaging.Tests/Contracts/HandlerRoutingContractTests.cs Adds routing contract tests for exact, prefix, correlation, and typed handler behaviors.
test/Orleans.DurableMessaging.Tests/Contracts/DurableEnvelopeContractTests.cs Adds envelope/builder/serializer contract tests including reply-to and context.
test/Orleans.DurableMessaging.Tests/Contracts/DeliveryAndOptionsContractTests.cs Adds contracts for delivery result/status and options validation defaults/boundaries.
test/Orleans.DurableJobs.Tests/DurableJobs/JobShardTests.cs Adds shard-level test distinguishing reschedule reset vs failure retry attempt semantics.
test/Orleans.DurableJobs.Tests/DurableJobs/JobShardManagerTestsRunner.cs Adds reassignment test ensuring reschedule reset semantics persist through shard reassignment.
test/Orleans.DurableJobs.Tests/DurableJobs/DurableJobFeatureHandlerTests.cs Adds coverage for registry behavior and durable job run-result compatibility semantics.
test/Orleans.Core.Tests/Orleans.Core.Tests.csproj Links HierarchicalKey tests into core test project.
test/Orleans.Core.Tests/DurableJobs/ShardExecutorTests.cs Adds executor tests for RetryAt, unknown disposition handling, and legacy shard behavior.
test/Orleans.Core.Tests/DurableJobs/DurableJobsExtensionsTests.cs Adds DI scope/registry behavior tests and explicit replacement/decoration rejection.
test/Orleans.Core.Tests/DurableJobs/DurableJobReceiverExtensionTests.cs Updates execution identity to (JobId, RunId), adds retention/dedup + feature handler precedence tests.
test/Benchmarks/Journaling/DurableListJournalBenchmarks.cs Updates benchmark test double to implement new revert/rollback surface.
src/Orleans.Journaling/JournaledStateManager.cs Adds rollback support flag, observer notifications, recovery fencing, and revert-pending-changes support.
src/Orleans.Journaling/IJournaledStateObserver.cs Introduces observer interface for write/recovery boundaries.
src/Orleans.Journaling/IJournaledStateManager.cs Adds SupportsRollback, observer registration, and RevertPendingChangesAsync surface.
src/Orleans.Journaling/IJournaledGrainParticipant.cs Introduces participant initialization hook for composed journaling features.
src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableSetCommandCodec.cs Fixes snapshot replay to reset reference scope per element.
src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableQueueCommandCodec.cs Fixes snapshot replay to reset reference scope per element.
src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableListCommandCodec.cs Fixes snapshot replay to reset reference scope per element.
src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableDictionaryCommandCodec.cs Fixes snapshot replay to reset reference scope per key/value entry.
src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryCommandCodecHelpers.cs Adds helper to read values using independent serializer sessions (reference-scope reset).
src/Orleans.Journaling/DurableGrain.cs Initializes composed participants before recovery and documents WriteStateAsync.
src/Orleans.DurableMessaging/RoutePrefixHandler.cs Adds prefix-based routing handler base class and documentation.
src/Orleans.DurableMessaging/RouteKeyHandler.cs Adds exact-route routing handler base class and documentation.
src/Orleans.DurableMessaging/README.md Adds package readme with configuration and semantics overview.
src/Orleans.DurableMessaging/Orleans.DurableMessaging.csproj Adds new packable Durable Messaging project.
src/Orleans.DurableMessaging/InboxHandlerContext.cs Adds handler context implementation for sending outbox messages and building envelopes.
src/Orleans.DurableMessaging/IInboxHandlerContext.cs Adds public handler context contract.
src/Orleans.DurableMessaging/IInboxHandler.cs Adds handler interfaces including typed handler adapter and documentation.
src/Orleans.DurableMessaging/IDurableOutbox.cs Adds outbox public contract.
src/Orleans.DurableMessaging/IDurableMessagingDiagnostics.cs Adds diagnostics contract + internal implementation for dead letters.
src/Orleans.DurableMessaging/IDurableInboxExtension.cs Adds grain extension contract for durable inbox delivery.
src/Orleans.DurableMessaging/IDurableInbox.cs Adds inbox public contract including handler registration and lookup.
src/Orleans.DurableMessaging/Hosting/DurableMessagingExtensions.cs Adds ISiloBuilder/IServiceCollection registration, option validation, and rollback requirement checks.
src/Orleans.DurableMessaging/DurableMessagingPumpResults.cs Adds pump execution result tracking and one-shot timer handle helper.
src/Orleans.DurableMessaging/DurableMessagingInstruments.cs Adds meters/counters/histograms for inbox/outbox behavior and depth tracking.
src/Orleans.DurableMessaging/DurableMessagingGrainParticipant.cs Ensures messaging services materialize via journaling participant initialization.
src/Orleans.DurableMessaging/DurableMessageState.cs Adds durable state models for inbox/outbox attempts and dead letters.
src/Orleans.DurableMessaging/DurableInbox.cs Adds inbox implementation with handler registration/selection and storage-backed message access.
src/Orleans.DurableMessaging/DurableEnvelopeData.cs Adds deferred (slice-based) envelope body/context storage and accessors.
src/Orleans.DurableMessaging/DurableEnvelope.cs Adds durable envelope type (routing, correlation, reply-to, metadata, and payload).
src/Orleans.DurableMessaging/DeliveryStatus.cs Adds delivery status enum contract.
src/Orleans.DurableMessaging/DeliveryResult.cs Adds delivery result struct contract and factories.
src/Orleans.DurableMessaging/CorrelationHandler.cs Adds correlation-hierarchy-based routing handler base class.
src/Orleans.DurableMessaging/Configuration/DurableInboxOptions.cs Adds durable messaging options and validation contract.
src/Orleans.DurableJobs/ShardExecutor.cs Adds RetryAt handling with reset-rescheduling path and explicit legacy-shard failure.
src/Orleans.DurableJobs/JournaledJobShard.cs Implements reset-rescheduling in journaled shard via IResettableJobShard.
src/Orleans.DurableJobs/JobShard.cs Adds reset-rescheduling support to base shard and introduces IResettableJobShard.
src/Orleans.DurableJobs/IDurableJobReceiverExtension.cs Updates receiver extension for feature-handler lookup, new execution identity, and retention semantics.
src/Orleans.DurableJobs/IDurableJobHandlerRegistry.cs Adds activation-scoped feature handler registry + lookup implementation.
src/Orleans.DurableJobs/Hosting/DurableJobsOptions.cs Adds completed-attempt retention option and validation.
src/Orleans.DurableJobs/Hosting/DurableJobsExtensions.cs Registers registry in DI and explicitly rejects replacement/decoration.
src/Orleans.DurableJobs/DurableJobRunResult.cs Adds RetryAt disposition and related fields/compatibility semantics.
src/api/Orleans.Journaling/Orleans.Journaling.cs Updates generated API surface for journaling observer/participant/rollback additions.
src/api/Orleans.DurableJobs/Orleans.DurableJobs.cs Updates generated API surface for feature handlers/registry and RetryAt semantics.
src/api/Orleans.Core.Abstractions/Orleans.Core.Abstractions.cs Updates generated API surface to include HierarchicalKey type and codec.
Orleans.slnx Adds Durable Messaging projects (src + tests) to the solution.
docs/site/src/data/unpublished-api-packages.json Marks Durable Messaging package as unpublished.
docs/site/src/data/external-link-allowlist.json Allow-lists NuGet link for unpublished Durable Messaging package.
docs/site/src/content/docs/toc.yml Adds Durable Messaging doc page to Journaling section.
docs/site/src/content/docs/resources/nuget-packages.md Adds Durable Messaging to NuGet package catalog and updates guidance blurb.
docs/site/src/content/docs/grains/durable-messaging.md Adds end-user documentation for durable messaging semantics and requirements.

Review details

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

  • Files reviewed: 83/83 changed files
  • Comments generated: 4
  • Review effort level: Lite

Comment thread src/Orleans.DurableMessaging/DurableInbox.cs Outdated
Comment thread src/Orleans.Journaling/IJournaledStateObserver.cs Outdated
Comment thread src/Orleans.DurableMessaging/RouteKeyHandler.cs Outdated
Comment thread src/Orleans.DurableMessaging/InboxHandlerContext.cs Outdated
Copilot AI review requested due to automatic review settings August 20, 2026 00:46

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (4)

Previously missed (1) — in code that hasn't changed since the last review.

src/Orleans.DurableMessaging/DurableInbox.cs:43

  • Exact-route handler storage uses the default string comparer, which is culture-sensitive. Route keys are protocol identifiers and should use ordinal comparisons (matching the rest of this PR’s routing semantics). Use StringComparer.Ordinal for _exactRouteHandlers to avoid surprising behavior under non-invariant cultures (e.g., Turkish-I issues).
        _inbox = inbox;
        _processed = processed;
        _handlers = new List<IInboxHandler>();
        _exactRouteHandlers = new Dictionary<string, IInboxHandler>();
        _capacity = capacity;

src/Orleans.Journaling/IJournaledStateObserver.cs:16

  • The remarks contradict the interface contract: they say the manager does not notify observers before a write, but this interface defines OnWriteStarted and JournaledStateManager invokes it before capturing state. This is confusing for implementers; update the remarks to reflect the actual callback sequence.
    src/Orleans.DurableMessaging/RouteKeyHandler.cs:21
  • The remarks suggest implementing IInboxHandler directly for prefix-based routing, but this PR also introduces RoutePrefixHandler for that exact scenario. Updating the docs helps steer users to the supported helper type and keeps the guidance consistent.
    src/Orleans.DurableMessaging/InboxHandlerContext.cs:125
  • The XML docs reference IStateMachineManager.WriteStateAsync(), which doesn't exist in this package and appears to be a leftover name. This should point to IJournaledStateManager.WriteStateAsync() to avoid misleading API consumers.
  • Files reviewed: 83/83 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings August 20, 2026 01:02

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (2)

Previously missed (2) — in code that hasn't changed since the last review.

src/Orleans.DurableJobs/IDurableJobReceiverExtension.cs:185

  • LongPollGetJobStatusAsync creates a CancellationTokenSource but never cancels it. That means the Task.Delay will always run to completion even when the job finishes quickly, causing avoidable timer allocations/work per call.
                using var cts = new CancellationTokenSource();
                var longPollDuration = TimeSpan.FromTicks(Math.Min(_shared.MessagingOptions.ResponseTimeout.Divide(2).Ticks, _shared.Options.JobStatusPollInterval.Ticks));
                await Task.WhenAny(Task.Delay(longPollDuration, cts.Token), state.Task);

                if (!state.Task.IsCompleted)
                {
                    return DurableJobRunResult.PollAfter(_shared.Options.JobStatusPollInterval);
                }

src/Orleans.DurableMessaging/Hosting/DurableMessagingExtensions.cs:58

  • The DurableInboxOptions validation predicate swallows all exceptions and turns them into a generic validation failure. That can hide unexpected exceptions (eg, NullReferenceException) and makes misconfiguration harder to diagnose. Consider only converting known validation exceptions into a failed predicate and letting other exceptions bubble.
        optionsBuilder.Validate(
            options =>
            {
                try
                {
                    options.Validate();
                    return true;
                }
                catch
                {
                    return false;
                }
            },
            "DurableInboxOptions validation failed.");
  • Files reviewed: 83/83 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings August 20, 2026 02:53

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (1)

Previously missed (1) — in code that hasn't changed since the last review.

test/Orleans.DurableMessaging.Tests/Support/SnapshotProbe.cs:33

  • WaitAsync adds a waiter to the per-grain list and then returns waiter.Completion.Task.WaitAsync(...), but if that wait times out/throws, the waiter remains in _waiters indefinitely. This can leak memory and can also keep evaluating stale predicates on subsequent Publish calls. Consider awaiting with a cleanup finally that removes the waiter (and optionally removes the empty list from _waiters).
  • Files reviewed: 85/85 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings August 20, 2026 02:59

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (4)

Previously missed (3) — in code that hasn't changed since the last review.

src/Orleans.Journaling/JournaledStateManager.cs:309

  • WorkLoop captures a stable observer snapshot into the local observers array, but OnWriteStarted() is invoked by iterating _observers directly. If an observer is registered during OnWritePreparingAsync (or from within another observer callback), it can change _observers mid-iteration, causing inconsistent callback sequencing and potential InvalidOperationException due to collection mutation during enumeration. Use the captured observers snapshot for OnWriteStarted to keep the write boundary consistent and mutation-safe.
    src/Orleans.Journaling/JournaledStateManager.cs:493
  • OnWriteCompleted() is invoked by iterating _observers directly, even though a stable observers snapshot was captured earlier for this write. If observers are registered during the write pipeline, enumerating the live HashSet can throw (collection modified) or cause some observers to see OnWriteCompleted without the corresponding OnWritePreparingAsync/OnWriteStarted. Use the per-write snapshot for completion callbacks.

This issue also appears on line 516 of the same file.
src/Orleans.Journaling/JournaledStateManager.cs:826

  • Recovery completion notifies observers by enumerating _observers directly inside a lock. If any observer callback registers another observer (re-entrantly) this can mutate the HashSet during enumeration and throw, potentially breaking recovery. Take a snapshot of _observers before iterating so callbacks are mutation-safe and sequencing is stable.

src/Orleans.Journaling/JournaledStateManager.cs:516

  • In the no-op write path (!hasCommittedBuffer), OnWriteCompleted() is called by iterating _observers directly. This has the same collection-mutation risk as the committed write path and can also notify observers which were registered after OnWritePreparingAsync ran for this write. Prefer using the stable observers snapshot captured at the start of the work item.
  • Files reviewed: 85/85 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings August 20, 2026 03:09
@ReubenBond

Copy link
Copy Markdown
Member Author

Addressed the latest observer-snapshot review findings in \3d4829b8d. Each write now uses one stable observer snapshot for preparation, start, and completion (including no-op writes), and recovery snapshots observers before callbacks so re-entrant registration cannot mutate an active enumeration. Focused regression coverage passes on net8.0 and net10.0.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (1)

Previously missed (1) — in code that hasn't changed since the last review.

test/Orleans.DurableMessaging.Tests/Support/SnapshotProbe.cs:33

  • WaitAsync adds a waiter to the per-grain list but never removes it if the 30s timeout elapses (WaitAsync throws). That can leak waiters across a test run and slow down Publish() due to ever-growing lists.
  • Files reviewed: 85/85 changed files
  • Comments generated: 1
  • Review effort level: Lite

Comment thread src/Orleans.Journaling/IJournaledStateManager.cs Outdated
Copilot AI review requested due to automatic review settings August 20, 2026 03:13
@ReubenBond

Copy link
Copy Markdown
Member Author

Addressed the remaining suppressed review findings in two focused commits:\n\n- \�222b4b77\ removes timed-out SnapshotProbe waiters so stale predicates are not retained or evaluated; regression coverage passes on net8.0 and net10.0.\n- \�d214b07b\ narrows durable inbox options validation to the documented \ArgumentOutOfRangeException, allowing unexpected faults to surface; the options-contract test passes on net8.0 and net10.0.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (1)

src/Orleans.Journaling/IJournaledStateManager.cs:58

  • IJournaledStateManager is a public interface and RevertPendingChangesAsync is added as a required member (no default implementation). This is a breaking change for any external implementations of IJournaledStateManager. Consider providing a default interface implementation (similar to RegisterObserver) which throws NotSupportedException, and use SupportsRollback to indicate capability.
  • Files reviewed: 86/86 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings August 20, 2026 03:19

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

  • Files reviewed: 86/86 changed files
  • Comments generated: 2
  • Review effort level: Lite

Comment thread src/Orleans.Journaling/IJournaledStateManager.cs
Comment thread src/Orleans.DurableMessaging/DurableInbox.cs Outdated
Copilot AI review requested due to automatic review settings August 20, 2026 03:24

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review details

Suppressed comments (1)

Previously missed (1) — in code that hasn't changed since the last review.

src/Orleans.DurableMessaging/DurableOutbox.cs:168

  • The XML docs reference IStateMachineManager.WriteStateAsync(), but durable messaging commits are driven by journaling. This should reference IJournaledStateManager.WriteStateAsync() to match the rest of the package docs and avoid misleading API consumers.
    /// The message is persisted atomically with grain state when IStateMachineManager.WriteStateAsync()
    /// is called. The background pump will deliver the message to the target grain ONLY AFTER
    /// the message has been durably persisted.
  • Files reviewed: 86/86 changed files
  • Comments generated: 0 new
  • Review effort level: Lite

Copilot AI review requested due to automatic review settings August 20, 2026 03:30

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

One or more issues must be addressed before approval.

2 open findings

🧠 Review effort: Lite

Comment on lines +268 to +269
pending = new(envelope);
_pendingMessages.Add(envelope.MessageId, pending);
{
internal const byte LegacyFramingVersion = OrleansBinaryV0JournalReader.FramingVersion;
internal const byte FramingVersion = 1;
internal const byte FramingVersion = 0;

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants