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:

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:

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.

State Manager example: an order grain saves its submitted state and outbox together. The outbox processor publishes to the order feed through a stream and calls the notification grain through RPC. Both receivers save message progress with their state; the feed uses StreamManager to attach its handler.

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:

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.

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