fix(journaling): make snapshot commits and rollback atomic - #10681
Conversation
There was a problem hiding this comment.
Pull request overview
This PR enhances Orleans Journaling to support composable grain-scoped participants and to make snapshot/recovery behavior more atomic and deterministic, while preserving OrleansBinary snapshot wire compatibility.
Changes:
- Introduces
IJournaledGrainParticipant.Initialize()and calls it fromDurableGrainso participant-provided grain-scoped services are materialized before recovery begins. - Adds an explicit state-manager capability to revert pending changes, and fences writes while recovery/revert is incomplete or has failed.
- Fixes OrleansBinary snapshot replay to reset serializer reference scopes at the same boundaries used during snapshot encoding (per dictionary key/value and per collection element), without changing the payload format; adds focused regression/golden-payload tests.
Show a summary per file
| File | Description |
|---|---|
| test/Orleans.Journaling.Tests/StateManagerTests.cs | Updates snapshot expectations (no append during snapshot replace) and adds coverage for snapshot-failure retry atomicity, revert behavior, and write fencing. |
| test/Orleans.Journaling.Tests/OrleansBinaryCommandCodecTests.cs | Adds golden-payload and aliasing tests to validate independent reference scopes during snapshot replay. |
| test/Orleans.Journaling.Tests/JournaledGrainParticipantTests.cs | Adds integration tests validating participant initialization ordering/once-per-activation and failure propagation. |
| test/Orleans.Journaling.Tests/DurableListDirectWriteTests.cs | Updates a test stub to satisfy the expanded IJournaledStateManager contract. |
| test/Orleans.Journaling.Tests/DurableCollectionDirectWriteTests.cs | Updates a test stub to satisfy the expanded IJournaledStateManager contract. |
| test/Benchmarks/Journaling/DurableListJournalBenchmarks.cs | Updates a benchmark stub to satisfy the expanded IJournaledStateManager contract. |
| src/Orleans.Journaling/JournaledStateManager.cs | Implements write fencing, adds RevertPendingChangesAsync, and adjusts snapshot write semantics to treat replace as a single logical commit. |
| src/Orleans.Journaling/IJournaledStateManager.cs | Adds RevertPendingChangesAsync to the public state-manager contract. |
| src/Orleans.Journaling/IJournaledGrainParticipant.cs | Adds the participant initialization contract to enable composable journaling features. |
| src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableSetCommandCodec.cs | Uses independent-value reads for snapshot items to reset reference scopes per element. |
| src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableQueueCommandCodec.cs | Uses independent-value reads for snapshot items to reset reference scopes per element. |
| src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableListCommandCodec.cs | Uses independent-value reads for snapshot items to reset reference scopes per element. |
| src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryDurableDictionaryCommandCodec.cs | Uses independent-value reads for snapshot keys/values to reset reference scopes per key/value. |
| src/Orleans.Journaling/Formats/OrleansBinary/OrleansBinaryCommandCodecHelpers.cs | Adds ReadIndependentValue which resets the serializer session after each independently encoded snapshot value. |
| src/Orleans.Journaling/DurableGrain.cs | Initializes all registered IJournaledGrainParticipant instances before lifecycle recovery. |
| src/api/Orleans.Journaling/Orleans.Journaling.cs | Updates generated public API surface for the new participant and state-manager members. |
Review details
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
- Files reviewed: 16/16 changed files
- Comments generated: 0
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (1)
src/Orleans.Journaling/JournaledStateManager.cs:859
- WriteStateAsync enqueues work items but does not verify that the manager is initialized (or even started) and does not check shutdown cancellation. Since Start() is only invoked from InitializeAsync, calling WriteStateAsync when _state is still Unknown (or after shutdown) can leave the returned task waiting forever because the work loop will never process the queue.
string operation;
lock (_lock)
{
ThrowIfWritesFenced();
var isSnapshot = _migrationSnapshotRequired || _storage.IsCompactionRequested;
- Files reviewed: 16/16 changed files
- Comments generated: 1
- 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.
src/Orleans.Journaling/JournaledStateManager.cs:927
- ThrowIfWritesFenced throws the same "Call RevertPendingChangesAsync" message for both ManagerState.Recovering (recovery in progress) and ManagerState.Fenced (recovery previously failed). In the Recovering case, this guidance is misleading (there may be nothing to "retry" yet) and can confuse callers who raced with an in-progress recovery.
private void ThrowIfWritesFenced()
{
if (_state is ManagerState.Recovering or ManagerState.Fenced)
{
throw new InvalidOperationException(
"Journaled state writes are fenced until recovery completes successfully. Call RevertPendingChangesAsync to retry recovery.");
}
}
- Files reviewed: 16/16 changed files
- Comments generated: 0 new
- Review effort level: Lite
There was a problem hiding this comment.
Review details
Suppressed comments (1)
src/Orleans.Journaling/JournaledStateManager.cs:916
- ThrowIfWriteOperationsUnavailable only checks that the work loop has started (_workLoop != null). If InitializeAsync starts the loop but recovery/initialization fails, _state remains ManagerState.Unknown, and subsequent WriteStateAsync/DeleteStateAsync calls will be allowed to enqueue work even though the manager has not successfully initialized. This contradicts the new “write operations require initialization” contract and can lead to writes being processed against an unrecovered state once recovery later succeeds.
Consider treating ManagerState.Unknown as uninitialized for all write operations.
private void ThrowIfWriteOperationsUnavailable()
{
_shutdownCancellation.Token.ThrowIfCancellationRequested();
if (_workLoop is null)
{
throw new InvalidOperationException("The journaled state manager has not been initialized.");
}
- Files reviewed: 16/16 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/Orleans.Journaling.Tests/StateManagerTests.cs:2174
- Throwing the captured exception using
throw exception;resets the stack trace, which makes diagnosing failures harder. Use ExceptionDispatchInfo to rethrow while preserving the original stack trace (still throwing the same exception instance).
if (NextReplaceException is { } exception)
{
NextReplaceException = null;
throw exception;
}
- Files reviewed: 17/17 changed files
- Comments generated: 1
- 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.
src/Orleans.Journaling/JournaledStateManager.cs:972
- RevertPendingChangesAsync checks _shutdownCancellation.Token before taking _lock, so StopAsync can cancel between that check and the enqueue/state transition. This can allow enqueuing a recovery work item and setting _state=Recovering after shutdown has begun, which is inconsistent with the atomic shutdown guards used by RegisterObserver/WriteStateAsync/DeleteStateAsync.
cancellationToken.ThrowIfCancellationRequested();
_shutdownCancellation.Token.ThrowIfCancellationRequested();
Task pendingRecovery;
bool didEnqueue;
lock (_lock)
{
if (_state is ManagerState.Unknown)
{
throw new InvalidOperationException("The journaled state manager has not been initialized.");
}
pendingRecovery = EnqueueOrGetPendingWorkItem<RevertPendingChangesWorkItem>(out didEnqueue);
_state = ManagerState.Recovering;
}
- Files reviewed: 14/14 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.
src/Orleans.Journaling/DurableGrain.cs:18
- DurableGrain initializes participants while iterating the resolved service collection. If the DI container were to resolve the IEnumerable lazily, a later participant constructor exception could occur after earlier participants have already been initialized, leaving partial side effects before activation fails. Materializing the participant list first ensures either all participants construct successfully before any Initialize() runs, or none do.
foreach (var feature in ServiceProvider.GetServices<IJournaledGrainParticipant>())
{
feature.Initialize();
}
- Files reviewed: 15/15 changed files
- Comments generated: 0 new
- Review effort level: Lite
Problem
Snapshot replacement currently flushes pending journal bytes as a separate append before replacing storage. A failed replacement can therefore expose an intermediate durable boundary, consume pending data too early, and notify state more than once. Callers also lack an explicit way to discard provisional mutations and reload the last durable state.
Solution
RevertPendingChangesAsyncto restore the last durable state explicitly.Rationale
These changes give volatile and durable state a single commit boundary and make rollback outcomes explicit. The implementation remains contained within Orleans Journaling and does not alter journal wire formats.
The independent OrleansBinary snapshot reference-scope replay fix was split into #10800.