A messaging toolbox for Microsoft Orleans
Orleans handles grain identity, routing, and activation. Reliable workflows still need application code to coordinate state changes and message delivery. Two familiar distributed systems problems need to be addressed:
- Ambiguous outcomes. After a timeout, a caller cannot tell whether the receiver never received its call, is still processing it, or finished but the reply was lost. A stream provider can also accept a message without its acknowledgement reaching the sender. Retrying may repeat work.
- Dual-write problem. Saving state and sending a message are two separate operations. Send first and the state save may fail; save first and the process may stop before sending. Ordinary Orleans persistence does not make the two operations atomic.
I built Egil.Orleans.Messaging to help with this without distributed transactions, using the established outbox pattern: save business state, receiver progress, and outgoing messages in one atomic local commit, then dispatch. Retries keep the same outbox identity so receivers can deduplicate grain calls and stream publications. The goal is at-least-once delivery with visible, recoverable failures.
The v1 release is available on NuGet.
There are five tools. Each can be used on its own, or in the combination a particular grain needs:
- State manager: Failed state writes and ambiguous outcomes can leave a grain’s managed state and other in-memory fields unreliable. By default,
IStateManager<T>blocks further state changes and storage I/O and requests grain deactivation.StateandHasUnsavedChangesremain readable for inspection, but can describe stale or unsaved values. Deactivation discards that grain activation’s working memory; a new grain activation reloads persisted state. Read-back recovery is an alternative when restoring the stored object is sufficient; rereading it alone does not rebuild other fields or properties in a grain. - Outbox: To address the dual-write problem, save
Outbox<T>together with the grain’s business state.OutboxProcessor<T>then dispatches its messages and manages retry timers and reminders. Failed sends surface throughAcknowledgeFailuresAsync; the grain decides whether to keep them pending or move them to dead-letter state. - Message tracker: Retries after ambiguous outcomes can repeat work.
MessageTrackertracks stream positions by default. A grain can opt intoOutboxIdentityto recognize an outbox message republished at a new position, storing a receipt for each logical message and stream source. Grain calls carrying delivery tokens use sender high-water marks. Persist tracking changes with the business state so that progress survives grain reactivation. - Stream manager: Orleans keeps subscriptions across grain activations, but your grain still has to attach its handlers again.
StreamManagerhandles attaching and resuming subscriptions, using saved positions where the stream provider supports replay. It also correlates consumer and producer spans when the provider carries trace context, making distributed workflows easier to follow. - Journaling (preview): For journaled grains using the new Journaling storage in Orleans, the Journaling companion addresses the dual-write problem and retries after ambiguous outcomes through
IDurableOutbox<T>andIDurableMessageTracker. They work alongside Orleans’IDurableValue<T>so receiver progress, business state, and outgoing messages can share one journal write, with incremental messaging changes. This is an alternative to theIStateManager<T>persistence path.
A grain that only receives calls might use just the state manager, adding a tracker for retried calls carrying delivery tokens. A sender can pair the state manager with an outbox, though the outbox also works with ordinary IPersistentState<T>. A stream consumer can use StreamManager alone, adding tracking and state management when it needs persisted progress.
For journaled grains, use the preview tracker and outbox with journaled business state, without IStateManager<T>. The companion’s example grain requests grain deactivation after write failures and recovers persisted state on a later grain activation. Its tracker supports the same tracking modes and retention settings described below.
The default state-manager recovery policy mirrors the approach in Orleans journaling on main: stop further persistence, request grain deactivation, and leave local state readable for inspection. That Orleans behavior has not yet been released.
An order, a stream, and a grain call
The processing pattern is one atomic local save of receiver progress, business state, and outgoing messages, followed by dispatch and later acknowledgement writes.
In this State Manager example, an order grain publishes OrderSubmitted to an OrderFeedGrain through an Orleans stream, and calls a NotificationGrain directly to create an in-app notification. Each order is submitted once; each receiver gets one logical message from it.
I have omitted using directives, Orleans serialization attributes, grain interfaces, and hosting from these excerpts. The immutable collections are from System.Collections.Immutable.
The two outgoing messages share an outbox. The stream carries OrderSubmitted directly; AddStreamPostman sends the outbox identity alongside it in Orleans request context:
public interface IOrderMessage { }
public sealed record OrderSubmitted(Guid OrderId) : IOrderMessage;
public sealed record CreateNotification(Guid OrderId) : IOrderMessage;
Save the order, then dispatch
The constructor registers the delivery routes and an acknowledgement callback. That callback removes and saves only the messages successfully dispatched.
public sealed record OrderState
{
public bool Submitted { get; init; }
public Outbox<IOrderMessage> Outbox { get; init; } = [];
}
public sealed class OrderGrain : Grain, IOrderGrain, IOutboxGrain
{
private readonly IStateManager<OrderState> storage;
private readonly OutboxProcessor<IOrderMessage> processor;
public OrderGrain(
[PersistentState(stateName: "order", storageName: "state")]
IStateManager<OrderState> storage)
{
this.storage = storage;
processor = this.RegisterOutboxProcessor(
outboxAccessor: () => storage.State.Outbox,
configure: (OutboxProcessorOptions<IOrderMessage> options) =>
{
// Establish a fallback before the grain accepts business calls.
options.ReminderPolicy = OutboxReminderPolicy.KeepRegistered;
options.AcknowledgePostedAsync = async (
ImmutableArray<OutboxMessageEnvelope<IOrderMessage>> posted,
CancellationToken ct) =>
{
// Persist removal of successful sends; failures stay pending.
await storage.WriteAsync(storage.State with
{
Outbox = storage.State.Outbox.RemoveRange(posted)
}, cancellationToken: ct);
};
})
// Carry the outbox identity alongside the plain stream payload.
.AddStreamPostman<OrderSubmitted>(
streamProviderName: "events",
streamId: (OrderSubmitted message) => StreamId.Create(ns: "orders", key: message.OrderId))
.AddPostman<CreateNotification>(
postman: (CreateNotification message, OutboxSequenceToken token,
IGrainFactory grains) =>
grains.GetGrain<INotificationGrain>(message.OrderId)
.ApplyAsync(message: message, token: token));
}
// Resume pending delivery after reactivation.
public override async Task OnActivateAsync(CancellationToken cancellationToken)
=> await processor.PostInBackgroundAsync(cancellationToken: cancellationToken);
public async Task SubmitAsync()
{
if (storage.State.Submitted)
return;
Guid id = this.GetPrimaryKey();
// Persist business state and outgoing intent together, before dispatch.
await storage.WriteAsync(storage.State with
{
Submitted = true,
Outbox = storage.State.Outbox
.Add(new OrderSubmitted(id))
.Add(new CreateNotification(id))
});
// Schedule dispatch; the processor manages subsequent retries.
await processor.PostInBackgroundAsync();
}
}
OutboxProcessor uses grain timers for retries while the grain is active. Reminders provide a durable way to wake it up later, and IOutboxGrain forwards their callbacks automatically.
The OutboxProcessorOptions<TMessage>.ReminderPolicy property has two options:
OnDeactivation(default): Waits until orderly grain deactivation to register a reminder when recovery is needed. This minimizes reminder API calls. The trade-off is that an abrupt silo crash can bypass registration, and a registration failure during deactivation only produces the warning log eventOutboxDeactivationReminderFailedthat can be monitored for. Deactivation continues. If the grain is not activated by other means later, pending outbox messages can remain undelivered indefinitely. Manual reactivation may be needed; theOnActivateAsynchook above starts posting automatically. Without that hook, recovery also needs an explicit post.KeepRegistered(used above): Establishes a fallback reminder during grain activation, keeps it across batches, and removes it on deactivation when the outbox is empty and no storage failure has fenced the grain. This makes more reminder API calls: initial registration, updates when switching between active and idle periods, and eventual removal. Constructor registration also makes grain activation depend on the reminder store: if the initial registration fails, activation fails. SetIdleReminderPeriodlonger than the grain’s idle collection age, with a margin, because reminder ticks reset idleness and can otherwise keep the grain alive.
Both policies account for storage failures that fenced the grain. If any recorded fencing failure is UnknownOutcome or DidNotPersist, the processor ensures an active recovery reminder, even when the local outbox looks empty. If all recorded fencing failures are Conflict, it leaves reminders unchanged; recovery depends on another owner or a later grain activation. See Choosing a recovery policy for the reasoning.
The OutboxDeactivationReminderFailed warning includes the grain, reminder, and observed OutboxItemCount; null means the count could not be read. The count reflects local state, so even zero does not rule out durable work after fencing. For diagnostics, storage.LastFailureKind and storage.Options.RecoveryPolicy expose the recorded failure classification and effective recovery policy, even after fencing.
Publish to several streams
One outbox item can also publish to several streams on the same provider. For example, replace the OrderSubmitted registration above to send each event to both the order feed and an audit stream:
.AddStreamPostman<OrderSubmitted>(
streamProviderName: "events",
streamIds: (OrderSubmitted message) =>
[
StreamId.Create(ns: "orders", key: message.OrderId),
StreamId.Create(ns: "order-audit", key: message.OrderId)
])
The processor publishes in destination order and acknowledges the item only after all publications succeed. If an attempt fails partway through, the retry starts again with the same outbox identity. Receivers using OutboxIdentity can reject repeated deliveries while their receipts are retained. Each stream has its own receipt; publication across streams is not atomic. Keep destination selection stable across retries, and establish explicit subscriptions before sending.
Receive through a stream
An Orleans StreamSequenceToken represents a position within a stream. StreamManager gives the handler a StreamCursor that also identifies the provider and full StreamId (namespace and key), plus the stable outbox identity when present.
MessageTracker defaults to StreamPosition, which tracks provider positions. This suits idempotent handlers, where applying the same event again leaves the same result. The feed below appends an order to a list, so this example sets OutboxIdentity globally to recognize a repeated publication even at a new position.
Set the shared defaults once during silo setup:
siloBuilder.ConfigureMessageTracker(configure: (MessageTrackerOptions options) =>
{
options.StreamTrackingMode = StreamTrackingMode.OutboxIdentity;
options.RetentionPeriod = null;
});
These defaults apply to newly created and deserialized trackers. A grain can override them with Tracker.Configure; reapply those overrides after each state load because tracker settings are not persisted. See global settings and per-grain overrides for the configuration hooks.
StreamPosition stores less tracking data, but a republication at a new position can repeat work. OutboxIdentity also stores a receipt for each logical message on each source. Here, RetentionPeriod = null keeps entries until explicit eviction, so old retries remain recognizable but the tracker can keep growing. A positive retention period allows automatic cleanup at the cost of accepting an old message again after its entry expires.
- Keep independent sources separate. A projection might consume the
ordersnamespace from two regional Event Hubs providers. Position 500 on one provider must not make position 20 on the other look like an old message. Each provider and full stream ID gets its own high-water mark; different order IDs in the same namespace are separate too. Identity tracking adds receipts scoped to that same source. - Save progress for resumption. As the grain accepts messages, it advances the appropriate source’s high-water mark and saves it with the business update. After grain reactivation,
StreamManageruses that source’s saved token to resume where the provider supports replay. InOutboxIdentitymode, an unseen outbox message can still be accepted at an older provider position without moving the retained high-water mark backwards. - Recognize an outbox retry at a new position. If publication succeeds but the sender’s acknowledgement write fails, the retry can receive a new
StreamSequenceToken. WithOutboxIdentityenabled, the cursor’s stable identity lets the retained receipt prevent another business update. Once that receipt expires or is evicted, the message can be accepted again. Receipts are scoped to each source; they do not deduplicate the same event across different providers.
public sealed record FeedState
{
public MessageTracker Tracker { get; init; } = new MessageTracker();
public ImmutableList<Guid> Orders { get; init; } = [];
}
public sealed class OrderFeedGrain : Grain, IOrderFeedGrain
{
private readonly IStateManager<FeedState> storage;
private readonly StreamManager streams;
public OrderFeedGrain(
[PersistentState(stateName: "feed", storageName: "state")]
IStateManager<FeedState> storage)
{
this.storage = storage;
streams = this.RegisterStreamManager(getTracker: () => storage.State.Tracker)
.ConfigureExplicitSubscription<OrderSubmitted>(
streamProviderName: "events",
streamId: StreamId.Create(ns: "orders", key: this.GetPrimaryKey()),
onNextAsync: ApplyAsync);
}
// Attach or resume handlers each time Orleans activates this grain.
public override Task OnActivateAsync(CancellationToken cancellationToken)
=> streams.EnsureExplicitSubscriptionsAsync(cancellationToken: cancellationToken);
private async Task ApplyAsync(OrderSubmitted message, StreamCursor cursor)
{
// Reject a retry of an outbox message already applied to this stream.
if (!storage.State.Tracker.TryAcceptMessage(cursor, out MessageTracker tracker))
return;
await storage.WriteAsync(storage.State with
{
// Save the receipt, provider progress, and feed update together.
Tracker = tracker,
Orders = storage.State.Orders.Add(message.OrderId)
});
}
}
The explicit subscription must be established before publishing; here, that means activating the feed first. On later activations, StreamManager reattaches the handler and uses saved positions where the provider supports replay. If no matching checkpoint remains, it attaches without a resume token and lets Orleans and the provider choose the starting position.
Implicit subscriptions are supported too. Orleans can activate the receiving grain when an event arrives, so the publisher does not need to activate the feed first. For the same orders namespace and order-ID key, replace the explicit example’s class setup with this, keeping FeedState and ApplyAsync unchanged:
[ImplicitStreamSubscription(streamNamespace: "orders")]
public sealed class OrderFeedGrain : Grain, IOrderFeedGrain, IImplicitStreamGrain
{
private readonly IStateManager<FeedState> storage;
public OrderFeedGrain(
[PersistentState(stateName: "feed", storageName: "state")]
IStateManager<FeedState> storage)
{
this.storage = storage;
this.RegisterStreamManager(getTracker: () => storage.State.Tracker)
.ConfigureImplicitSubscription<OrderSubmitted>(
streamNamespace: "orders",
onNextAsync: ApplyAsync);
}
// Keep the ApplyAsync method from the explicit example.
}
IImplicitStreamGrain forwards Orleans’ subscription callbacks to the manager. The streams field and OnActivateAsync override from the explicit example are no longer needed. See Orleans implicit subscriptions for the namespace-to-grain mapping and provider configuration.
Receive through a grain call
The direct call passes an OutboxSequenceToken to the notification grain. It saves that token and the notification together, so a retry cannot add the notification twice while the sender’s tracking entry is retained.
public sealed record NotificationState
{
public MessageTracker Tracker { get; init; } = new MessageTracker();
public ImmutableList<string> Notifications { get; init; } = [];
}
public sealed class NotificationGrain(
[PersistentState(stateName: "notifications", storageName: "state")]
IStateManager<NotificationState> storage)
: Grain, INotificationGrain
{
public async Task ApplyAsync(
CreateNotification message, OutboxSequenceToken token)
{
// Reject a duplicate or stale delivery token.
if (!storage.State.Tracker.TryAcceptMessage(token, out MessageTracker tracker))
return;
// Save sender progress and the notification together.
await storage.WriteAsync(storage.State with
{
Tracker = tracker,
Notifications = storage.State.Notifications.Add(
$"Order {message.OrderId} submitted.")
});
}
}
These excerpts assume durable storage and reminders, using System.Text.Json or Orleans binary serialization for the stored messaging types. Orleans’ built-in stream providers preserve request context, including the memory, Azure Queue, and Event Hubs providers. Custom adapters must preserve it too. Durability, replay, and delivery guarantees still depend on the provider, and a successful publication does not mean every subscriber has applied the event.
Tracker size depends on the mode. Position tracking keeps one checkpoint per provider and full stream ID, while identity tracking also keeps a receipt for each accepted logical message on that source. A growing number of stream IDs can grow the tracker in either mode.
Automatic eviction is disabled by default. Set a positive RetentionPeriod globally through ConfigureMessageTracker or per grain through Tracker.Configure to remove old receipts, stream checkpoints, and RPC sender entries as messages are accepted. Retention uses receiver acceptance time; duplicate attempts do not extend it. Choose a period that covers the application’s retry and replay window: after an entry expires or is evicted, an old message can be accepted again. Cleanup is saved with the tracking and business changes. Idle trackers do no background cleanup, and a time window is not a size limit.
For explicit cleanup, EvictStreamReceipts removes receipts while preserving stream checkpoints and RPC positions. Switching to StreamPosition stops adding receipts but does not erase existing ones. Persist the returned tracker to save the cleanup. See tracking modes and retention for the options and eviction methods.
Ordinary stream publishers without outbox metadata still use provider positions for tracking. The grain-call token path keeps a high-water mark per sender, so receivers handling multiple messages from one sender need an ordering or recovery strategy. Separately created outbox entries have separate identities: the order’s Submitted check prevents a repeated business command from creating another batch.
This messaging toolbox is used in production today in a system that handles tens of millions of messages per month.
See the library documentation for the current API and configuration, and the Orleans docs for persistence, message delivery guarantees, and stream subscriptions.
Hope this helps.
Comments