diff --git a/docs/site/src/content/docs/grains/journaling/runtime-behavior.md b/docs/site/src/content/docs/grains/journaling/runtime-behavior.md index dddcba9ea91..a1d1d278664 100644 --- a/docs/site/src/content/docs/grains/journaling/runtime-behavior.md +++ b/docs/site/src/content/docs/grains/journaling/runtime-behavior.md @@ -62,13 +62,11 @@ Concurrent calls made while the same kind of write is queued can share that queu ## Safe-to-commit staging All interleaved callers share the manager's pending journal. Prepare fallible work, external acknowledgements, -and proposed output in operation-local data. After establishing that an outcome is safe to commit, apply its -mutations to durable state and initiate a write. Coordinate that transition with other interleaved operations -which can affect the same decision. Any caller's write can include staged mutations from other calls. - -If an application error occurs after staging and makes those mutations unsafe to commit, end the activation's -use of the manager and request deactivation. In-flight methods can retain local decisions and references -across awaits; a fresh activation reconstructs both application and durable state together. +and proposed output in operation-local data. After the final preparation await, check the relevant +preconditions and apply the complete safe-to-commit update synchronously, then request an ordinary write. +Orleans executes that synchronous block on a single activation thread. Another grain turn can run when +the operation awaits, so keep shared state safe to commit at each await. Any caller's write can include +staged mutations from other calls. ## Consistency and competing writers @@ -80,7 +78,7 @@ Design commands to tolerate retries at the application boundary. Use operation i ## Storage failures -A failed append, snapshot replacement, delete, or initialization permanently fences that manager instance. +A failed append, snapshot replacement, or delete permanently fences that manager instance. Queued operations fault, and later write, delete, registration, and initialization requests fail explicitly. Existing in-memory state remains available to in-flight calls until deactivation completes. The grain runtime starts deactivation as part of handling the failure. @@ -99,8 +97,44 @@ assigned lifetimes. Cancelling a write's cancellation token stops the caller's wait. An already queued write continues to its storage outcome, so the caller reconciles that outcome before retrying the command. -An initialization failure preserves stored data for diagnosis. Restore the required format/codec registration -or repair the backing data before creating a fresh manager or retrying activation. +An initialization failure reports its error to that attempt's callers and leaves the manager uninitialized. +The caller can retry after a transient +failure or after restoring the required format, codec, or backing data. Each attempt resets recovery +bookkeeping and replays the journal from the beginning using the same registered state machines. +Concurrent callers share the active attempt; writes and deletion become available after initialization +succeeds. State registration stays closed once initialization has begun. +One work-loop task owns recovery attempts and subsequent journal work for the manager's lifetime. +After a failed attempt it waits for an explicit initialization request before retrying. + +Owner shutdown which cancels initial recovery cancels all initialization waiters and leaves the manager +stopped. Disposal waits for the owned read to finish before releasing journal resources. Cancelling an +individual initialization caller's token ends only its wait; owned recovery continues for other callers. + +## Custom state lifecycle + +Custom implementations share the manager's single logical execution thread. +Supply command codecs as constructor dependencies. The registration factory selects codecs keyed by the +same write-format key used to configure the journal owner. Activation-owned state factories resolve those +dependencies from the activation's services; standalone callers supply codecs with the appropriate lifetime. +The codec is available when the state is constructed, including for an empty journal. +During replay, selects the codec +for each entry's stored format. + +States synchronously encode their pending changes through +or their current contents through . +After storage acknowledges captured bytes, +performs durable-completion bookkeeping. A zero-byte write completes without this callback. + +The journal owner keeps feature operations quiescent through deletion's storage and reset outcome, +including when a caller cancels its wait. Successful deletion calls +before completing deletion waiters. + +The manager records the first write or delete failure, fences further persistence, faults current +and queued manager waiters, and requests grain deactivation. Features observe their write failures and +complete their own operation waiters and resource cleanup through their operation and lifecycle ownership. +Standalone callers own that cleanup explicitly. A previously captured write retains its actual storage +outcome and acknowledgement bookkeeping. Owner-canceled initial recovery and idle shutdown complete +through normal shutdown; admitted write/delete storage cancellation remains terminal. ## Compaction diff --git a/src/Orleans.Journaling/IJournaledStateManager.cs b/src/Orleans.Journaling/IJournaledStateManager.cs index d619a8f37f7..b6363a1b4cc 100644 --- a/src/Orleans.Journaling/IJournaledStateManager.cs +++ b/src/Orleans.Journaling/IJournaledStateManager.cs @@ -19,7 +19,11 @@ public interface IJournaledStateManager : IAsyncDisposable /// Initializes the state manager by replaying its journal. /// /// - /// A failed initialization permanently fences this instance. Recover by creating a new manager and new state instances. + /// A recovery failure fails the current initialization attempt and leaves this instance uninitialized. + /// A subsequent call retries recovery from the beginning, resetting and replaying the registered state machines. + /// Writes become available after initialization succeeds. A manager fenced by a persistence failure requires a new instance. + /// Owner shutdown which cancels recovery cancels initialization and leaves this instance stopped. + /// Cancelling the caller's token ends only that caller's wait while owned recovery continues. /// /// The cancellation token. /// A which represents the operation. @@ -58,7 +62,8 @@ public interface IJournaledStateManager : IAsyncDisposable /// Resets this instance, removing any persistent state. /// /// - /// Quiesce other operations before deleting state: deletion resets every registered state machine. + /// The caller keeps other operations quiescent through completion: deletion resets every registered state machine. + /// Cancellation ends the caller's wait; an already queued deletion continues to its storage and reset outcome. /// A failed deletion permanently fences the manager and requests deactivation of its owning grain. /// /// The cancellation token. diff --git a/src/Orleans.Journaling/IStateMachine.cs b/src/Orleans.Journaling/IStateMachine.cs index 39f824a6db7..bf3fd61204e 100644 --- a/src/Orleans.Journaling/IStateMachine.cs +++ b/src/Orleans.Journaling/IStateMachine.cs @@ -25,14 +25,15 @@ namespace Orleans.Journaling; /// rather than treating in-memory mutations as durable. /// /// -/// A failed journal operation permanently fences the manager and requests grain deactivation. -/// A new manager initializes new state instances by calling and replaying durable entries. +/// A failed write or delete permanently fences the manager and requests grain deactivation. +/// Recovery calls before replaying durable entries, including when initialization is retried. /// /// /// /// Application code prepares fallible work in operation-local data and stages only mutations which are /// safe to commit. Staged mutations are shared by all interleaved callers using the same manager. -/// Storage acknowledgement establishes durability; recovery takes place in a fresh manager and state instances. +/// Storage acknowledgement establishes durability. A new activation creates a fresh manager and state instances; +/// retrying failed initialization resets and replays the existing instances. /// /// public interface IStateMachine diff --git a/src/Orleans.Journaling/JournaledStateManager.cs b/src/Orleans.Journaling/JournaledStateManager.cs index 720e6789d5e..cbf6be752bc 100644 --- a/src/Orleans.Journaling/JournaledStateManager.cs +++ b/src/Orleans.Journaling/JournaledStateManager.cs @@ -156,11 +156,11 @@ public void RegisterStateMachine(string name, IStateMachine stateMachine) public async ValueTask InitializeAsync(CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); - _shutdownCancellation.Token.ThrowIfCancellationRequested(); Task task; bool didEnqueue; lock (_lock) { + _shutdownCancellation.Token.ThrowIfCancellationRequested(); ThrowIfFenced(); if (_workLoop is null) { @@ -187,14 +187,50 @@ private Task Start() private async Task WorkLoop() { await Task.CompletedTask.ConfigureAwait(ConfigureAwaitOptions.ContinueOnCapturedContext | ConfigureAwaitOptions.ForceYielding); - try - { - await RecoverAsync(_shutdownCancellation.Token).ConfigureAwait(true); - } - catch (Exception exception) + while (!_shutdownCancellation.IsCancellationRequested) { - Fence(exception); - return; + try + { + await RecoverAsync(_shutdownCancellation.Token).ConfigureAwait(true); + _workSignal.Signal(); + break; + } + catch (OperationCanceledException) when (_shutdownCancellation.IsCancellationRequested) + { + return; + } + catch (Exception exception) + { + try + { + LogErrorProcessingWorkItems(_shared.Logger, exception); + } + finally + { + lock (_lock) + { + FaultQueuedWorkItemsUnderLock(exception); + } + } + } + + // Signals can remain from the failed attempt. Retry only for newly queued initialization work. + while (true) + { + await _workSignal.WaitAsync().ConfigureAwait(true); + if (_shutdownCancellation.IsCancellationRequested) + { + return; + } + + lock (_lock) + { + if (_workQueue.Count > 0) + { + break; + } + } + } } while (!_shutdownCancellation.Token.IsCancellationRequested) @@ -548,6 +584,10 @@ private async Task WorkLoop() } } } + catch (OperationCanceledException) when (_shutdownCancellation.IsCancellationRequested) + { + return; + } catch (Exception exception) { Fence(exception); @@ -560,6 +600,11 @@ private void Fence(Exception exception) { lock (_lock) { + if (_state is ManagerState.Fenced) + { + return; + } + _state = ManagerState.Fenced; _failure = exception; } diff --git a/src/Orleans.Journaling/README.md b/src/Orleans.Journaling/README.md index 549abe06179..671bbda74e9 100644 --- a/src/Orleans.Journaling/README.md +++ b/src/Orleans.Journaling/README.md @@ -183,8 +183,9 @@ provider and requested state name. - `OnRecoveryCompleted`: finish reconstruction before application use. - `OnWriteCompleted`: publish effects that depend on storage acknowledgement. -The recovery model uses fresh instances and replay. `JournalReplayContext.ResolveStateMachine` -routes entries to the state machine for their stream. +Recovery resets state machines and replays durable entries. A failed initialization can be retried on the +same manager and registered states. `JournalReplayContext.ResolveStateMachine` routes entries to the +state machine for their stream. `IJournaledStateManager` is independent of the grain-facing `IDurableStateManager` and extends `IAsyncDisposable`. Its owner API provides `RegisterStateMachine`, `TryGetStateMachine`, @@ -255,13 +256,18 @@ durable state and await `WriteStateAsync`. One acknowledgement covers the manage batch, including changes staged by interleaved callers. Applications are responsible for sequencing that transition with other interleaved operations and for making uncertain-outcome retries idempotent. -A failed journal operation permanently fences the manager, faults queued operations, and requests +A failed write or delete permanently fences the manager, faults queued operations, and requests deactivation of the associated grain. In-flight calls retain their existing in-memory state while subsequent state-manager operations fail explicitly. A new activation recovers the actual durable outcome. For a manager created through `IJournaledStateManagerFactory`, dispose the failed instance and create another manager for the same `JournalId`, explicitly constructing and registering fresh state components before initialization and retiring the old components and dependencies according to their assigned lifetimes. +An initialization failure reports its error to the attempt's callers and leaves the manager uninitialized. +Call `InitializeAsync` again to retry from the beginning using the existing reset/replay contract. +Concurrent callers share the active attempt, and writes become available after recovery succeeds. +State registration stays closed after initialization first begins. + Cancelling a caller's wait leaves an already queued write running. Observe durability through write acknowledgement or a fresh activation before deciding whether to retry an application command. diff --git a/test/Orleans.Journaling.Tests/KeyedJournalingRegistrationTests.cs b/test/Orleans.Journaling.Tests/KeyedJournalingRegistrationTests.cs index 3f72fdc0e4f..6c994b1f2b0 100644 --- a/test/Orleans.Journaling.Tests/KeyedJournalingRegistrationTests.cs +++ b/test/Orleans.Journaling.Tests/KeyedJournalingRegistrationTests.cs @@ -382,6 +382,126 @@ public void DurableService_ResolvesCommandCodecFromJournalFormatKey() scope.ServiceProvider.GetRequiredService().ObservableLifecycle).Subscriptions); } + [Fact] + public async Task StateConstruction_InjectsActivationScopedCodecBeforeRecovery() + { + var builder = CreateNamedProviderBuilder(); + builder.AddVolatileJournalStorage(); + builder.Services.Configure(options => + options.JournalFormatKey = OrleansBinaryJournalFormat.JournalFormatKey); + builder.Services.AddScoped(services => + { + var context = Substitute.For(); + context.GrainId.Returns(GrainId.Create("codec-scope", Guid.NewGuid().ToString("N"))); + context.ActivationServices.Returns(services); + context.ObservableLifecycle.Returns(new CompositionTestLifecycle()); + return context; + }); + builder.Services.AddKeyedScoped>(OrleansBinaryJournalFormat.JournalFormatKey, + static (services, _) => new OrleansBinaryDurableValueCommandCodec( + services.GetRequiredService().GetCodec(), + services.GetRequiredService())); + builder.Services.AddStateMachine(static (services, _) => + new CodecState(services.GetRequiredKeyedService>(OrleansBinaryJournalFormat.JournalFormatKey))); + await using var services = builder.Services.BuildServiceProvider(validateScopes: true); + await using var first = services.CreateAsyncScope(); + await using var second = services.CreateAsyncScope(); + var owner = first.ServiceProvider.GetRequiredService(); + var manager = first.ServiceProvider.GetRequiredService(); + Assert.Same(owner, manager); + var state = manager.GetOrAddState("value"); + var codec = state.Codec; + Assert.Same(first.ServiceProvider.GetRequiredKeyedService>(OrleansBinaryJournalFormat.JournalFormatKey), codec); + Assert.NotSame(codec, second.ServiceProvider.GetRequiredService() + .GetOrAddState("value").Codec); + Assert.Throws(() => + services.GetRequiredKeyedService>(OrleansBinaryJournalFormat.JournalFormatKey)); + Assert.Same(state, manager.GetOrAddState("value")); + Assert.Same(state, first.ServiceProvider.GetRequiredKeyedService("value")); + Assert.True(owner.TryGetStateMachine("value", out var registered)); + Assert.Same(state, registered); + var lifecycle = Assert.IsType(first.ServiceProvider.GetRequiredService().ObservableLifecycle); + Assert.Equal(1, lifecycle.Subscriptions); + await lifecycle.OnStart(TestContext.Current.CancellationToken); + Assert.Equal(0, state.Value); + state.Value = 42; + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Equal(0, owner.PendingWriteByteCount); + await lifecycle.OnStop(TestContext.Current.CancellationToken); + } + + [Fact] + public async Task StateConstruction_InjectsNamedFormatCodecOnEmptyJournal() + { + var builder = CreateNamedProviderBuilder(); + builder.AddVolatileJournalStorage(); + var customStorage = new VolatileJournalStorage(CustomFormatKey); + builder.Services.AddKeyedSingleton(CustomFormatKey, static (services, _) => + new NamedBinaryJournalFormat(services.GetRequiredService())); + builder.Services.AddKeyedSingleton(typeof(IDurableDictionaryCommandCodec<,>), CustomFormatKey, + typeof(OrleansBinaryDurableDictionaryCommandCodec<,>)); + builder.Services.AddKeyedScoped>(CustomFormatKey, static (services, _) => + new OrleansBinaryDurableValueCommandCodec( + services.GetRequiredService().GetCodec(), + services.GetRequiredService())); + builder.Services.AddKeyedSingleton("custom", (services, _) => + new JournaledStateManagerFactory( + new JournaledStateManagerShared( + services.GetRequiredService>(), + Options.Create(new JournaledStateManagerOptions { JournalFormatKey = CustomFormatKey }), + TimeProvider.System, + services), + new TestJournalStorageProvider(customStorage))); + await using var services = builder.Services.BuildServiceProvider(validateScopes: true); + await using var dependencies = services.CreateAsyncScope(); + var factory = services.GetRequiredKeyedService("custom"); + await using var defaultManager = services.GetRequiredService().CreateStandalone(new JournalId("default")); + await using var customManager = factory.CreateStandalone(new JournalId("custom")); + IJournaledStateManager delegating = new DelegatingStateManager(customManager); + var codec = dependencies.ServiceProvider.GetRequiredKeyedService>(CustomFormatKey); + Assert.Throws(() => + services.GetRequiredKeyedService>(CustomFormatKey)); + var defaultCodec = services.GetRequiredKeyedService>(JsonLinesJournalFormat.JournalFormatKey); + Assert.NotSame(defaultCodec, codec); + var state = new CodecState(codec); + Assert.Same(codec, state.Codec); + delegating.RegisterStateMachine("value", state); + var defaultState = new CodecState(defaultCodec); + defaultManager.RegisterStateMachine("value", defaultState); + await defaultManager.InitializeAsync(TestContext.Current.CancellationToken); + await delegating.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(0, state.Value); + Assert.Empty(customStorage.Segments); + state.Value = 42; + defaultState.Value = 7; + await delegating.WriteStateAsync(TestContext.Current.CancellationToken); + await defaultManager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Single(customStorage.Segments); + await using var recovered = factory.CreateStandalone(new JournalId("custom")); + var recoveredState = new CodecState(codec); + recovered.RegisterStateMachine("value", recoveredState); + await recovered.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(42, recoveredState.Value); + + await using var recoveredDefault = services.GetRequiredService().CreateStandalone(new JournalId("default")); + var recoveredDefaultState = new CodecState(defaultCodec); + recoveredDefault.RegisterStateMachine("value", recoveredDefaultState); + await recoveredDefault.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(7, recoveredDefaultState.Value); + } + + [Fact] + public void StateConstruction_MissingFormatCodecFails() + { + var builder = CreateNamedProviderBuilder(); + builder.AddJournaling(); + using var services = builder.Services.BuildServiceProvider(); + Assert.NotNull(services.GetRequiredKeyedService>(JsonLinesJournalFormat.JournalFormatKey)); + var exception = Assert.Throws(() => + new CodecState(services.GetRequiredKeyedService>(CustomFormatKey))); + Assert.Contains(nameof(IDurableValueCommandCodec), exception.Message); + } + [Fact] public async Task StateManagerFactory_CreatesManagerForJournalId() { @@ -421,6 +541,36 @@ private sealed class TestJournalStorageProvider(IJournalStorage storage) : IJour public IJournalStorage CreateStorage(JournalId journalId) => storage; } + private sealed class NamedBinaryJournalFormat(IJournalFormat inner) : IJournalFormat + { + public string FormatKey => CustomFormatKey; + public string? MimeType => inner.MimeType; + public JournalBufferWriter CreateWriter() => inner.CreateWriter(); + public void Replay(JournalBufferReader input, JournalReplayContext context) => inner.Replay(input, context); + } + + private sealed class CodecState(IDurableValueCommandCodec codec) : IStateMachine, IDurableValueCommandHandler + { + public IDurableValueCommandCodec Codec { get; } = codec; + public int Value { get; set; } + public void Reset(JournalStreamWriter writer) => Value = 0; + public void WritePendingEntries(JournalStreamWriter writer) => Codec.WriteSet(Value, writer); + public void WriteSnapshot(JournalStreamWriter writer) => Codec.WriteSet(Value, writer); + public void ReplayEntry(JournalEntry entry, JournalReplayContext context) => + context.GetRequiredCommandCodec(entry.FormatKey, Codec).Apply(entry.Reader, this); + public void ApplySet(int value) => Value = value; + } + + private sealed class DelegatingStateManager(IJournaledStateManager inner) : IJournaledStateManager + { + public ValueTask InitializeAsync(CancellationToken cancellationToken) => inner.InitializeAsync(cancellationToken); + public void RegisterStateMachine(string name, IStateMachine state) => inner.RegisterStateMachine(name, state); + public bool TryGetStateMachine(string name, [System.Diagnostics.CodeAnalysis.NotNullWhen(true)] out IStateMachine? state) + => inner.TryGetStateMachine(name, out state); + public ValueTask WriteStateAsync(CancellationToken cancellationToken) => inner.WriteStateAsync(cancellationToken); + public ValueTask DeleteStateAsync(CancellationToken cancellationToken) => inner.DeleteStateAsync(cancellationToken); + } + private sealed class LifecycleJournalStorageProvider : IJournalStorageProvider, IJournalStorageCatalog, ILifecycleParticipant { private readonly VolatileJournalStorageProvider _storage = new(); diff --git a/test/Orleans.Journaling.Tests/StateManagerLifecycleTests.cs b/test/Orleans.Journaling.Tests/StateManagerLifecycleTests.cs new file mode 100644 index 00000000000..abe57f4be4d --- /dev/null +++ b/test/Orleans.Journaling.Tests/StateManagerLifecycleTests.cs @@ -0,0 +1,773 @@ +using System.Reflection; +using System.Threading.Channels; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; +using NSubstitute; +using Xunit; + +namespace Orleans.Journaling.Tests; + +public partial class StateManagerTests +{ + [Fact] + public async Task StateMachineProtocol_WritesAndResetsAfterDeletion() + { + var storage = new CapturingStorage(); + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(1, state.ResetCount); + Assert.Equal(1, state.RecoveryCompletedCount); + + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Single(storage.Appends); + Assert.Equal(1, state.CaptureCount); + Assert.Equal(1, state.WriteCompletedCount); + + await manager.DeleteStateAsync(TestContext.Current.CancellationToken); + Assert.Equal(1, storage.DeleteCount); + Assert.Equal(2, state.ResetCount); + Assert.Equal(1, state.WriteCompletedCount); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task Delete_ResetsStatesBeforeCompletion(bool cancelWait) + { + var storage = new BlockingDeleteStorage(); + await using var manager = CreateTestSystem(storage).Manager; + var first = new LifecycleState(); + var second = new LifecycleState(); + manager.RegisterStateMachine("first", first); + manager.RegisterStateMachine("second", second); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + using var caller = new CancellationTokenSource(); + var deletion = manager.DeleteStateAsync(caller.Token).AsTask(); + await WaitFor(storage.FirstDeleteStarted.Task); + Assert.Equal(1, first.ResetCount); + Assert.Equal(1, second.ResetCount); + Assert.False(deletion.IsCompleted); + if (cancelWait) + { + caller.Cancel(); + await Assert.ThrowsAnyAsync(() => deletion); + } + + var resets = new List(); + var resetCompleted = NewSignal(); + first.ResetAction = () => + { + if (!cancelWait) Assert.False(deletion.IsCompleted); + resets.Add("first"); + }; + second.ResetAction = () => + { + if (!cancelWait) Assert.False(deletion.IsCompleted); + resets.Add("second"); + resetCompleted.SetResult(); + }; + storage.AllowFirstDelete.SetResult(); + await WaitFor(resetCompleted.Task); + if (!cancelWait) await WaitFor(deletion); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(["first", "second"], resets); + Assert.Equal(2, first.ResetCount); + Assert.Equal(2, second.ResetCount); + Assert.Equal(1, storage.DeleteCount); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task OperationLocalPreparation_ThenSynchronousUpdatesPersistTogether(bool snapshot) + { + var context = new QueuedSynchronizationContext(); + await context.Run(async () => + { + var storage = new CapturingStorage { IsCompactionRequested = snapshot }; + await using var manager = CreateTestSystem(storage).Manager; + var dictionary = new DurableDictionary("items", manager, CreateDictionaryCodec()); + var total = new DurableValue("total", manager, CreateValueCodec()); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + dictionary.Add("existing", 1); + total.Value = 1; + + var preparing = NewSignal(); + var release = NewSignal(); + var update = PrepareAndApply(); + await WaitFor(preparing.Task); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.False(update.IsCompleted); + Assert.Single(dictionary); + Assert.Equal(1, total.Value); + await AssertRecovered(1); + + release.SetResult(); + await WaitFor(update); + Assert.Equal(2, dictionary.Count); + Assert.Equal(43, total.Value); + await AssertRecovered(43); + + async Task PrepareAndApply() + { + var proposed = (Key: "prepared", Value: 42); + preparing.SetResult(); + await WaitFor(release.Task); + Assert.False(dictionary.ContainsKey(proposed.Key)); + dictionary.Add(proposed.Key, proposed.Value); + total.Value += proposed.Value; + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + } + + async Task AssertRecovered(int expectedTotal) + { + await using var recovered = CreateTestSystem(storage).Manager; + var recoveredItems = new DurableDictionary("items", recovered, CreateDictionaryCodec()); + var recoveredTotal = new DurableValue("total", recovered, CreateValueCodec()); + await recovered.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(expectedTotal, recoveredTotal.Value); + Assert.Equal(expectedTotal == 1 ? 1 : 2, recoveredItems.Count); + Assert.Equal(1, recoveredItems["existing"]); + if (expectedTotal != 1) Assert.Equal(42, recoveredItems["prepared"]); + } + }); + } + + [Fact] + public async Task ZeroByteWrite_DoesNotRepeatAcknowledgement() + { + var storage = new CapturingStorage(); + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState { EmitEntry = false }; + manager.RegisterStateMachine("state", state); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + // Initialization records the state name even when the state emits no entries. + Assert.True(manager.PendingWriteByteCount > 0); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + var directoryEntry = Assert.Single(ReadBinaryEntries(Assert.Single(storage.Appends))); + Assert.Equal(0u, directoryEntry.StreamId.Value); + Assert.Equal(0, manager.PendingWriteByteCount); + Assert.Equal(1, state.WriteCompletedCount); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Equal(2, state.CaptureCount); + Assert.Equal(1, state.WriteCompletedCount); + Assert.Single(storage.Appends); + } + + [Fact] + public async Task LaterStorageFailure_PreservesAlreadyCapturedWriteAcknowledgement() + { + var storage = new CapturingStorage { BlockNextAppend = true }; + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + var capturedWrite = manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask(); + await WaitFor(storage.BlockedAppendStarted.Task); + var expected = new IOException("Failure of the next append."); + storage.NextAppendException = expected; + var failingWrite = manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask(); + Assert.Equal(0, state.WriteCompletedCount); + storage.ReleaseAppend.SetResult(); + await WaitFor(capturedWrite); + Assert.Same(expected, await Record.ExceptionAsync(() => WaitFor(failingWrite))); + Assert.Equal(2, state.CaptureCount); + Assert.Equal(1, state.WriteCompletedCount); + Assert.Single(storage.Appends); + var rejected = await Assert.ThrowsAsync(() => manager.WriteStateAsync(CancellationToken.None).AsTask()); + Assert.Same(expected, rejected.InnerException); + } + + [Fact] + public async Task DeactivationFailure_PreservesOriginalStorageFailure() + { + var expected = new IOException("Original storage failure."); + var secondary = new InvalidOperationException("Deactivation failed."); + var storage = new CapturingStorage { BlockNextAppend = true, NextAppendException = expected }; + var provider = Substitute.For(); + var context = Substitute.For(); + context.GrainId.Returns(GrainId.Create("test-grain", "deactivation-failure")); + context.ActivationServices.Returns(ServiceProvider); + provider.CreateStorage(JournalId.FromGrainId(context.GrainId)).Returns(storage); + context.When(value => value.Deactivate(Arg.Any(), Arg.Any())) + .Do(_ => throw secondary); + var shared = new JournaledStateManagerShared( + ServiceProvider.GetRequiredService>(), + Options.Create(ManagerOptions), TimeProvider.System, ServiceProvider); + await using var manager = new JournaledStateManager(shared, provider, context); + manager.RegisterStateMachine("state", new LifecycleState()); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + var current = manager.WriteStateAsync(CancellationToken.None).AsTask(); + await WaitFor(storage.BlockedAppendStarted.Task); + var queued = manager.WriteStateAsync(CancellationToken.None).AsTask(); + storage.ReleaseAppend.SetResult(); + Assert.Same(expected, await Record.ExceptionAsync(() => WaitFor(current))); + Assert.Same(expected, await Record.ExceptionAsync(() => WaitFor(queued))); + var late = await Assert.ThrowsAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask()); + Assert.Same(expected, late.InnerException); + context.Received(1).Deactivate( + Arg.Is(reason => ReferenceEquals(reason.Exception, expected)), Arg.Any()); + } + + [Fact] + public async Task InitializationFailure_ReportsOriginalCauseAndAllowsRetry() + { + var expected = new IOException("Recovery failed."); + var storage = new CapturingStorage { NextReadException = expected }; + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + Assert.Same(expected, await Record.ExceptionAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask())); + Assert.Equal(0, state.RecoveryCompletedCount); + Assert.True(manager.TryGetStateMachine("state", out var registered)); + Assert.Same(state, registered); + var write = await Assert.ThrowsAsync(() => manager.WriteStateAsync(CancellationToken.None).AsTask()); + var delete = await Assert.ThrowsAsync(() => manager.DeleteStateAsync(CancellationToken.None).AsTask()); + Assert.Contains("not been initialized", write.Message); + Assert.Contains("not been initialized", delete.Message); + var lateRegistration = Assert.Throws(() => manager.RegisterStateMachine("late", new LifecycleState())); + Assert.Contains("initialization has begun", lateRegistration.Message); + + await manager.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(1, state.RecoveryCompletedCount); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Single(storage.Appends); + Assert.Equal(1, state.WriteCompletedCount); + } + + [Fact] + public async Task RecoveryRetries_ReuseSingleWorkLoopTask() + { + var storage = new MutableReadStorage(3, [1, 2, 3], [1, 2, 3], CreatePersistedValueBytes("value", 42)); + await using var manager = CreateTestSystem(storage).Manager; + var value = new DurableValue("value", manager, CreateValueCodec()); + var initial = manager.InitializeAsync(CancellationToken.None).AsTask(); + var workLoop = GetWorkLoop(); + await Assert.ThrowsAsync(() => WaitFor(initial)); + Assert.Equal(1, storage.ReadCount); + + var retry = manager.InitializeAsync(CancellationToken.None).AsTask(); + await Assert.ThrowsAsync(() => WaitFor(retry)); + Assert.Same(workLoop, GetWorkLoop()); + Assert.False(workLoop.IsCompleted); + Assert.Equal(2, storage.ReadCount); + + const int callerCount = 8; + var start = NewSignal(); + var allEnqueued = NewSignal(); + var enqueued = 0; + var callers = Enumerable.Range(0, callerCount).Select(_ => Task.Run(async () => + { + await WaitFor(start.Task); + var initialization = manager.InitializeAsync(CancellationToken.None).AsTask(); + if (Interlocked.Increment(ref enqueued) == callerCount) allEnqueued.SetResult(); + await initialization; + }, TestContext.Current.CancellationToken)).ToArray(); + start.SetResult(); + await WaitFor(allEnqueued.Task); + await WaitFor(storage.BlockedReadStarted.Task); + Assert.Same(workLoop, GetWorkLoop()); + Assert.False(workLoop.IsCompleted); + Assert.Equal(3, storage.ReadCount); + Assert.All(callers, caller => Assert.False(caller.IsCompleted)); + + storage.AllowBlockedRead.SetResult(); + await WaitFor(Task.WhenAll(callers)); + Assert.Equal(42, value.Value); + value.Value = 43; + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Same(workLoop, GetWorkLoop()); + Assert.False(workLoop.IsCompleted); + Assert.Equal(3, storage.ReadCount); + + await WaitFor(manager.DisposeAsync().AsTask()); + Assert.Same(workLoop, GetWorkLoop()); + Assert.True(workLoop.IsCompletedSuccessfully); + + Task GetWorkLoop() => Assert.IsAssignableFrom( + typeof(JournaledStateManager).GetField("_workLoop", BindingFlags.Instance | BindingFlags.NonPublic)!.GetValue(manager)); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task RecoveryFailure_WaitsForExplicitRetryOrShutdown(bool lifecycleStop) + { + var context = new QueuedSynchronizationContext(); + await context.Run(async () => + { + var storage = new MutableReadStorage([1, 2, 3], []); + var sut = CreateTestSystem(storage); + await using var manager = sut.Manager; + manager.RegisterStateMachine("state", new LifecycleState()); + await Assert.ThrowsAsync(() => lifecycleStop + ? sut.Lifecycle.OnStart(CancellationToken.None) + : manager.InitializeAsync(CancellationToken.None).AsTask()); + await Task.Yield(); + Assert.Equal(1, storage.ReadCount); + Assert.Empty(storage.OperationLog); + Assert.True(manager.TryGetStateMachine("state", out _)); + Assert.Throws(() => manager.RegisterStateMachine("late", new LifecycleState())); + + await WaitFor(lifecycleStop + ? sut.Lifecycle.OnStop(TestContext.Current.CancellationToken) + : manager.DisposeAsync().AsTask()); + Assert.Equal(1, storage.ReadCount); + }); + } + + [Fact] + public async Task RecoveryRetry_CoalescesCallersAndOwnsRead() + { + var storage = new MutableReadStorage(2, [1, 2, 3], CreatePersistedValueBytes("value", 42)); + await using var manager = CreateTestSystem(storage).Manager; + var value = new DurableValue("value", manager, CreateValueCodec()); + await Assert.ThrowsAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask()); + Assert.Equal(1, storage.ReadCount); + + using var caller = new CancellationTokenSource(); + var canceledWaiter = manager.InitializeAsync(caller.Token).AsTask(); + await WaitFor(storage.BlockedReadStarted.Task); + var remainingWaiter = manager.InitializeAsync(CancellationToken.None).AsTask(); + Assert.Equal(2, storage.ReadCount); + var write = await Assert.ThrowsAsync(() => manager.WriteStateAsync(CancellationToken.None).AsTask()); + Assert.Contains("not been initialized", write.Message); + var registration = Assert.Throws(() => manager.RegisterStateMachine("late", new LifecycleState())); + Assert.Contains("initialization has begun", registration.Message); + caller.Cancel(); + await Assert.ThrowsAnyAsync(() => WaitFor(canceledWaiter)); + Assert.False(storage.ReadToken.IsCancellationRequested); + Assert.False(remainingWaiter.IsCompleted); + + storage.AllowBlockedRead.SetResult(); + await WaitFor(remainingWaiter); + Assert.Equal(42, value.Value); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(2, storage.ReadCount); + value.Value = 43; + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Equal(["append"], storage.OperationLog); + } + + [Fact] + public async Task RecoveryCompletionFailure_RetryResetsRegisteredState() + { + var storage = new CapturingStorage(); + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + var expected = new InvalidOperationException("Recovery completion failed."); + state.RecoveryCompletedAction = () => throw expected; + Assert.Same(expected, await Record.ExceptionAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask())); + Assert.Equal(1, state.ResetCount); + Assert.Equal(1, state.RecoveryCompletedCount); + Assert.Empty(storage.Appends); + + state.RecoveryCompletedAction = null; + await manager.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal(2, state.ResetCount); + Assert.Equal(2, state.RecoveryCompletedCount); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Single(storage.Appends); + } + + [Fact] + public async Task InitializationCallerCancellation_LeavesOwnedRecoveryRunning() + { + var storage = new MutableReadStorage(1, Array.Empty()); + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + using var caller = new CancellationTokenSource(); + var canceledWaiter = manager.InitializeAsync(caller.Token).AsTask(); + await WaitFor(storage.BlockedReadStarted.Task); + var remainingWaiter = manager.InitializeAsync(CancellationToken.None).AsTask(); + caller.Cancel(); + await Assert.ThrowsAnyAsync(() => WaitFor(canceledWaiter)); + Assert.False(storage.ReadToken.IsCancellationRequested); + Assert.False(remainingWaiter.IsCompleted); + Assert.Equal(0, state.RecoveryCompletedCount); + + storage.AllowBlockedRead.SetResult(); + await WaitFor(remainingWaiter); + Assert.Equal(1, state.RecoveryCompletedCount); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Equal(["append"], storage.OperationLog); + Assert.Equal(1, state.WriteCompletedCount); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task RecoveryShutdown_CancelsInitializationWithoutCompletingRecovery(bool lifecycleStop) + { + var storage = new MutableReadStorage(1, Array.Empty()); + var sut = CreateTestSystem(storage); + await using var manager = sut.Manager; + var first = new LifecycleState(); + var second = new LifecycleState(); + manager.RegisterStateMachine("first", first); + manager.RegisterStateMachine("second", second); + var startup = lifecycleStop + ? sut.Lifecycle.OnStart(CancellationToken.None) + : manager.InitializeAsync(CancellationToken.None).AsTask(); + await WaitFor(storage.BlockedReadStarted.Task); + var another = manager.InitializeAsync(CancellationToken.None).AsTask(); + Assert.False(startup.IsCompleted); + Assert.False(another.IsCompleted); + + await WaitFor(lifecycleStop + ? sut.Lifecycle.OnStop(TestContext.Current.CancellationToken) + : manager.DisposeAsync().AsTask()); + foreach (var waiter in new[] { startup, another }) + { + await Assert.ThrowsAnyAsync(() => WaitFor(waiter)); + Assert.True(waiter.IsCanceled); + } + + Assert.True(storage.ReadToken.IsCancellationRequested); + Assert.False(storage.AllowBlockedRead.Task.IsCompleted); + Assert.Empty(storage.OperationLog); + Assert.All(new[] { first, second }, state => + { + Assert.Equal(0, state.RecoveryCompletedCount); + Assert.Equal(0, state.CaptureCount); + Assert.Equal(0, state.WriteCompletedCount); + }); + if (lifecycleStop) + { + Assert.True(manager.TryGetStateMachine("first", out var registered)); + Assert.Same(first, registered); + await Assert.ThrowsAnyAsync(() => manager.WriteStateAsync(CancellationToken.None).AsTask()); + } + + await WaitFor(manager.DisposeAsync().AsTask()); + await Assert.ThrowsAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask()); + } + + [Theory] + [InlineData("provider-cancellation")] + [InlineData("provider-io")] + [InlineData("replay")] + public async Task RecoveryFailure_FaultsAttemptWaitersAndAllowsRetry(string failure) + { + var context = new QueuedSynchronizationContext(); + await context.Run(async () => + { + Exception expected = failure == "provider-cancellation" + ? new OperationCanceledException("Provider read cancellation.", new CancellationToken(canceled: true)) + : new IOException("Provider read failure."); + IJournalStorage storage = failure == "replay" + ? new MutableReadStorage([1, 2, 3], []) + : new CapturingStorage { NextReadException = expected }; + await using var manager = CreateTestSystem(storage).Manager; + var states = new[] { new LifecycleState(), new LifecycleState() }; + manager.RegisterStateMachine("first", states[0]); + manager.RegisterStateMachine("second", states[1]); + var waiters = new[] + { + manager.InitializeAsync(CancellationToken.None).AsTask(), + manager.InitializeAsync(CancellationToken.None).AsTask() + }; + + var observed = await Record.ExceptionAsync(() => WaitFor(waiters[0])); + if (failure == "replay") + { + Assert.Contains("Failed to recover journaling state", Assert.IsType(observed).Message); + Assert.NotNull(observed.InnerException); + } + else + { + Assert.Same(expected, observed); + } + + Assert.Same(observed, await Record.ExceptionAsync(() => WaitFor(waiters[1]))); + Assert.All(states, state => Assert.Equal(0, state.RecoveryCompletedCount)); + var retry = manager.InitializeAsync(CancellationToken.None).AsTask(); + var concurrent = manager.InitializeAsync(CancellationToken.None).AsTask(); + await WaitFor(Task.WhenAll(retry, concurrent)); + Assert.All(states, state => Assert.Equal(1, state.RecoveryCompletedCount)); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.All(states, state => Assert.Equal(1, state.WriteCompletedCount)); + }); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task RecoveryRetry_ShutdownCancelsAllWaitingCallers(bool lifecycleStop) + { + var storage = new MutableReadStorage(2, [1, 2, 3], []); + var sut = CreateTestSystem(storage); + await using var manager = sut.Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await Assert.ThrowsAsync(() => lifecycleStop + ? sut.Lifecycle.OnStart(CancellationToken.None) + : manager.InitializeAsync(CancellationToken.None).AsTask()); + var first = manager.InitializeAsync(CancellationToken.None).AsTask(); + await WaitFor(storage.BlockedReadStarted.Task); + var second = manager.InitializeAsync(CancellationToken.None).AsTask(); + await WaitFor(lifecycleStop + ? sut.Lifecycle.OnStop(TestContext.Current.CancellationToken) + : manager.DisposeAsync().AsTask()); + + await Assert.ThrowsAnyAsync(() => WaitFor(first)); + await Assert.ThrowsAnyAsync(() => WaitFor(second)); + Assert.True(storage.ReadToken.IsCancellationRequested); + Assert.Equal(2, storage.ReadCount); + Assert.Equal(0, state.RecoveryCompletedCount); + Assert.Empty(storage.OperationLog); + if (lifecycleStop) + { + await Assert.ThrowsAnyAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask()); + } + + await manager.DisposeAsync(); + await Assert.ThrowsAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask()); + } + + [Fact] + public async Task RecoveryIoFailure_DuringShutdownPreservesOriginalCause() + { + var expected = new IOException("Read failure during shutdown."); + var entered = NewSignal(); + var storage = Substitute.For(); + storage.ReadAsync(Arg.Any(), Arg.Any()) + .Returns(call => ReadAsync(call.Arg())); + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + var initializing = manager.InitializeAsync(CancellationToken.None).AsTask(); + await WaitFor(entered.Task); + var another = manager.InitializeAsync(CancellationToken.None).AsTask(); + await WaitFor(manager.DisposeAsync().AsTask()); + Assert.Same(expected, await Record.ExceptionAsync(() => WaitFor(initializing))); + Assert.Same(expected, await Record.ExceptionAsync(() => WaitFor(another))); + Assert.Equal(0, state.RecoveryCompletedCount); + + async ValueTask ReadAsync(CancellationToken token) + { + entered.SetResult(); + try + { + await Task.Delay(Timeout.InfiniteTimeSpan, token); + } + catch (OperationCanceledException) when (token.IsCancellationRequested) + { + throw expected; + } + } + } + + [Fact] + public async Task IdleShutdown_CompletesWithoutFencing() + { + var sut = CreateTestSystem(); + await using var manager = sut.Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await sut.Lifecycle.OnStart(TestContext.Current.CancellationToken); + await sut.Lifecycle.OnStop(TestContext.Current.CancellationToken); + Assert.True(manager.TryGetStateMachine("state", out var registered)); + Assert.Same(state, registered); + await Assert.ThrowsAnyAsync(() => manager.WriteStateAsync(CancellationToken.None).AsTask()); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task AdmittedShutdownCancellation_FencesAndFaultsWaiters(bool snapshot) + { + var storage = new CapturingStorage + { + IsCompactionRequested = snapshot, + BlockNextAppend = !snapshot, + BlockNextReplace = snapshot + }; + var sut = CreateTestSystem(storage); + await using var manager = sut.Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await sut.Lifecycle.OnStart(TestContext.Current.CancellationToken); + var write = manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask(); + await WaitFor(snapshot ? storage.ReplaceEntered.Task : storage.BlockedAppendStarted.Task); + var queued = manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask(); + await sut.Lifecycle.OnStop(TestContext.Current.CancellationToken); + var exception = await Assert.ThrowsAnyAsync(() => WaitFor(write)); + Assert.Same(exception, await Record.ExceptionAsync(() => WaitFor(queued))); + var rejected = Assert.Throws(() => manager.TryGetStateMachine("state", out _)); + Assert.Same(exception, rejected.InnerException); + Assert.Equal(0, state.WriteCompletedCount); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task CallerCancellation_PreservesOwnedStorageAndAck(bool snapshot) + { + var storage = new CapturingStorage + { + IsCompactionRequested = snapshot, + BlockNextAppend = !snapshot, + BlockNextReplace = snapshot + }; + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + var acknowledged = NewSignal(); + state.WriteCompletedAction = () => acknowledged.SetResult(); + manager.RegisterStateMachine("state", state); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + using var caller = new CancellationTokenSource(); + var write = manager.WriteStateAsync(caller.Token).AsTask(); + await WaitFor(snapshot ? storage.ReplaceEntered.Task : storage.BlockedAppendStarted.Task); + caller.Cancel(); + await Assert.ThrowsAnyAsync(() => write); + Assert.False(acknowledged.Task.IsCompleted); + state.EmitEntry = false; + storage.IsCompactionRequested = false; + var nextWrite = manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask(); + if (snapshot) storage.ReleaseReplace.SetResult(); + else storage.ReleaseAppend.SetResult(); + await WaitFor(acknowledged.Task); + await WaitFor(nextWrite); + Assert.Equal(2, state.CaptureCount); + Assert.Equal(1, state.WriteCompletedCount); + Assert.Equal(0, manager.PendingWriteByteCount); + Assert.Equal(snapshot ? 0 : 1, storage.Appends.Count); + Assert.Equal(snapshot ? 1 : 0, storage.Replaces.Count); + } + + [Fact] + public async Task AdmittedDeleteShutdownCancellation_FencesAndFaultsWaiters() + { + var storage = new BlockingDeleteStorage(); + var sut = CreateTestSystem(storage); + await using var manager = sut.Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await sut.Lifecycle.OnStart(TestContext.Current.CancellationToken); + var deleting = manager.DeleteStateAsync(CancellationToken.None).AsTask(); + await WaitFor(storage.FirstDeleteStarted.Task); + var queued = manager.WriteStateAsync(CancellationToken.None).AsTask(); + await WaitFor(sut.Lifecycle.OnStop(TestContext.Current.CancellationToken)); + var exception = await Assert.ThrowsAnyAsync(() => WaitFor(deleting)); + Assert.Same(exception, await Record.ExceptionAsync(() => WaitFor(queued))); + var rejected = Assert.Throws(() => manager.TryGetStateMachine("state", out _)); + Assert.Same(exception, rejected.InnerException); + Assert.Equal(1, state.ResetCount); + Assert.Equal(0, state.WriteCompletedCount); + } + + [Fact] + public async Task CallerCancellation_PreservesQueuedWrite() + { + var context = new QueuedSynchronizationContext(); + await context.Run(async () => + { + var storage = new CapturingStorage(); + await using var manager = CreateTestSystem(storage).Manager; + var state = new LifecycleState(); + manager.RegisterStateMachine("state", state); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + using var caller = new CancellationTokenSource(); + var write = manager.WriteStateAsync(caller.Token).AsTask(); + Assert.Equal(0, state.CaptureCount); + caller.Cancel(); + var remainingWaiter = manager.WriteStateAsync(CancellationToken.None).AsTask(); + await Assert.ThrowsAnyAsync(() => write); + await remainingWaiter; + Assert.Single(storage.Appends); + Assert.Equal(1, state.CaptureCount); + Assert.Equal(1, state.WriteCompletedCount); + }); + } + + private static TaskCompletionSource NewSignal() => new(TaskCreationOptions.RunContinuationsAsynchronously); + + private static Task WaitFor(Task task) => task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + + private sealed class QueuedSynchronizationContext : SynchronizationContext + { + private readonly Channel<(SendOrPostCallback Callback, object? State)> _queue = + Channel.CreateUnbounded<(SendOrPostCallback, object?)>(); + + public override void Post(SendOrPostCallback callback, object? state) => Assert.True(_queue.Writer.TryWrite((callback, state))); + + public async Task Run(Func action) + { + Task task = Task.CompletedTask; + Invoke(_ => task = action(), null); + while (!task.IsCompleted) + { + var next = await _queue.Reader.ReadAsync(TestContext.Current.CancellationToken).AsTask().WaitAsync(TimeSpan.FromSeconds(10)); + Invoke(next.Callback, next.State); + } + + await task; + } + + private void Invoke(SendOrPostCallback callback, object? state) + { + var previous = Current; + SetSynchronizationContext(this); + try + { + callback(state); + } + finally + { + SetSynchronizationContext(previous); + } + } + } + + private sealed class LifecycleState : IStateMachine + { + public bool EmitEntry { get; set; } = true; + public Action? ResetAction { get; set; } + public Action? WriteCompletedAction { get; set; } + public Action? RecoveryCompletedAction { get; set; } + public int ResetCount { get; private set; } + public int CaptureCount { get; private set; } + public int WriteCompletedCount { get; private set; } + public int RecoveryCompletedCount { get; private set; } + + public void OnRecoveryCompleted() + { + RecoveryCompletedCount++; + RecoveryCompletedAction?.Invoke(); + } + + public void Reset(JournalStreamWriter writer) + { + ResetCount++; + ResetAction?.Invoke(); + } + + public void WritePendingEntries(JournalStreamWriter writer) + { + CaptureCount++; + if (EmitEntry) + { + using var entry = writer.BeginEntry(); + entry.Writer.GetSpan(1)[0] = 1; + entry.Writer.Advance(1); + entry.Commit(); + } + } + + public void WriteSnapshot(JournalStreamWriter writer) => WritePendingEntries(writer); + + public void OnWriteCompleted() + { + WriteCompletedCount++; + WriteCompletedAction?.Invoke(); + } + + public void ReplayEntry(JournalEntry entry, JournalReplayContext context) => throw new NotSupportedException(); + } +} diff --git a/test/Orleans.Journaling.Tests/StateManagerTests.cs b/test/Orleans.Journaling.Tests/StateManagerTests.cs index 73e072838c6..2c9abd55015 100644 --- a/test/Orleans.Journaling.Tests/StateManagerTests.cs +++ b/test/Orleans.Journaling.Tests/StateManagerTests.cs @@ -26,7 +26,7 @@ namespace Orleans.Journaling.Tests; [TestSuite("BVT")] [TestProvider("None")] [TestCategory("BVT")] -public class StateManagerTests : JournalingTestBase +public partial class StateManagerTests : JournalingTestBase { /// /// Tests the registration and basic operation of multiple states. @@ -121,8 +121,8 @@ public async Task StateManager_WriteOperations_RequireSuccessfulInitialization() var deleteException = await Assert.ThrowsAsync( () => sut.Manager.DeleteStateAsync(TestContext.Current.CancellationToken).AsTask()); - Assert.Contains("fenced", writeException.Message, StringComparison.Ordinal); - Assert.Contains("fenced", deleteException.Message, StringComparison.Ordinal); + Assert.Contains("not been initialized", writeException.Message, StringComparison.Ordinal); + Assert.Contains("not been initialized", deleteException.Message, StringComparison.Ordinal); } [Fact] @@ -617,7 +617,7 @@ await Assert.ThrowsAsync( [InlineData("append")] [InlineData("replace")] [InlineData("delete")] - public async Task StateManager_FailureDeactivatesOwningGrain(string operation) + public async Task StateManager_FailureHandling_RespectsRecoveryBoundary(string operation) { var expected = new IOException("Expected journal operation failure."); var storage = new CapturingStorage(); @@ -635,6 +635,14 @@ public async Task StateManager_FailureDeactivatesOwningGrain(string operation) if (operation == "initialize") { storage.NextReadException = expected; + Assert.Same(expected, await Assert.ThrowsAsync(() => + manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask())); + context.DidNotReceive().Deactivate(Arg.Any(), Arg.Any()); + await manager.InitializeAsync(TestContext.Current.CancellationToken); + value.Value = 42; + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Single(storage.Appends); + return; } else { @@ -648,7 +656,6 @@ public async Task StateManager_FailureDeactivatesOwningGrain(string operation) var failedOperation = operation switch { - "initialize" => manager.InitializeAsync(TestContext.Current.CancellationToken), "delete" => manager.DeleteStateAsync(TestContext.Current.CancellationToken), _ => manager.WriteStateAsync(TestContext.Current.CancellationToken) }; @@ -1064,7 +1071,7 @@ public async Task StateManager_Recovery_RejectsMalformedTrailingData() } [Fact] - public async Task StateManager_FreshRecovery_ReplaysFixedStorage() + public async Task StateManager_RecoveryRetry_ReplaysFixedStorage() { var validBytes = CreatePersistedValueBytes("value", 42); var storage = new MutableReadStorage([.. validBytes, 1, 2, 3], validBytes); @@ -1075,11 +1082,6 @@ await Assert.ThrowsAsync( () => sut.Lifecycle.OnStart(TestContext.Current.CancellationToken) .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken)); - await Assert.ThrowsAsync( - () => sut.Manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask()); - await sut.Manager.DisposeAsync(); - sut = CreateTestSystem(storage: storage); - value = new DurableValue("value", sut.Manager, CreateValueCodec()); await sut.Manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask() .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); await sut.Manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask() @@ -1090,7 +1092,37 @@ await sut.Manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask( } [Fact] - public async Task StateManager_FreshRecovery_PreservesUnknownStreamOnce() + public async Task StateManager_RecoveryRetry_ReplaysListWithoutDuplicatingEntries() + { + var seedStorage = new CapturingStorage(); + await using (var seed = CreateTestSystem(seedStorage).Manager) + { + var source = new DurableList("items", seed, + new OrleansBinaryDurableListCommandCodec(CodecProvider.GetCodec(), SessionPool)); + await seed.InitializeAsync(TestContext.Current.CancellationToken); + source.Add(1); + source.Add(2); + await seed.WriteStateAsync(TestContext.Current.CancellationToken); + } + + var bytes = seedStorage.RecoverableBytes; + var storage = new MutableReadStorage([.. bytes, 1, 2, 3], bytes); + await using var manager = CreateTestSystem(storage).Manager; + var items = new DurableList("items", manager, + new OrleansBinaryDurableListCommandCodec(CodecProvider.GetCodec(), SessionPool)); + await Assert.ThrowsAsync(() => manager.InitializeAsync(CancellationToken.None).AsTask()); + Assert.Equal([1, 2], items); + + await manager.InitializeAsync(TestContext.Current.CancellationToken); + Assert.Equal([1, 2], items); + Assert.Equal(2, storage.ReadCount); + items.Add(3); + await manager.WriteStateAsync(TestContext.Current.CancellationToken); + Assert.Equal([1, 2, 3], items); + } + + [Fact] + public async Task StateManager_RecoveryRetry_PreservesUnknownStreamOnce() { var validBytes = CreateUnknownStreamBytes(new JournalStreamId(99), [1, 2, 3]); var storage = new MutableReadStorage([.. validBytes, 1, 2, 3], validBytes) { IsCompactionRequested = true }; @@ -1100,10 +1132,6 @@ await Assert.ThrowsAsync( () => sut.Lifecycle.OnStart(TestContext.Current.CancellationToken) .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken)); - await Assert.ThrowsAsync( - () => sut.Manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask()); - await sut.Manager.DisposeAsync(); - sut = CreateTestSystem(storage: storage); await sut.Manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask() .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); await sut.Manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask() @@ -1116,7 +1144,7 @@ await sut.Manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask( } [Fact] - public async Task StateManager_FreshRecovery_RemovesStaleRetiredPlaceholder() + public async Task StateManager_RecoveryRetry_RemovesStaleRetiredPlaceholder() { var storage = new MutableReadStorage([.. CreateNamedUnknownStreamBytes("stale", new JournalStreamId(8), [1, 2, 3]), 1, 2, 3], []); var sut = CreateTestSystem(storage: storage); @@ -1125,10 +1153,6 @@ await Assert.ThrowsAsync( () => sut.Lifecycle.OnStart(TestContext.Current.CancellationToken) .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken)); - await Assert.ThrowsAsync( - () => sut.Manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask()); - await sut.Manager.DisposeAsync(); - sut = CreateTestSystem(storage: storage); await sut.Manager.InitializeAsync(TestContext.Current.CancellationToken).AsTask() .WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); await sut.Manager.WriteStateAsync(TestContext.Current.CancellationToken).AsTask() @@ -2276,6 +2300,8 @@ private sealed class MutableReadStorage : IJournalStorage private byte[] _bytes; private int _readCount; + public int ReadCount => Volatile.Read(ref _readCount); + public MutableReadStorage(params byte[][] readSnapshots) : this(blockedReadNumber: 0, readSnapshots) { } @@ -2330,9 +2356,12 @@ public byte[] Bytes public TaskCompletionSource AllowBlockedRead { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + public CancellationToken ReadToken { get; private set; } + public async ValueTask ReadAsync(IJournalStorageConsumer consumer, CancellationToken cancellationToken) { ArgumentNullException.ThrowIfNull(consumer); + ReadToken = cancellationToken; cancellationToken.ThrowIfCancellationRequested(); if (Interlocked.Increment(ref _readCount) == _blockedReadNumber) {