| name | orleans-grainservice-cache |
| description | Use when implementing or validating the Orleans pattern where a per-silo GrainService caches data by subscribing to a data grain that persists observer references, so updates continue after grain deactivation or silo restarts/migration. |
Orleans GrainService Cache Subscription Pattern
Use this skill to implement or verify the pattern where a data grain persists references to per-silo grain services so cache updates survive grain lifecycle events.
Problem this solves
You need an in-memory cache of data from a "data grain" that stays current when the data changes. Orleans Observers are ideal, but normal POCO services (e.g., singletons) lack an Orleans address and must resubscribe on an interval to maintain subscriptions. This creates overhead and causes stale data between resubscriptions when the data grain deactivates or its silo crashes.
GrainServices are addressable, so the data grain can persist their GrainIds and notify them directly. Subscriptions survive grain deactivation, silo crashes, silo restarts, and red/green deployments. Failed notifications (e.g., silo down) are handled gracefully — the GrainService automatically resubscribes on startup. GrainServices also unsubscribe cleanly during normal shutdowns.
Pattern summary
- A single data grain owns the authoritative value and persists subscriber/observer IDs (GrainIds).
- Each silo hosts a
GrainService that subscribes once on startup and writes to a POCO singleton cache (IDataCache).
- The data grain rehydrates subscriptions on activation and notifies observers via
ObserverManager.
- Subscriptions survive grain deactivation, silo crashes, silo restarts, and red/green deployments.
- Any type can implement
IDataGrainObserver, which implements IGrainObserver, so observers are not limited to a GrainService.
Contracts
public interface IDataGrain : IGrainWithStringKey
{
Task Subscribe(IDataGrainObserver subscriber);
Task Unsubscribe(IDataGrainObserver subscriber);
Task UpdateValue(string value);
Task<string?> GetValue();
}
public interface IDataGrainObserver : IGrainObserver
{
[AlwaysInterleave]
Task OnDataUpdated(string grainKey, string? value);
}
Data grain (authoritative + persistent subscribers)
[GenerateSerializer]
public sealed class DataGrainState
{
[Id(0)]
public string? Value { get; set; }
[Id(1)]
public ImmutableHashSet<GrainId> Subscribers
{
get => field ?? [];
set => field = value ?? [];
}
}
public sealed class DataGrain([PersistentState("data")] IPersistentState<DataGrainState> state,
ILogger<DataGrain> logger) : Grain, IDataGrain
{
private readonly ObserverManager<GrainId, IDataGrainObserver> observerManager =
new(TimeSpan.FromDays(365 * 10), logger);
public override async Task OnActivateAsync(CancellationToken cancellationToken)
{
foreach (var subscriberId in state.State.Subscribers)
{
var observer = GrainFactory.GetGrain<IDataGrainObserver>(subscriberId);
observerManager.Subscribe(subscriberId, observer);
}
await base.OnActivateAsync(cancellationToken);
}
public async Task Subscribe(IDataGrainObserver subscriber)
{
var subscriberId = subscriber.GetGrainId();
observerManager.Subscribe(subscriberId, subscriber);
state.State.Subscribers = state.State.Subscribers.Add(subscriberId);
await state.WriteStateAsync();
await subscriber.OnDataUpdated(this.GetPrimaryKeyString(), state.State.Value);
}
public async Task Unsubscribe(IDataGrainObserver subscriber)
{
var subscriberId = subscriber.GetGrainId();
observerManager.Unsubscribe(subscriberId);
var storedSubs = state.State.Subscribers.Remove(subscriberId);
if (state.State.Subscribers != storedSubs)
{
state.State.Subscribers = storedSubs;
await state.WriteStateAsync();
}
}
public Task<string?> GetValue() => Task.FromResult(state.State.Value);
public async Task UpdateValue(string value)
{
state.State.Value = value;
await state.WriteStateAsync();
await NotifySubscribersAsync(value);
}
private async Task NotifySubscribersAsync(string? value)
{
if (observerManager.Count == 0)
{
return;
}
var subscribersChangedDuringNotification = false;
await observerManager.Notify(async observer =>
{
try
{
await observer.OnDataUpdated(this.GetPrimaryKeyString(), value);
}
catch (SiloUnavailableException)
{
state.State.Subscribers = state.State.Subscribers.Remove(observer.GetGrainId());
subscribersChangedDuringNotification = true;
throw;
}
});
if (subscribersChangedDuringNotification)
{
await state.WriteStateAsync();
}
}
}
POCO singleton cache (recommended)
Use a regular singleton that can be injected into any grain or POCO/non-Orleans type.
public interface IDataCache
{
Task<string?> GetValue(string grainKey);
}
public sealed class DataCache(IGrainFactory grainFactory) : IDataCache
{
private readonly ConcurrentDictionary<string, Task<string?>> cache = new(StringComparer.Ordinal);
public Task<string?> GetValue(string grainKey)
=> cache.GetOrAdd(grainKey, async key =>
{
var grain = grainFactory.GetGrain<IDataGrain>(key);
var value = await grain.GetValue();
return value;
});
public void OnDataUpdated(string grainKey, string? value)
=> cache.AddOrUpdate(
grainKey,
_ => Task.FromResult(value),
(_, _) => Task.FromResult(value));
}
Per-silo grain service (addressable)
public interface ICacheGrainService : IDataGrainObserver, Orleans.Services.IGrainService
{
}
public sealed partial class CacheGrainService(
DataCache dataCache,
GrainId grainId,
Silo silo,
IGrainFactory grainFactory,
ILogger<CacheGrainService> logger)
: GrainService(grainId, silo, NullLoggerFactory.Instance), ICacheGrainService
{
private IDataGrainObserver? observerReference;
public override async Task Start()
{
await base.Start();
await SubscribeToDataGrainAsync();
}
public override async Task Stop()
{
if (observerReference is not null)
{
var grain = grainFactory.GetGrain<IDataGrain>(DataGrainConstants.GrainKey);
await grain.Unsubscribe(observerReference);
}
await base.Stop();
}
public Task OnDataUpdated(string grainKey, string? value)
{
dataCache.OnDataUpdated(grainKey, value);
LogCacheUpdatedLog(this.GetPrimaryKeyString(), grainKey, value);
return Task.CompletedTask;
}
private async Task SubscribeToDataGrainAsync()
{
observerReference = this.AsReference<IDataGrainObserver>();
var grain = grainFactory.GetGrain<IDataGrain>(DataGrainConstants.GrainKey);
await grain.Subscribe(observerReference);
LogSubscribed(this.GetPrimaryKeyString());
}
[LoggerMessage(LogLevel.Information, Message = "Cache updated on {Silo} for grain {GrainKey} -> {Value}")]
private partial void LogCacheUpdatedLog(string silo, string grainKey, string? value);
[LoggerMessage(LogLevel.Information, Message = "Subscribed cache grain service on {Silo}")]
private partial void LogSubscribed(string silo);
}
Grain key constant
public static class DataGrainConstants
{
public const string GrainKey = "global-cache";
}
Registration
siloBuilder.AddGrainService<CacheGrainService>();
siloBuilder.Services.AddSingleton<DataCache>();
siloBuilder.Services.AddSingleton<IDataCache>(services => services.GetRequiredService<DataCache>());
Validation checklist
- Initial subscription: each silo cache sees the initial
null value for the shared grain key, or whatever makes sense for the cache value type.
- Update fan-out:
UpdateValue notifies all per-silo services.
- Deactivation survival: deactivate the data grain and verify updates still reach services without resubscription.
- Silo crash tolerance: stop a non-hosting silo; updates still flow to remaining services.
Test setup (replicate for validation)
Use an in-process TestCluster with two silos, register the grain service, and assert that updates survive deactivation and silo loss.
Cluster fixture
public sealed class ClusterFixture
{
public TestCluster Cluster { get; private set; } = null!;
public IGrainFactory GrainFactory => Cluster.GrainFactory;
public async ValueTask<ClusterFixture> InitializeAsync()
{
var builder = new TestClusterBuilder();
builder.Options.InitialSilosCount = 2;
builder.AddSiloBuilderConfigurator<SiloConfigurator>();
Cluster = builder.Build();
await Cluster.DeployAsync();
return this;
}
public async ValueTask DisposeAsync()
{
await Cluster.StopAllSilosAsync();
await Cluster.DisposeAsync();
}
public IEnumerable<T> GetServiceFromActiveSilos<T>()
{
return Cluster
.GetActiveSilos()
.SelectMany(siloHandle => Cluster
.GetSiloServiceProvider(siloHandle.SiloAddress)
.GetService<IEnumerable<T>>() ?? Enumerable.Empty<T>());
}
public async ValueTask<SiloAddress> GetHostingSiloAsync(IGrain grain)
{
var grainId = grain.GetGrainId();
var managementGrain = GrainFactory.GetGrain<IManagementGrain>(0);
var activations = await managementGrain.GetDetailedGrainStatistics();
var grainSiloAddress = activations.FirstOrDefault(stat => stat.GrainId == grainId)?.SiloAddress;
Assert.NotNull(grainSiloAddress);
return grainSiloAddress;
}
private sealed class SiloConfigurator : ISiloConfigurator
{
public void Configure(ISiloBuilder siloBuilder)
{
siloBuilder.AddMemoryGrainStorageAsDefault();
siloBuilder.AddGrainService<CacheGrainService>();
siloBuilder.Services.AddSingleton<DataCache>();
siloBuilder.Services.AddSingleton<IDataCache>(services => services.GetRequiredService<DataCache>());
}
}
}
Tests
public sealed class CacheGrainServiceSubscriptionTests
{
[Fact]
public async Task GrainService_subscribing_to_data_grain_and_receives_notifications()
{
await using var fixture = await new ClusterFixture().InitializeAsync();
var dataGrain = fixture.GrainFactory.GetGrain<IDataGrain>(DataGrainConstants.GrainKey);
var siloDataCaches = fixture.GetServiceFromActiveSilos<IDataCache>();
await Assert.AllAsync(siloDataCaches, async cache =>
Assert.Null(await cache.GetValue(DataGrainConstants.GrainKey)));
await dataGrain.UpdateValue("v1");
await Assert.AllAsync(siloDataCaches, async cache =>
Assert.Equal("v1", await cache.GetValue(DataGrainConstants.GrainKey)));
}
[Fact]
public async Task GrainService_subscriptions_survive_data_grain_deactivation()
{
await using var fixture = await new ClusterFixture().InitializeAsync();
var dataGrain = fixture.GrainFactory.GetGrain<IDataGrain>(DataGrainConstants.GrainKey);
var siloDataCaches = fixture.GetServiceFromActiveSilos<IDataCache>();
await dataGrain.UpdateValue("v1");
await dataGrain.Cast<IGrainManagementExtension>().DeactivateOnIdle();
await dataGrain.UpdateValue("v2");
await Assert.AllAsync(siloDataCaches, async cache =>
Assert.Equal("v2", await cache.GetValue(DataGrainConstants.GrainKey)));
}
[Fact]
public async Task Subscriptions_and_notifications_handle_silo_crashes()
{
await using var fixture = await new ClusterFixture().InitializeAsync();
var dataGrain = fixture.GrainFactory.GetGrain<IDataGrain>(DataGrainConstants.GrainKey);
var dataGrainSiloAddress = await fixture.GetHostingSiloAsync(dataGrain);
var silo = fixture.Cluster.GetActiveSilos().First(x => !x.SiloAddress.Equals(dataGrainSiloAddress));
await silo.StopSiloAsync(stopGracefully: false);
await dataGrain.UpdateValue("v1");
Assert.Equal("v1", await fixture.GrainFactory
.GetGrain<ICacheDataDependentGrain>(DataGrainConstants.GrainKey)
.GetValue());
}
}
Gotchas
- Use
IGrainObserver for the observer interface so regular grains/POCOs can also subscribe.
- The grain service still needs
IGrainService on its public interface to register with Orleans.
- Persist
GrainId references, not direct observer references.
ObserverManager removes failed observers from its in-memory collection. Catch SiloUnavailableException in the notification handler to also remove them from persisted state.
References