From f89a54c9a37cb7e7a8efae88d9098d66cd457542 Mon Sep 17 00:00:00 2001 From: TedHartMS <15467143+TedHartMS@users.noreply.github.com> Date: Wed, 9 Sep 2026 16:43:10 -0700 Subject: [PATCH 1/6] Document the state machine. --- .../Index/Checkpointing/StateMachineDriver.cs | 63 ++- website/docs/dev/tsavorite/intro.md | 5 +- website/docs/dev/tsavorite/state-machine.md | 521 ++++++++++++++++++ website/sidebars.js | 2 +- 4 files changed, 575 insertions(+), 16 deletions(-) create mode 100644 website/docs/dev/tsavorite/state-machine.md diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineDriver.cs b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineDriver.cs index f06a926ae76..5c9fba08bae 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineDriver.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineDriver.cs @@ -15,21 +15,41 @@ namespace Tsavorite.core /// public class StateMachineDriver { + // Globally published phase and version. SystemState systemState; + + // The single state machine currently owning this driver; null while idle. IStateMachine stateMachine; + + // Already-started tasks that must complete before the driver can leave the current phase. + // ProcessWaitingListAsync awaits these only after the transition-in epoch barrier completes. readonly List<(Task task, StateMachineTaskType type)> waitingList; + + // Completion source for the entire state-machine run, not an individual phase transition. TaskCompletionSource stateMachineCompleted; - // All threads have entered the given state + + // Semaphore associated with the currently published state. MakeTransitionWorker releases it after + // prior-epoch participants have advanced or suspended and GlobalAfterEnteringState has completed. + // This does not include completion of tasks in waitingList. SemaphoreSlim waitForTransitionIn; + + // GlobalAfterEnteringState may run on an arbitrary epoch-drain thread, so its exception is captured + // here and rethrown by ProcessWaitingListAsync on the state-machine driver path. Exception waitForTransitionInException; - // All threads have exited the given state + + // Semaphore associated with the currently published state. GlobalStateMachineStep releases it as + // soon as the next state is published; it does not wait for the transition-in epoch barrier. SemaphoreSlim waitForTransitionOut; - // Transactions drained in last version + + // Version whose active transactions must drain before the state machine can advance. long lastVersion; TaskCompletionSource lastVersionTransactionsDone; + List callbacks; readonly LightEpoch epoch; readonly ILogger logger; + + // Active transaction counts are indexed by version parity; only two adjacent versions can be active. readonly long[] NumActiveTransactions; public SystemState SystemState => SystemState.Copy(ref systemState); @@ -157,6 +177,7 @@ public void EndTransaction(long txnVersion) internal void AddToWaitingList(Task waiter, StateMachineTaskType type) { + // Callers start the operation before registering it. The driver awaits it after transition-in. if (waiter != null) waitingList.Add((waiter, type)); } @@ -219,22 +240,25 @@ void GlobalStateMachineStep(SystemState expectedState) var nextState = stateMachine.NextState(systemState); + // Run task-specific work while systemState still identifies the previous phase. stateMachine.GlobalBeforeEnteringState(nextState, this); - // Execute any additional registered callbacks + // External callbacks have the same before-publication ordering as the state-machine task hooks. if (callbacks != null) { foreach (var callback in callbacks) callback.BeforeEnteringState(nextState); } - // Write new phase + // Publish the new phase and version so subsequent session refreshes observe nextState. systemState.Word = nextState.Word; - // Release waiters for new phase + // Release the semaphore associated with the phase just exited. _ = waitForTransitionOut?.Release(int.MaxValue); - // Write new semaphores + // Install semaphores for the newly published phase. Its transition-out semaphore is released + // when the following phase is published; its transition-in semaphore is released by the + // epoch-drain callback below. These assignments occur after systemState is published. waitForTransitionOut = new SemaphoreSlim(0); waitForTransitionIn = new SemaphoreSlim(0); @@ -244,6 +268,9 @@ void GlobalStateMachineStep(SystemState expectedState) try { epoch.Resume(); + + // Associate MakeTransitionWorker with the prior epoch. It becomes eligible only after + // participants still announcing that epoch have advanced through ProtectAndDrain or suspended. epoch.BumpCurrentEpoch(() => MakeTransitionWorker(nextState)); } finally @@ -259,6 +286,7 @@ void GlobalStateMachineStep(SystemState expectedState) /// public async Task WaitForStateChange(SystemState currentState) { + // Capture before rechecking state so a racing transition that releases this semaphore is observed. var _waitForTransitionOut = waitForTransitionOut; if (SystemState.Equal(currentState, systemState)) { @@ -273,7 +301,11 @@ public async Task WaitForStateChange(SystemState currentState) /// public async Task WaitForCompletion(SystemState currentState) { + // First wait until currentState is no longer published. await WaitForStateChange(currentState).ConfigureAwait(false); + + // Then capture the newly published state and wait until its epoch transition and + // GlobalAfterEnteringState hooks complete. Phase waiting-list tasks are not included. currentState = systemState; var _waitForTransitionIn = waitForTransitionIn; if (SystemState.Equal(currentState, systemState)) @@ -286,29 +318,35 @@ void MakeTransitionWorker(SystemState nextState) { try { + // This is an epoch-drain action and may execute synchronously from BumpCurrentEpoch + // or later on any thread that advances or suspends epoch protection. stateMachine.GlobalAfterEnteringState(nextState, this); } catch (Exception e) { - // Store the exception to be thrown by state machine driver - // We do not throw here as this epoch action may be executed in a different thread context + // Propagate on the driver path rather than throwing on an arbitrary epoch-drain thread. waitForTransitionInException = e; logger?.LogError(e, "Exception in state machine transition worker"); } finally { + // Signal that the epoch transition and all after-transition hooks have finished. waitForTransitionIn.Release(int.MaxValue); } } async Task ProcessWaitingListAsync(CancellationToken token = default) { + // Do not process phase tasks until the prior epoch has drained and after-transition hooks finish. await waitForTransitionIn.WaitAsync(token).ConfigureAwait(false); if (waitForTransitionInException != null) { throw waitForTransitionInException; } + + // These tasks were started by state-machine hooks and may have progressed concurrently with + // the epoch transition. Awaiting them here prevents the driver from publishing the next phase. foreach (var (task, type) in waitingList) { try @@ -331,6 +369,7 @@ async Task RunStateMachine(CancellationToken token = default) { do { + // Publish one transition, then wait for both transition-in and its registered phase work. GlobalStateMachineStep(systemState); await ProcessWaitingListAsync(token).ConfigureAwait(false); } while (systemState.Phase != Phase.REST); @@ -377,14 +416,14 @@ void FastForwardStateMachineToRest() if (waitForTransitionOut?.CurrentCount == 0) _ = waitForTransitionOut?.Release(int.MaxValue); - // Clear semaphores + // Failure recovery does not execute skipped transition hooks. Discard their synchronization state. waitForTransitionOut = null; waitForTransitionIn = null; - // Clear exception if any + // Clear any exception captured from an after-transition hook. waitForTransitionInException = null; - // Clear waiting list + // The failed run no longer waits for phase-specific asynchronous work. waitingList.Clear(); } } diff --git a/website/docs/dev/tsavorite/intro.md b/website/docs/dev/tsavorite/intro.md index f0455d1e613..1d574c54104 100644 --- a/website/docs/dev/tsavorite/intro.md +++ b/website/docs/dev/tsavorite/intro.md @@ -10,6 +10,5 @@ Garnet’s storage layer, called Tsavorite, was forked from our prior open-sourc tiered storage support (memory, SSD, and cloud storage), fast non-blocking checkpointing, recovery, operation logging for durability, multi-key [locking](locking.md) and transaction support, and better memory management and [space reuse](reviv.md). - - - +Checkpointing, version changes, and index growth are coordinated by the +[state machine driver](state-machine.md). diff --git a/website/docs/dev/tsavorite/state-machine.md b/website/docs/dev/tsavorite/state-machine.md new file mode 100644 index 00000000000..543d729b8a3 --- /dev/null +++ b/website/docs/dev/tsavorite/state-machine.md @@ -0,0 +1,521 @@ +--- +id: state-machine +sidebar_label: State Machine Driver +title: State Machine Driver +--- + +# Tsavorite state machine driver + +Tsavorite uses `StateMachineDriver` to coordinate operations that require all +sessions to move through a sequence of globally visible phases. Checkpointing +and index growth use this mechanism to publish phase and version changes, +synchronize participating threads through `LightEpoch`, run phase-specific +actions, and wait for asynchronous work. + +This page describes the current implementation. The principal files are under +`libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/`. + +## Concepts and types + +The state machine has three layers: + +| Layer | Responsibility | +|---|---| +| `IStateMachine` | Defines the phase graph through `NextState`. | +| `IStateMachineTask` | Performs work before and after each transition. | +| `StateMachineDriver` | Owns the current state, publishes transitions, coordinates the epoch barrier, and waits for phase work. | + +`SystemState` stores the current `Phase` and `Version` in one 64-bit `Word`. +The high 8 bits contain the phase and the remaining bits contain the version. +The driver starts in `REST`, version 1. + +`IStateMachine` extends `IStateMachineTask`: + +```cs +public interface IStateMachine : IStateMachineTask +{ + SystemState NextState(SystemState currentState); +} +``` + +This gives every state machine three operations: + +- `NextState(currentState)` returns the state that should be published next. +- `GlobalBeforeEnteringState(nextState, driver)` runs before `nextState` is + published. Epoch participants still belong to the previous transition. +- `GlobalAfterEnteringState(nextState, driver)` runs after the state is + published and holders of the prior epoch have advanced or suspended. + +`StateMachineBase` composes an ordered array of `IStateMachineTask` instances. +It calls every task's before-transition method in array order, and later calls +every task's after-transition method in the same order. For example, a full +checkpoint is constructed with the index checkpoint task first and the hybrid +log backend second, so index work precedes hybrid-log work within each hook. + +## Defining a state machine + +A concrete state machine implements `NextState` as a phase graph. The base +version-change graph is: + +```text +REST(v) -> PREPARE(v) -> IN_PROGRESS(v + 1) -> REST(v + 1) +``` + +Checkpoint and index-growth state machines extend or replace this graph: + +| State machine | Phase sequence | +|---|---| +| `VersionChangeSM` | `REST -> PREPARE -> IN_PROGRESS -> REST` | +| `HybridLogCheckpointSM` | `REST -> PREPARE -> IN_PROGRESS -> WAIT_FLUSH -> PERSISTENCE_CALLBACK -> REST` | +| `FullCheckpointSM` | `REST -> PREPARE -> IN_PROGRESS -> WAIT_INDEX_CHECKPOINT -> WAIT_FLUSH -> PERSISTENCE_CALLBACK -> REST` | +| `IndexCheckpointSM` | `REST -> PREPARE -> WAIT_INDEX_CHECKPOINT -> WAIT_FLUSH -> PERSISTENCE_CALLBACK -> REST` | +| `StreamingSnapshotCheckpointSM` | `REST -> PREPARE -> IN_PROGRESS -> WAIT_FLUSH -> REST` | +| `IndexResizeSM` | `REST -> PREPARE_GROW -> IN_PROGRESS_GROW -> REST` | + +`VersionChangeSM` increments the version when it returns `IN_PROGRESS`. +Derived checkpoint state machines retain that increment. Index-only +checkpointing and index growth do not increment the `SystemState` version. + +The static `Checkpoint` factory creates the state machine and its ordered task +set: + +- `Checkpoint.Full` combines `IndexCheckpointSMTask` with either + `FoldOverSMTask` or `SnapshotCheckpointSMTask`. +- `Checkpoint.IndexOnly` creates an `IndexCheckpointSM` containing index + checkpoint tasks. +- `Checkpoint.HybridLogOnly` creates a `HybridLogCheckpointSM` containing the + selected hybrid-log backend. +- `Checkpoint.Streaming` creates a `StreamingSnapshotCheckpointSM` containing + streaming snapshot tasks. + +The factory assigns one checkpoint GUID to all tasks in the operation. Its +two-store overloads place both stores' tasks in the same state machine so they +advance through the same global phase sequence. + +Index growth is constructed directly from `IndexResizeSMTask` and +`IndexResizeSM`. + +## Launching a state machine + +Checkpoint APIs first construct the state machine and then call +`StateMachineDriver.Register`: + +```cs +var stateMachine = Checkpoint.Full(this, checkpointType, out token); +return stateMachineDriver.Register(stateMachine, cancellationToken); +``` + +`Register` uses `Interlocked.CompareExchange` to install the state machine only +if no other state machine is active. If another checkpoint, index checkpoint, +or index growth operation already owns the driver, `Register` returns `false`. +On success it: + +1. Creates `stateMachineCompleted`, a + `TaskCompletionSource` whose continuations run asynchronously. +2. Starts `RunStateMachine` on a thread-pool task. +3. Returns `true` as soon as the operation has been accepted. + +The `TryInitiate*Checkpoint` APIs expose this non-blocking behavior. +`Take*CheckpointAsync` wraps it by calling `CompleteCheckpointAsync` when +registration succeeds. + +`CompleteCheckpointAsync` delegates to `StateMachineDriver.CompleteAsync`, +which awaits `stateMachineCompleted`. It cannot be called while the caller +holds epoch protection. `CompleteCheckpointAsync` resets the index and +hybrid-log checkpoint structures if that wait throws or is canceled, then +rethrows. + +`StateMachineDriver.RunAsync` uses the same compare-exchange ownership check +and the same driver loop, but directly awaits `RunStateMachine` rather than +launching it through `Task.Run`. Index growth uses this path. Direct +checkpoint consumers can use it as well; Garnet's `DatabaseManagerBase`, for +example, constructs a checkpoint state machine and passes it to `RunAsync`. +That direct path does not pass through `CompleteCheckpointAsync` and therefore +does not receive its catch/reset behavior. + +The cancellation token passed to `CompleteAsync` is also registered to cancel +the shared completion source. Canceling a completion wait can therefore make +other completion waiters observe cancellation. If this token differs from the +token used to launch the driver, canceling the wait does not itself guarantee +that the driver has stopped before `CompleteCheckpointAsync` resets the +checkpoint structures. + +## The driver loop + +`RunStateMachine` repeats two operations: + +```cs +do +{ + GlobalStateMachineStep(systemState); + await ProcessWaitingListAsync(token).ConfigureAwait(false); +} while (systemState.Phase != Phase.REST); +``` + +`GlobalStateMachineStep` publishes one transition. +`ProcessWaitingListAsync` waits for that transition to finish and then waits +for asynchronous work registered for the phase. The next transition does not +begin until both operations complete. + +## One state transition, step by step + +The following sequence occurs for every transition. + +### 1. Verify the expected state + +`GlobalStateMachineStep` receives the state observed by the driver loop. It +compares that state with the live `systemState` and returns without doing +anything if they differ. + +### 2. Calculate the next state + +The driver calls: + +```cs +var nextState = stateMachine.NextState(systemState); +``` + +Only the state machine defines the graph. The driver does not have +checkpoint-specific phase logic. + +### 3. Run before-transition task hooks + +The driver calls: + +```cs +stateMachine.GlobalBeforeEnteringState(nextState, this); +``` + +For a `StateMachineBase`, this invokes each constituent +`IStateMachineTask.GlobalBeforeEnteringState` in construction order. +`nextState` is not yet globally visible. + +These hooks currently initialize checkpoint state, capture addresses, publish +checkpoint-manager version-shift notifications, start I/O, and add already +created tasks to the driver's waiting list. + +### 4. Run optional external callbacks + +Callbacks installed through `UnsafeRegisterCallback` receive +`BeforeEnteringState(nextState)` after the state-machine tasks and before the +state is published. Registration is not safe concurrently with Tsavorite +operations, and expensive callback work delays the transition. + +### 5. Publish the new state + +The driver assigns: + +```cs +systemState.Word = nextState.Word; +``` + +Sessions that subsequently refresh their local execution context can now +observe the new phase and version. + +### 6. Release transition-out waiters + +The driver releases the existing `waitForTransitionOut` semaphore. That +semaphore belongs to the state being exited, so its waiters can now observe +that the state has changed. + +### 7. Create semaphores for the new state + +The driver creates new zero-count `waitForTransitionOut` and +`waitForTransitionIn` semaphores. Once both assignments complete, they are +intended to describe the newly published state: + +- the new `waitForTransitionOut` will be released when the driver publishes + the following state; +- the new `waitForTransitionIn` will be released when the epoch transition + into this state is complete. + +### 8. Establish the epoch boundary + +The state-machine driver temporarily resumes epoch protection and calls: + +```cs +epoch.BumpCurrentEpoch(() => MakeTransitionWorker(nextState)); +``` + +`BumpCurrentEpoch` associates `MakeTransitionWorker` with the prior epoch. +The action becomes eligible only after threads that still announce that prior +epoch have suspended or advanced through `ProtectAndDrain`. + +The action may run synchronously from `BumpCurrentEpoch`, or later on any +thread that drains the epoch. It must not rely on thread-affine state. + +This is a safe-point barrier, not a count of completed API calls. A participant +can advance its announced epoch at an explicit `ProtectAndDrain` point without +returning from the outer API call. Code after such a point must follow the +newly visible state or otherwise preserve the state machine's invariants. + +### 9. Run after-transition task hooks + +Once the prior epoch is safe, `MakeTransitionWorker` calls: + +```cs +stateMachine.GlobalAfterEnteringState(nextState, this); +``` + +The current after-transition work includes: + +- tracking transactions from the previous checkpoint version after entering + `IN_PROGRESS`; +- tracking transactions before index growth proceeds; +- splitting all buckets after entering `IN_PROGRESS_GROW`. + +An exception cannot be thrown directly from this callback because it may be +running on an unrelated epoch participant. `MakeTransitionWorker` stores it in +`waitForTransitionInException` and logs it. + +### 10. Release transition-in waiters + +In a `finally` block, `MakeTransitionWorker` releases +`waitForTransitionIn`. This occurs whether the after-transition hook succeeds +or fails, provided the worker runs before failure recovery clears the shared +transition fields. See [Completion, cancellation, and +failure](#completion-cancellation-and-failure) for the current cancellation +race. + +### 11. Wait for the transition and phase tasks + +`ProcessWaitingListAsync` first waits on `waitForTransitionIn`. It then: + +1. rethrows any exception captured from the after-transition hook; +2. awaits every task in `waitingList` in insertion order; +3. logs and propagates any non-cancellation task failure; +4. clears the waiting list. + +The tasks are normally already running before this method awaits them. +Sequential awaits therefore do not imply that the underlying I/O was issued +sequentially. + +```mermaid +sequenceDiagram + participant D as StateMachineDriver + participant S as IStateMachine + participant E as LightEpoch + participant P as Epoch participants + + D->>S: NextState(current) + D->>S: GlobalBeforeEnteringState(next) + D->>D: Publish systemState = next + D->>D: Release old transitionOut + D->>D: Create next transitionOut / transitionIn + D->>E: BumpCurrentEpoch(MakeTransitionWorker) + E-->>P: Prior epoch must drain + P-->>E: Suspend or ProtectAndDrain + E->>S: GlobalAfterEnteringState(next) + E->>D: Release transitionIn + D->>D: Await waitingList +``` + +## Transition-out and transition-in waiters + +The two semaphores answer different questions. + +### `waitForTransitionOut` + +`waitForTransitionOut` is intended to mean "the driver has published a state +different from this one." `WaitForStateChange` captures the current semaphore +and then checks that the caller's state is still current: + +```cs +var transitionOut = waitForTransitionOut; +if (SystemState.Equal(currentState, systemState)) + await transitionOut.WaitAsync(); +``` + +Capturing before checking avoids one missed-release race. However, the current +publisher writes `systemState` before releasing the old transition-out +semaphore and installing the new semaphores. During that window, a concurrent +caller can observe the new state with the previous semaphore (or with a null +semaphore during the first transition). These methods are therefore not a +linearizable state-wait API in the current implementation. + +### `waitForTransitionIn` + +`waitForTransitionIn` is intended to mean "the epoch callback and all +`GlobalAfterEnteringState` hooks for the published state have finished." It +does not mean that asynchronous checkpoint I/O in the waiting list has +finished. + +`WaitForCompletion` first waits to leave the supplied state. It then samples +the newly current state and its transition-in semaphore, rechecks that the +state is unchanged, and waits for transition-in completion. It is subject to +the same publication window described above. + +The driver itself also waits on `waitForTransitionIn` at the start of +`ProcessWaitingListAsync`. + +## The waiting list + +`waitingList` is a list of `(Task, StateMachineTaskType)` pairs. State-machine +tasks call `AddToWaitingList` after starting asynchronous work. The type is +used to identify failures in logs. + +Current waiting-list entries are: + +| Type | Work being awaited | +|---|---| +| `LastVersionTransactionsDone` | Transactions still active in the previous version | +| `IndexCheckpointSMTaskMainIndexCheckpoint` | Main hash-index checkpoint I/O | +| `IndexCheckpointSMTaskOverflowBucketsCheckpoint` | Overflow-bucket checkpoint I/O | +| `FoldOverSMTaskHybridLogFlushed` | Fold-over hybrid-log flush | +| `SnapshotCheckpointSMTaskHybridLogFlushed` | Snapshot hybrid-log flush | + +A task can be added from a before-transition or after-transition hook. For +example, index checkpoint I/O starts during `PREPARE`, and its existing tasks +are added to the waiting list before entering `WAIT_INDEX_CHECKPOINT`. +Previous-version transaction tracking is added from the after-transition hook +for `IN_PROGRESS`. + +In the current snapshot implementation, +`SnapshotCheckpointSMTask.GlobalBeforeEnteringState(WAIT_FLUSH)` creates and +initializes the snapshot devices, calls `AsyncFlushPagesForSnapshot`, and adds +the returned flush task to the list. Consequently, snapshot flush issuance +begins before `WAIT_FLUSH` is published; `ProcessWaitingListAsync` later waits +for its completion after transition-in. + +## Session participation + +Safe context operations call `UnsafeResumeThread` before entering Tsavorite +and `UnsafeSuspendThread` in a `finally` block. Resume acquires epoch +protection and calls `InternalRefresh`, which: + +1. calls `epoch.ProtectAndDrain`; +2. copies the driver's `SystemState` into the session execution context; +3. applies phase-specific handling. + +For example, after the global state enters `IN_PROGRESS`, an active +transaction whose version is still the previous version receives an effective +local state of `PREPARE` at that older version. `PREPARE_GROW` prevents +non-transactional sessions from proceeding until index growth reaches a phase +they can enter. + +Unsafe contexts manage their epoch lifetime explicitly, but participate in the +same epoch transitions. + +For more detail about acquisition, suspension, refresh, and drain actions, see +[Epoch Protection](epochprotection.md). + +## Transaction tracking + +Transactions may span individual context operations, so epoch participation +alone does not describe their full lifetime. `StateMachineDriver` therefore +tracks active transaction counts by version. + +The transaction sequence is: + +1. `AcquireTransactionVersion` reads the current system version and increments + its active count. +2. The transaction acquires its key locks. +3. `VerifyTransactionVersion` checks whether a version transition occurred + during lock acquisition. If so, it moves the active count to the new + version. +4. `EndTransaction` decrements the final version's count. + +After entering checkpoint `IN_PROGRESS`, +`HybridLogCheckpointSMTask.GlobalAfterEnteringState` calls +`TrackLastVersion`. If transactions remain in the old version, the driver +creates `lastVersionTransactionsDone` and adds it to the waiting list. +Therefore the driver does not proceed to `WAIT_FLUSH` until those transactions +finish. New-version transactions can continue. + +Index growth uses the same mechanism but treats `PREPARE_GROW` as a full +barrier that prevents new transactions from starting. + +## Current checkpoint phase work + +The phase graph determines ordering, while tasks determine what each phase +does. + +| Entering phase | Current checkpoint work | +|---|---| +| `PREPARE` | Initialize checkpoint state, record the start and begin addresses, initialize the index device where applicable, and start the fuzzy index checkpoint for index/full checkpoints. Snapshot and fold-over backends initialize hybrid-log metadata; snapshot log devices are not initialized until the `WAIT_FLUSH` before-hook. Streaming snapshot starts phase-one scanning. | +| `IN_PROGRESS` | Notify the checkpoint manager that the version shift is starting, issue the version-shift trigger, publish the incremented version, and then track transactions that remain in the old version. | +| `WAIT_INDEX_CHECKPOINT` | Add the already-running main-index and overflow-bucket checkpoint tasks to the waiting list. | +| `WAIT_FLUSH` | End the version shift, issue the flush-begin trigger, verify old-version transactions are drained, and capture the final fuzzy-region address. Snapshot starts copying pages to snapshot devices; fold-over shifts the read-only address and waits for its flush; streaming snapshot performs phase-two scanning. | +| `PERSISTENCE_CALLBACK` | Commit index and hybrid-log metadata, capture final object-log positions, and dispose snapshot devices. | +| `REST` | Clean old checkpoint artifacts, issue the checkpoint-completed trigger, dispose/reset checkpoint state, and advance the checkpoint completion chain. | + +All task-specific work in this table is invoked from +`GlobalBeforeEnteringState` except old-version transaction tracking, which is +invoked from `GlobalAfterEnteringState(IN_PROGRESS)`. State publication itself +is performed by the driver between the before-transition and after-transition +hooks. + +## Index growth phase work + +Index growth uses a shorter graph: + +1. Before entering `PREPARE_GROW`, capture the current version. +2. After entering `PREPARE_GROW`, track existing transactions. New + transactions are prevented from starting in this phase. +3. Before entering `IN_PROGRESS_GROW`, verify both transaction-version counts + are zero, allocate and publish the new hash-table version, and initialize + split tracking. +4. After entering `IN_PROGRESS_GROW`, split all buckets. +5. Return to `REST`. + +## Completion, cancellation, and failure + +When the state machine reaches `REST`, `RunStateMachine` leaves its loop. Its +`finally` block: + +- clears `stateMachineCompleted` from the driver; +- atomically releases ownership of the active `stateMachine`; +- completes the saved completion source successfully, as canceled, or with the + captured exception. + +The loop can exit abnormally in four ways: + +- a before-transition hook throws directly from `GlobalStateMachineStep`; +- an after-transition hook stores its exception in + `waitForTransitionInException`, which `ProcessWaitingListAsync` rethrows; +- a waiting-list task faults while `ProcessWaitingListAsync` awaits it. +- the driver token cancels the transition-in wait or a waiting-list task wait. + +`RunStateMachine` catches these failures and calls +`FastForwardStateMachineToRest`. Fast-forwarding repeatedly calls +`NextState` and publishes only the resulting `SystemState.Word` values until +the phase is `REST`; it does not invoke the skipped before-transition or +after-transition hooks. It then resets old-version transaction tracking, +releases transition-out waiters, clears transition state and exceptions, and +clears the waiting list. + +The original exception is logged, rethrown, and stored on the completion +source so a caller awaiting checkpoint completion observes the failure. + +There is a current failure-ordering limitation when cancellation occurs while +`MakeTransitionWorker` is still queued in the epoch drain list. +`ProcessWaitingListAsync` can observe cancellation and fast-forward to `REST`, +which clears the shared `stateMachine` and `waitForTransitionIn` fields without +releasing transition-in. If the queued worker subsequently runs, it uses those +shared fields rather than captured stable references. Consequently, an +external transition-in waiter is not guaranteed to be released on this path, +and the delayed worker can encounter cleared state. Normal transition +completion and exceptions thrown directly by an executing after-transition +hook do release transition-in through the worker's `finally` block. + +## Rules for state-machine task code + +- Put work that must occur before a state becomes visible in + `GlobalBeforeEnteringState`. +- Treat `GlobalAfterEnteringState` as an epoch drain action. It can run + synchronously or on an arbitrary thread, so it must be thread-agnostic. +- Do not block while holding an epoch needed by the transition being awaited. +- Add only valid, already-created tasks to `waitingList`; a null task is + ignored. +- Remember that the driver waits for transition-in before awaiting the waiting + list, but the listed tasks may have started before state publication. +- Preserve task construction order when one task's phase work depends on + another task. +- Ensure `NextState` always provides a path back to `REST`; failure recovery + follows that graph without invoking task hooks. + +## Related topics + +- [Epoch Protection](epochprotection.md) +- [Locking](locking.md) +- [Store Functions](storefunctions.md) diff --git a/website/sidebars.js b/website/sidebars.js index 08666b099a3..244b2405bcb 100644 --- a/website/sidebars.js +++ b/website/sidebars.js @@ -24,7 +24,7 @@ const sidebars = { {type: 'category', label: 'Server Extensions', items: ["extensions/overview", "extensions/raw-strings", "extensions/objects", "extensions/transactions", "extensions/procedure", "extensions/module"]}, {type: 'category', label: 'Cluster Mode', items: ["cluster/overview", "cluster/replication", "cluster/key-migration"]}, {type: 'category', label: 'Developer Guide', items: ["dev/onboarding", "dev/code-structure", "dev/configuration", "dev/network", "dev/processing", "dev/garnet-api", - {type: 'category', label: 'Tsavorite - Storage Layer', collapsed: true, items: ["dev/tsavorite/intro", "dev/tsavorite/reviv", "dev/tsavorite/locking", "dev/tsavorite/readcache", "dev/tsavorite/storefunctions", "dev/tsavorite/epochprotection", "dev/tsavorite/logrecord", "dev/tsavorite/object-allocator", "dev/tsavorite/buffer-pool"]}, + {type: 'category', label: 'Tsavorite - Storage Layer', collapsed: true, items: ["dev/tsavorite/intro", "dev/tsavorite/reviv", "dev/tsavorite/locking", "dev/tsavorite/readcache", "dev/tsavorite/storefunctions", "dev/tsavorite/epochprotection", "dev/tsavorite/state-machine", "dev/tsavorite/logrecord", "dev/tsavorite/object-allocator", "dev/tsavorite/buffer-pool"]}, {type: 'category', label: 'Chunked Record Layouts', collapsed: true, items: ["dev/aof-record-layout", "dev/migration-replication-record-layout"]}, "dev/device-tuning", "dev/transactions", From 78f361812eaf938db6fbb43a35371ca42dec2974 Mon Sep 17 00:00:00 2001 From: TedHartMS <15467143+TedHartMS@users.noreply.github.com> Date: Thu, 10 Sep 2026 11:51:10 -0700 Subject: [PATCH 2/6] Make HeapObjectBase serialization exception-safe, and do not reset SerializationPhase in ClearSerializedObjectData --- .../cs/src/core/Allocator/HeapObjectBase.cs | 49 ++++++++++++++----- 1 file changed, 38 insertions(+), 11 deletions(-) diff --git a/libs/storage/Tsavorite/cs/src/core/Allocator/HeapObjectBase.cs b/libs/storage/Tsavorite/cs/src/core/Allocator/HeapObjectBase.cs index 85693a802d0..2110045b319 100644 --- a/libs/storage/Tsavorite/cs/src/core/Allocator/HeapObjectBase.cs +++ b/libs/storage/Tsavorite/cs/src/core/Allocator/HeapObjectBase.cs @@ -62,10 +62,18 @@ public void Serialize(BinaryWriter writer) // CopyUpdater does that, as it must ensure the object's (v1) data is not changed during the checkpoint. if (SerializationPhase == SerializationPhase.REST && MakeTransition(SerializationPhase.REST, SerializationPhase.SERIALIZING)) { - // Directly serialize to wire, do not cache serialized state - WriteType(writer, isNull: false); - DoSerialize(writer); - SerializationPhase = SerializationPhase.REST; + // Directly serialize to wire, do not cache serialized state. Restore REST even if DoSerialize + // throws: leaving the object in SERIALIZING would make every later Serialize() and + // CacheSerializedObjectData() spin on it forever, wedging the checkpoint pipeline. + try + { + WriteType(writer, isNull: false); + DoSerialize(writer); + } + finally + { + SerializationPhase = SerializationPhase.REST; + } return; } @@ -144,10 +152,20 @@ public void CacheSerializedObjectData(ref LogRecord dstLogRecord, ref RMWInfo rm { if (SerializationPhase == (int)SerializationPhase.REST && MakeTransition(SerializationPhase.REST, SerializationPhase.SERIALIZING)) { - using var ms = new MemoryStream(); - using var writer = new BinaryWriter(ms, Encoding.UTF8); - DoSerialize(writer); - serializedBytes = ms.ToArray(); + try + { + using var ms = new MemoryStream(); + using var writer = new BinaryWriter(ms, Encoding.UTF8); + DoSerialize(writer); + serializedBytes = ms.ToArray(); + } + catch + { + // Publish no partial capture, and do not strand the object in SERIALIZING. + serializedBytes = null; + SerializationPhase = SerializationPhase.REST; + throw; + } SerializationPhase = SerializationPhase.SERIALIZED; // This is the only place .SERIALIZED is set break; @@ -167,9 +185,18 @@ public void CacheSerializedObjectData(ref LogRecord dstLogRecord, ref RMWInfo rm /// public void ClearSerializedObjectData() { - // Do not disturb an unrelated object that is currently serializing. - if (Interlocked.Exchange(ref serializedBytes, null) is not null) - SerializationPhase = SerializationPhase.REST; + // Release the cached (v) bytes so they can be GC'd, but deliberately leave the phase terminal. + // + // Bytes are only ever cached by CacheSerializedObjectData for a source object that a CopyUpdate + // has just superseded, and Clone() is a shallow copy: the (v+1) record that replaced this one + // shares this object's internal collections and keeps mutating them. Returning to REST would let + // a later checkpoint take the direct-serialize path and enumerate those still-shared, still-live + // collections with no synchronization against the writer. + // + // SERIALIZED with null bytes is the state Serialize() already documents as "superseded after + // checkpoint completion": it writes a null indicator. That is safe because the superseding record + // always sits at a higher address and carries the live data. + _ = Interlocked.Exchange(ref serializedBytes, null); } } } \ No newline at end of file From 9aa4b27e1612eead511e0ec06ce7d1831a0220ac Mon Sep 17 00:00:00 2001 From: TedHartMS <15467143+TedHartMS@users.noreply.github.com> Date: Sat, 12 Sep 2026 09:45:49 -0700 Subject: [PATCH 3/6] Close three checkpoint-vs-mutation gaps in the SAD protocol Do not mutate a live object from the read path HashTimeToLive ran DeleteExpiredItems(), mutating hash, expirationTimes and expirationQueue while holding only a shared lock, so a concurrent snapshot could enumerate the dictionary as it changed. The purge is dropped; expired fields are already filtered by ContainsKey, so the response is unchanged. Do not dispose a source record that a checkpoint has frozen Disposal clears the record's heap fields, which returns the value's ObjectIdMap slot to the page free list for reuse by another record. The snapshot flush reads object ids from its page copy but resolves them against the live map, so disposing a frozen record let the flush serialize a freed, or recycled and therefore unrelated, object. OnDisposeSupersededSource() skips disposal while frozen; the value is released at page eviction instead. Applied to the four RMW expiration paths and the Delete path. Keep reused allocations above the checkpoint's fuzzy-region start The fuzzy-start clamp on minRevivAddress only covered the free record pool, so an allocation saved before the version shift and retried afterwards could place a (v+1) record below the fuzzy start, where the snapshot includes it regardless of its version bit. The clamp now covers both reuse paths. Elide the source record on Delete CreateNewRecordDelete already passed elideSourceRecord to the allocator, which relaxes minRevivAddress on the premise that the source leaves the tag chain, but it never actually elided it; the relaxed bound was therefore unsound. The tombstone now carries the source's PreviousAddress and the source is invalidated and freelisted, as CreateNewRecordUpsert and CreateNewRecordRMW do. Unlike those, Delete fires OnDispose(Deleted) before eliding, because Garnet keys range-index file cleanup and Vector Set deletion off that reason; the subsequent elide or freelist disposal only accounts for the key, so nothing is counted twice. PageEvictionFiresOnEvictForEveryLiveRecord assumed delete sources stay in the tag chain and are visited by eviction. It now derives the expected eviction count from the number of sources actually elided. --- libs/server/Objects/Hash/HashObjectImpl.cs | 5 ++-- .../Tsavorite/Implementation/BlockAllocate.cs | 17 +++++++----- .../Index/Tsavorite/Implementation/Helpers.cs | 21 +++++++++++++++ .../Implementation/InternalDelete.cs | 26 ++++++++++++++++--- .../Tsavorite/Implementation/InternalRMW.cs | 16 ++++++------ .../test.recordops/RecordLifecycleTests.cs | 24 +++++++++++------ 6 files changed, 82 insertions(+), 27 deletions(-) diff --git a/libs/server/Objects/Hash/HashObjectImpl.cs b/libs/server/Objects/Hash/HashObjectImpl.cs index 5156a8c3a41..bbb67db7275 100644 --- a/libs/server/Objects/Hash/HashObjectImpl.cs +++ b/libs/server/Objects/Hash/HashObjectImpl.cs @@ -456,8 +456,9 @@ private void HashExpire(ref ObjectInput input, ref ObjectOutput output, byte res private void HashTimeToLive(ref ObjectInput input, ref ObjectOutput output, byte respProtocolVersion) { - DeleteExpiredItems(); - + // This is a read operation, so it must not mutate the object: it runs under a shared lock and may execute + // while the flush path is serializing this same instance. Expired fields are reported as absent by + // ContainsKey (via GetExpiration) without being removed; the mutating paths purge them from the live object. var isMilliseconds = input.arg1 == 1; var isTimestamp = input.arg2 == 1; var numFields = input.parseState.Count; diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/BlockAllocate.cs b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/BlockAllocate.cs index cb5ea67e16a..2a27706286c 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/BlockAllocate.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/BlockAllocate.cs @@ -60,6 +60,17 @@ bool TryAllocateRecord(TSes if ((minRevivAddress <= stackCtx.hei.Address) && (!options.elideSourceRecord || stackCtx.hei.Address != stackCtx.recSrc.LogicalAddress)) minRevivAddress = stackCtx.hei.Address; + // A (v+1) record must never land below the checkpoint's fuzzy-region start: the snapshot includes every record + // below that address regardless of its version bit, so reusing a lower address would leak (v+1) data into the + // checkpoint image. This applies to both reuse paths below, including an allocation that was saved before the + // version shift and is only now being retried as (v+1). + if (sessionFunctions.Ctx.IsInV1) + { + var fuzzyStartAddress = _hybridLogCheckpoint.info.startLogicalAddress; + if (fuzzyStartAddress > minRevivAddress) + minRevivAddress = fuzzyStartAddress; + } + if (options.recycle && operationState.retryNewLogicalAddress != kInvalidAddress && GetAllocationForRetry(sessionFunctions, ref operationState, minRevivAddress, in sizeInfo, out newLogicalAddress, out newPhysicalAddress)) { @@ -68,12 +79,6 @@ bool TryAllocateRecord(TSes } if (RevivificationManager.UseFreeRecordPool) { - if (sessionFunctions.Ctx.IsInV1) - { - var fuzzyStartAddress = _hybridLogCheckpoint.info.startLogicalAddress; - if (fuzzyStartAddress > minRevivAddress) - minRevivAddress = fuzzyStartAddress; - } if (TryTakeFreeRecord(sessionFunctions, in sizeInfo, minRevivAddress, out newLogicalAddress, out newPhysicalAddress)) { new LogRecord(newPhysicalAddress).PrepareForRevivification(ref sizeInfo); diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/Helpers.cs b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/Helpers.cs index c5e437c85a2..dc7a439d25d 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/Helpers.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/Helpers.cs @@ -106,6 +106,27 @@ private bool IsFrozen(TSess internal long GetMinRevivifiableAddress() => RevivificationManager.GetMinRevivifiableAddress(hlogBase.GetTailAddress(), hlogBase.ReadOnlyAddress); + /// + /// Dispose the resources of an in-memory source record that a newly-CAS'd record has just superseded, unless an + /// ongoing checkpoint has frozen it. + /// + /// + /// Disposal clears the record's heap fields, which returns the value's slot to that page's + /// free list for reuse by another record. The snapshot flush reads object ids from its page copy but resolves them + /// against the live map, so disposing a frozen record lets the flush serialize a freed - or recycled, and therefore + /// unrelated - object in its place. A frozen record must keep its value until the checkpoint has captured it; the + /// value is then accounted for and released when the page is evicted. + /// + [MethodImpl(MethodImplOptions.AggressiveInlining)] + private void OnDisposeSupersededSource(TSessionFunctionsWrapper sessionFunctions, + ref OperationStackContext stackCtx, ref LogRecord logRecord) + where TSessionFunctionsWrapper : ISessionFunctionsWrapper + { + if (IsFrozen(sessionFunctions, ref stackCtx, logRecord.Info)) + return; + OnDispose(ref logRecord, DisposeReason.Deleted); + } + [MethodImpl(MethodImplOptions.NoInlining)] private (bool elided, bool added) TryElideAndTransferToFreeList(TSessionFunctionsWrapper sessionFunctions, ref OperationStackContext stackCtx, ref LogRecord logRecord) diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalDelete.cs b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalDelete.cs index a11a8fed589..e26b3048178 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalDelete.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalDelete.cs @@ -217,11 +217,14 @@ private OperationStatus CreateNewRecordDelete(sessionFunctions, ref stackCtx, srcLogRecord.Info) }; - // We know the existing record cannot be elided; it must point to a valid record; otherwise InternalDelete would have returned NOTFOUND. + // If the source record is elidable we will detach it from the tag chain below, after the CAS, by carrying its + // PreviousAddress on the new tombstone. CanElide() excludes checkpoint-frozen records. if (!TryAllocateRecord(sessionFunctions, ref operationState, ref stackCtx, ref sizeInfo, allocOptions, out var newLogicalAddress, out var newPhysicalAddress, out var status)) return status; var newLogRecord = WriteNewRecordInfo(key, hlogBase, newLogicalAddress, newPhysicalAddress, in sizeInfo, sessionFunctions.Ctx.InNewVersion, previousAddress: stackCtx.recSrc.LatestLogicalAddress); + if (allocOptions.elideSourceRecord) + newLogRecord.InfoRef.PreviousAddress = srcLogRecord.Info.PreviousAddress; newLogRecord.InfoRef.SetTombstone(); stackCtx.SetNewRecord(newLogicalAddress); @@ -256,10 +259,27 @@ private OperationStatus CreateNewRecordDelete= GetMinRevivifiableAddress()) + _ = TryTransferToFreeList(sessionFunctions, stackCtx.recSrc.LogicalAddress, ref srcLogRecord); + else + OnDispose(ref srcLogRecord, DisposeReason.Elided); + } + else if (stackCtx.recSrc.HasMainLogSrc) + { + // Dispose the superseded source record's resources, unless a checkpoint has frozen it. + OnDisposeSupersededSource(sessionFunctions, ref stackCtx, ref srcLogRecord); srcLogRecord.InfoRef.Seal(); // Not elided so Seal without invalidate } diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalRMW.cs b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalRMW.cs index 8b463401cd8..d06a146660d 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalRMW.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Tsavorite/Implementation/InternalRMW.cs @@ -413,17 +413,17 @@ private OperationStatus CreateNewRecordRMW(sessionFunctions, ref stackCtx, ref srcLogRecord.AsMemoryLogRecordRef()); doingCU = false; forExpiration = true; } else if (rmwInfo.Action == RMWAction.ExpireAndStop) { - // Immediately dispose all resources on the expired source record. + // Immediately dispose all resources on the expired source record, unless frozen by a checkpoint. if (stackCtx.recSrc.HasMainLogSrc) - OnDispose(ref srcLogRecord.AsMemoryLogRecordRef(), DisposeReason.Deleted); + OnDisposeSupersededSource(sessionFunctions, ref stackCtx, ref srcLogRecord.AsMemoryLogRecordRef()); if (allocOptions.elideSourceRecord) { @@ -522,9 +522,9 @@ private OperationStatus CreateNewRecordRMW(sessionFunctions, ref stackCtx, ref srcLogRecord.AsMemoryLogRecordRef()); addTombstone = true; newLogRecord.InfoRef.SetTombstone(); newLogRecord.InfoRef.SetModified(); @@ -533,9 +533,9 @@ private OperationStatus CreateNewRecordRMW(sessionFunctions, ref stackCtx, ref srcLogRecord.AsMemoryLogRecordRef()); doingCU = false; forExpiration = true; diff --git a/libs/storage/Tsavorite/cs/test/test.recordops/RecordLifecycleTests.cs b/libs/storage/Tsavorite/cs/test/test.recordops/RecordLifecycleTests.cs index 9109161d43f..4dfd14bff38 100644 --- a/libs/storage/Tsavorite/cs/test/test.recordops/RecordLifecycleTests.cs +++ b/libs/storage/Tsavorite/cs/test/test.recordops/RecordLifecycleTests.cs @@ -470,9 +470,9 @@ public void PendingReadFromDiskFiresOnDisposeDiskRecordOnce() /// /// Filling the log well beyond its mutable window forces page eviction. OnEvict must fire for - /// every non-tombstoned, non-invalid record evicted past HeadAddress — including sealed source - /// records from immutable-region deletes. Tombstoned records are skipped (heap was decremented - /// at the delete site). Invalid/elided records are skipped (already cleaned up). + /// every non-tombstoned, non-invalid record evicted past HeadAddress. Tombstoned records are + /// skipped (heap was decremented at the delete site). Invalid records are skipped, including the + /// source records that an immutable-region delete elided from the tag chain and already cleaned up. /// [Test, Category("TsavoriteKV")] public void PageEvictionFiresOnEvictForEveryLiveRecord() @@ -501,18 +501,26 @@ public void PageEvictionFiresOnEvictForEveryLiveRecord() "Each Delete should fire OnDispose(Deleted) exactly once"); var deletedDisposeCountBeforeEvict = tracker.DisposeCount(DisposeReason.Deleted); + // An immutable-region delete whose source is the only record in its tag chain elides that source: + // it is invalidated and either freelisted or disposed, so eviction will never visit it. + var elidedSources = tracker.DisposeCount(DisposeReason.Elided) + tracker.DisposeCount(DisposeReason.RevivificationFreeList); + ClassicAssert.LessOrEqual(elidedSources, immutableDeletes, + "Only immutable-delete source records can be elided"); + // Force all records out to disk. store.Log.FlushAndEvict(wait: true); // Precise count: // - (n - deleted) live records: visited by OnEvict. // - mutableDeletes records: tombstoned in-place, skipped by OnEvict. - // - immutableDeletes sealed source records: NOT tombstoned, visited by OnEvict. + // - elidedSources records: invalidated and cleaned up at the delete site, skipped by OnEvict. + // - the remaining immutableDeletes sealed source records: NOT tombstoned, visited by OnEvict. // - immutableDeletes new tombstone records at tail: tombstoned, skipped by OnEvict. - // Total = (n - deleted) + immutableDeletes = n - mutableDeletes. - ClassicAssert.AreEqual(n - mutableDeletes, tracker.EvictCount(EvictionSource.MainLog), - $"OnEvict(MainLog) must fire exactly {n - mutableDeletes} times: " + - $"{n - deleted} live + {immutableDeletes} sealed sources, skipping {mutableDeletes} in-place tombstones"); + var expectedEvictions = (n - deleted) + (immutableDeletes - elidedSources); + ClassicAssert.AreEqual(expectedEvictions, tracker.EvictCount(EvictionSource.MainLog), + $"OnEvict(MainLog) must fire exactly {expectedEvictions} times: " + + $"{n - deleted} live + {immutableDeletes - elidedSources} sealed sources, " + + $"skipping {mutableDeletes} in-place tombstones and {elidedSources} elided sources"); ClassicAssert.AreEqual(0, tracker.EvictCount(EvictionSource.ReadCache), "No read cache is configured, OnEvict(ReadCache) must never fire"); ClassicAssert.AreEqual(deletedDisposeCountBeforeEvict, tracker.DisposeCount(DisposeReason.Deleted), From f13aa834330f429fb11f2eef6818a87238c63655 Mon Sep 17 00:00:00 2001 From: TedHartMS <15467143+TedHartMS@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:34:04 -0700 Subject: [PATCH 4/6] Add regression tests for the HeapObjectBase serialization fixes CleanupOfCachedDataLeavesObjectTerminal covers #2101: it drives the CopyUpdate-during- checkpoint path that caches a superseded object's (v) bytes, runs the post-checkpoint cleanup that releases them, and asserts the object stays terminal. Returning it to REST let the next checkpoint take the direct-serialize path and enumerate the collections the superseding (v+1) record still shares and mutates. FailedSerializationRestoresRestPhase covers the exception-safety half: a DoSerialize that throws must leave the object in REST, not stranded in SERIALIZING where every later Serialize() and CacheSerializedObjectData() would spin on it. Both fail against the pre-fix code with exactly these diagnostics. --- .../test/test.recovery/ObjectLogScanTests.cs | 109 ++++++++++++++++++ 1 file changed, 109 insertions(+) diff --git a/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs b/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs index d3c1f39bf27..408629d18ef 100644 --- a/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs +++ b/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs @@ -138,6 +138,31 @@ public override void DoSerialize(BinaryWriter writer) public override void WriteType(BinaryWriter writer, bool isNull) => writer.Write(isNull); } + sealed class CountingSerializationHeapObject : HeapObjectBase + { + internal const string SerializeFailureMessage = "Injected DoSerialize failure"; + + // Counts entries into the direct-serialize path, and records which indicator Serialize() wrote. + internal int doSerializeCount; + internal bool? lastWriteTypeIsNull; + internal bool throwOnSerialize; + + public override IHeapObject Clone() => new CountingSerializationHeapObject(); + public override void Dispose() { } + public override void DoSerialize(BinaryWriter writer) + { + _ = Interlocked.Increment(ref doSerializeCount); + if (throwOnSerialize) + throw new InvalidOperationException(SerializeFailureMessage); + writer.Write(0); + } + public override void WriteType(BinaryWriter writer, bool isNull) + { + lastWriteTypeIsNull = isNull; + writer.Write(isNull); + } + } + [SetUp] public void Setup() { @@ -260,6 +285,90 @@ public async Task CleanupWithoutCachedDataTest() await serializationTask.ConfigureAwait(false); } + /// + /// Regression test for #2101. A CopyUpdate during a checkpoint caches the superseded (v) object's bytes + /// and leaves it SERIALIZED; post-checkpoint cleanup then releases those bytes. Because Clone() is a + /// shallow copy, the (v+1) record that superseded it shares and keeps mutating its collections, so the + /// object must stay terminal. Returning it to REST let the next checkpoint take the direct-serialize + /// path and enumerate those live collections, throwing "Collection was modified" out of DoSerialize and + /// wedging the checkpoint pipeline. + /// + [Test] + [Category("TsavoriteKV")] + public void CleanupOfCachedDataLeavesObjectTerminal() + { + log = Devices.CreateLogDevice(Path.Join(MethodTestDir, "CachedDataCleanup.log"), deleteOnClose: true); + objlog = Devices.CreateLogDevice(Path.Join(MethodTestDir, "CachedDataCleanup.obj.log"), deleteOnClose: true); + store = new(new() + { + IndexSize = 1L << 13, + LogDevice = log, + ObjectLogDevice = objlog, + MutableFraction = 0.1, + LogMemorySize = 1L << 15, + PageSize = MinKvLogPageSize + }, StoreFunctions.Create(comparer, () => new TrackingHeapObjectSerializer()) + , (allocatorSettings, storeFunctions) => new(allocatorSettings, storeFunctions) + ); + + using var session = store.NewSession(new TestObjectFunctions()); + var context = session.BasicContext; + var value = new CountingSerializationHeapObject(); + + var beginAddress = store.Log.TailAddress; + _ = context.Upsert(new OverflowTestKey(1), value, Empty.Default); + var endAddress = store.Log.TailAddress; + + // Drive the CopyUpdate-during-checkpoint path that caches the (v) bytes: marking the record as the + // new version is what makes CacheSerializedObjectData capture rather than just hand off the object. + var logRecord = store.hlogBase._wrapper.CreateLogRecord(beginAddress); + logRecord.InfoRef.SetIsInNewVersion(); + RMWInfo rmwInfo = default; + value.CacheSerializedObjectData(ref logRecord, ref rmwInfo, srcIsOnMemoryLog: true); + + // Caching serializes once, into the cached byte[] rather than to a writer. + Assert.That(value.doSerializeCount, Is.EqualTo(1), "CacheSerializedObjectData did not capture the (v) bytes"); + + // Post-checkpoint cleanup releases the cached bytes. + store.Log.ClearSerializedObjectData(beginAddress, endAddress); + + // This transition can only succeed if cleanup incorrectly reset the phase to REST. + Assert.That(value.MakeTransition(SerializationPhase.REST, SerializationPhase.SERIALIZING), Is.False, + "Cleanup returned a superseded object to REST"); + + using var stream = new MemoryStream(); + using var writer = new BinaryWriter(stream); + value.Serialize(writer); + + // SERIALIZED with no cached bytes means the object was superseded after the checkpoint completed: + // the superseding record sits at a higher address and carries the live data, so this one writes null. + Assert.That(value.doSerializeCount, Is.EqualTo(1), "A superseded object was re-serialized from its live state"); + Assert.That(value.lastWriteTypeIsNull, Is.True, "A superseded object with no cached bytes must write the null indicator"); + } + + /// + /// A failing DoSerialize must not strand the object in SERIALIZING: every later Serialize() and + /// CacheSerializedObjectData() spins waiting for that phase to clear, so the checkpoint pipeline would + /// hang rather than report the failure. + /// + [Test] + [Category("TsavoriteKV")] + public void FailedSerializationRestoresRestPhase() + { + var value = new CountingSerializationHeapObject { throwOnSerialize = true }; + + using var stream = new MemoryStream(); + using var writer = new BinaryWriter(stream); + + var failure = Assert.Throws(() => value.Serialize(writer)); + Assert.That(failure.Message, Is.EqualTo(CountingSerializationHeapObject.SerializeFailureMessage)); + Assert.That(value.doSerializeCount, Is.EqualTo(1)); + + // Only succeeds if the failed serialization restored REST. + Assert.That(value.MakeTransition(SerializationPhase.REST, SerializationPhase.SERIALIZING), Is.True, + "A failed serialization stranded the object outside REST"); + } + internal struct ObjectPushScanTestFunctions : IScanIteratorFunctions { internal long numRecords; From 2598100b54887ec6c1c29885b60ab755ee03cefa Mon Sep 17 00:00:00 2001 From: TedHartMS <15467143+TedHartMS@users.noreply.github.com> Date: Tue, 15 Sep 2026 15:50:42 -0700 Subject: [PATCH 5/6] Cover the cached-serialization failure path FailedSerializationRestoresRestPhase only calls Serialize, so it exercises the direct path's finally. The catch in CacheSerializedObjectData, which covers the other DoSerialize caller, had no test: a regression there could strand the object in SERIALIZING, hanging every later serialization attempt on it. FailedCachedSerializationRestoresRestPhase drives the CopyUpdate-during-checkpoint capture with a throwing DoSerialize and asserts the exception is rethrown, REST is restored, and no partial capture is published. It fails with 'stranded the object outside REST' if the catch is removed. The store setup shared with CleanupOfCachedDataLeavesObjectTerminal moves into a helper. --- .../test/test.recovery/ObjectLogScanTests.cs | 85 +++++++++++++++---- 1 file changed, 68 insertions(+), 17 deletions(-) diff --git a/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs b/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs index 408629d18ef..0a78fb37660 100644 --- a/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs +++ b/libs/storage/Tsavorite/cs/test/test.recovery/ObjectLogScanTests.cs @@ -286,19 +286,14 @@ public async Task CleanupWithoutCachedDataTest() } /// - /// Regression test for #2101. A CopyUpdate during a checkpoint caches the superseded (v) object's bytes - /// and leaves it SERIALIZED; post-checkpoint cleanup then releases those bytes. Because Clone() is a - /// shallow copy, the (v+1) record that superseded it shares and keeps mutating its collections, so the - /// object must stay terminal. Returning it to REST let the next checkpoint take the direct-serialize - /// path and enumerate those live collections, throwing "Collection was modified" out of DoSerialize and - /// wedging the checkpoint pipeline. + /// Creates the store used by the cached-serialization tests, upserts , and marks its + /// record as the new version, which is what makes CacheSerializedObjectData capture the (v) bytes rather + /// than simply hand the object off. Returns the address range covering that single record. /// - [Test] - [Category("TsavoriteKV")] - public void CleanupOfCachedDataLeavesObjectTerminal() + private (long beginAddress, long endAddress) CreateStoreWithCachedSerializationSource(string name, IHeapObject value, out long recordAddress) { - log = Devices.CreateLogDevice(Path.Join(MethodTestDir, "CachedDataCleanup.log"), deleteOnClose: true); - objlog = Devices.CreateLogDevice(Path.Join(MethodTestDir, "CachedDataCleanup.obj.log"), deleteOnClose: true); + log = Devices.CreateLogDevice(Path.Join(MethodTestDir, $"{name}.log"), deleteOnClose: true); + objlog = Devices.CreateLogDevice(Path.Join(MethodTestDir, $"{name}.obj.log"), deleteOnClose: true); store = new(new() { IndexSize = 1L << 13, @@ -313,16 +308,34 @@ public void CleanupOfCachedDataLeavesObjectTerminal() using var session = store.NewSession(new TestObjectFunctions()); var context = session.BasicContext; - var value = new CountingSerializationHeapObject(); var beginAddress = store.Log.TailAddress; _ = context.Upsert(new OverflowTestKey(1), value, Empty.Default); var endAddress = store.Log.TailAddress; - // Drive the CopyUpdate-during-checkpoint path that caches the (v) bytes: marking the record as the - // new version is what makes CacheSerializedObjectData capture rather than just hand off the object. + recordAddress = beginAddress; var logRecord = store.hlogBase._wrapper.CreateLogRecord(beginAddress); logRecord.InfoRef.SetIsInNewVersion(); + return (beginAddress, endAddress); + } + + /// + /// Regression test for #2101. A CopyUpdate during a checkpoint caches the superseded (v) object's bytes + /// and leaves it SERIALIZED; post-checkpoint cleanup then releases those bytes. Because Clone() is a + /// shallow copy, the (v+1) record that superseded it shares and keeps mutating its collections, so the + /// object must stay terminal. Returning it to REST let the next checkpoint take the direct-serialize + /// path and enumerate those live collections, throwing "Collection was modified" out of DoSerialize and + /// wedging the checkpoint pipeline. + /// + [Test] + [Category("TsavoriteKV")] + public void CleanupOfCachedDataLeavesObjectTerminal() + { + var value = new CountingSerializationHeapObject(); + var (beginAddress, endAddress) = CreateStoreWithCachedSerializationSource("CachedDataCleanup", value, out var recordAddress); + + // Drive the CopyUpdate-during-checkpoint path that caches the (v) bytes. + var logRecord = store.hlogBase._wrapper.CreateLogRecord(recordAddress); RMWInfo rmwInfo = default; value.CacheSerializedObjectData(ref logRecord, ref rmwInfo, srcIsOnMemoryLog: true); @@ -347,9 +360,9 @@ public void CleanupOfCachedDataLeavesObjectTerminal() } /// - /// A failing DoSerialize must not strand the object in SERIALIZING: every later Serialize() and - /// CacheSerializedObjectData() spins waiting for that phase to clear, so the checkpoint pipeline would - /// hang rather than report the failure. + /// A failing DoSerialize on the direct path (Serialize straight to the wire) must not strand the object + /// in SERIALIZING: every later Serialize() and CacheSerializedObjectData() spins waiting for that phase + /// to clear, so the checkpoint pipeline would hang rather than report the failure. /// [Test] [Category("TsavoriteKV")] @@ -369,6 +382,44 @@ public void FailedSerializationRestoresRestPhase() "A failed serialization stranded the object outside REST"); } + /// + /// The same guarantee for the other DoSerialize caller: the CopyUpdate-during-checkpoint path that captures + /// the (v) bytes. A failure there must rethrow, restore REST, and publish no partial capture; stranding the + /// object in SERIALIZING would hang every later serialization attempt on it. + /// + [Test] + [Category("TsavoriteKV")] + public void FailedCachedSerializationRestoresRestPhase() + { + var value = new CountingSerializationHeapObject { throwOnSerialize = true }; + _ = CreateStoreWithCachedSerializationSource("FailedCachedSerialization", value, out var recordAddress); + + var failure = Assert.Throws(() => + { + // The ref local must be created inside the lambda; ref locals cannot be captured. + var record = store.hlogBase._wrapper.CreateLogRecord(recordAddress); + RMWInfo info = default; + value.CacheSerializedObjectData(ref record, ref info, srcIsOnMemoryLog: true); + }); + Assert.That(failure.Message, Is.EqualTo(CountingSerializationHeapObject.SerializeFailureMessage)); + Assert.That(value.doSerializeCount, Is.EqualTo(1)); + + // Only succeeds if the failed capture restored REST rather than leaving SERIALIZING behind. + Assert.That(value.MakeTransition(SerializationPhase.REST, SerializationPhase.SERIALIZING), Is.True, + "A failed cached serialization stranded the object outside REST"); + Assert.That(value.MakeTransition(SerializationPhase.SERIALIZING, SerializationPhase.REST), Is.True); + + // No partial capture was published: a later Serialize re-serializes from the live object and writes it + // as present, rather than emitting a truncated cached buffer or the superseded null indicator. + value.throwOnSerialize = false; + using var stream = new MemoryStream(); + using var writer = new BinaryWriter(stream); + value.Serialize(writer); + + Assert.That(value.doSerializeCount, Is.EqualTo(2), "The object did not re-serialize from its live state"); + Assert.That(value.lastWriteTypeIsNull, Is.False, "A failed capture must not leave the object looking superseded"); + } + internal struct ObjectPushScanTestFunctions : IScanIteratorFunctions { internal long numRecords; From baa0f12dd2f4624fcdce76bf0af9179cd2a90ce9 Mon Sep 17 00:00:00 2001 From: TedHartMS <15467143+TedHartMS@users.noreply.github.com> Date: Sun, 20 Sep 2026 13:08:24 -0700 Subject: [PATCH 6/6] Reproduce #2101 end to end, and bound the HTTL retention claim The existing ObjectLogScanTests coverage asserts the serialization-phase invariant by calling CacheSerializedObjectData directly on an object that holds no collection. That is a valid unit test of the SAD protocol and it does fail without the fix, but it never runs InternalRMW, never has a checkpoint active, and can never produce 'Collection was modified'. It therefore does not reproduce the reported bug. CheckpointSharedObjectMutationTests drives the real path: a snapshot checkpoint is paused on entry to WAIT_FLUSH, an RMW in that window takes the CPR copy-update route through CreateNewRecordRMW into CacheSerializedObjectData, and the clone shares the source's List exactly as Garnet's shallow Clone() does. Post-checkpoint cleanup then runs, the surviving record mutates the shared list, and the superseded object is serialized again while a second mutation lands mid-enumeration. Against the pre-fix allocator that fails with InvalidOperationException: Collection was modified; enumeration operation may not execute. which is the exception from the issue, raised from the same direct serialize path. It passes with the fix. HashTimeToLive gains a comment recording why dropping DeleteExpiredItems does not leak without bound: a field only becomes expirable through HEXPIRE/HPEXPIRE, which is an RMW that purges first, so the expired-but-unpurged set cannot exceed the fields carrying a TTL at the last mutating operation and a read-only workload cannot grow it. --- libs/server/Objects/Hash/HashObjectImpl.cs | 6 + .../CheckpointSharedObjectMutationTests.cs | 251 ++++++++++++++++++ 2 files changed, 257 insertions(+) create mode 100644 libs/storage/Tsavorite/cs/test/test.recovery/CheckpointSharedObjectMutationTests.cs diff --git a/libs/server/Objects/Hash/HashObjectImpl.cs b/libs/server/Objects/Hash/HashObjectImpl.cs index bbb67db7275..f33b77597e3 100644 --- a/libs/server/Objects/Hash/HashObjectImpl.cs +++ b/libs/server/Objects/Hash/HashObjectImpl.cs @@ -459,6 +459,12 @@ private void HashTimeToLive(ref ObjectInput input, ref ObjectOutput output, byte // This is a read operation, so it must not mutate the object: it runs under a shared lock and may execute // while the flush path is serializing this same instance. Expired fields are reported as absent by // ContainsKey (via GetExpiration) without being removed; the mutating paths purge them from the live object. + // + // Retention is bounded, not unbounded: a field can only become expirable through HEXPIRE/HPEXPIRE, which is + // an RMW that itself calls DeleteExpiredItems() first. So the expired-but-unpurged set can never exceed the + // fields that were carrying a TTL at the last mutating operation, and a read-only workload cannot grow it. + // That memory is reclaimed by the next mutating operation, by HCOLLECT, or by the background collector + // (ExpiredObjectCollectionFrequencySecs). This matches SortedSetTimeToLive (ZTTL), which has never purged. var isMilliseconds = input.arg1 == 1; var isTimestamp = input.arg2 == 1; var numFields = input.parseState.Count; diff --git a/libs/storage/Tsavorite/cs/test/test.recovery/CheckpointSharedObjectMutationTests.cs b/libs/storage/Tsavorite/cs/test/test.recovery/CheckpointSharedObjectMutationTests.cs new file mode 100644 index 00000000000..125233bdbe5 --- /dev/null +++ b/libs/storage/Tsavorite/cs/test/test.recovery/CheckpointSharedObjectMutationTests.cs @@ -0,0 +1,251 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT license. + +using System; +using System.Collections.Generic; +using System.IO; +using System.Threading; +using System.Threading.Tasks; +using Garnet.test; +using NUnit.Framework; +using Tsavorite.core; +using static Tsavorite.test.TestUtils; + +namespace Tsavorite.test.recovery +{ + using ClassAllocator = ObjectAllocator>; + using ClassStoreFunctions = StoreFunctions; + + /// + /// End-to-end regression test for #2101, driving the real paths: an RMW that performs a CopyUpdate while a + /// checkpoint is active, the post-checkpoint cleanup that releases the cached bytes, and a later flush that + /// serializes the superseded record while the surviving record mutates the collection they share. + /// + [TestFixture] + internal class CheckpointSharedObjectMutationTests : TestBase + { + private TsavoriteKV store; + private IDevice log, objlog; + + /// + /// Models a Garnet collection object: is a SHALLOW copy, so the (v+1) record created by + /// a CopyUpdate shares this instance's list. enumerates that shared list, and can + /// be paused mid-enumeration so a test can mutate it and reproduce "Collection was modified". + /// + internal sealed class SharedListHeapObject : HeapObjectBase + { + internal readonly List Items; + + /// Set to pause inside the enumeration so the caller can mutate . + internal ManualResetEventSlim PauseDuringSerialize; + internal readonly ManualResetEventSlim ReachedSerialize = new(false); + + internal int DoSerializeCount; + + internal SharedListHeapObject(List items) + { + Items = items; + HeapMemorySize = 64; + } + + // Shallow copy: this is the crux of #2101 - the clone shares Items with the record it supersedes. + public override IHeapObject Clone() => new SharedListHeapObject(Items); + + public override void Dispose() { } + + public override void DoSerialize(BinaryWriter writer) + { + _ = Interlocked.Increment(ref DoSerializeCount); + writer.Write(Items.Count); + + // foreach over the shared list: throws InvalidOperationException if another thread mutates it. + var first = true; + foreach (var item in Items) + { + if (first && PauseDuringSerialize is not null) + { + first = false; + ReachedSerialize.Set(); + PauseDuringSerialize.Wait(TimeSpan.FromSeconds(10)); + } + writer.Write(item); + } + } + + public override void WriteType(BinaryWriter writer, bool isNull) => writer.Write(isNull); + } + + internal sealed class SharedListSerializer : BinaryObjectSerializer + { + public override void Deserialize(out IHeapObject obj) + { + var count = reader.ReadInt32(); + var items = new List(count); + for (var i = 0; i < count; i++) + items.Add(reader.ReadInt32()); + obj = new SharedListHeapObject(items); + } + + public override void Serialize(IHeapObject obj) => ((SharedListHeapObject)obj).DoSerialize(writer); + } + + /// Pauses the checkpoint state machine on entry to a chosen phase so the test can act in that window. + private sealed class PauseAtPhase(Phase pauseAt) : IStateMachineCallback + { + internal readonly ManualResetEventSlim Reached = new(false); + internal readonly ManualResetEventSlim Release = new(false); + + public void BeforeEnteringState(SystemState next) + { + if (next.Phase != pauseAt) + return; + Reached.Set(); + _ = Release.Wait(TimeSpan.FromSeconds(30)); + } + } + + [SetUp] + public void Setup() + { + DeleteDirectory(MethodTestDir, wait: true); + log = Devices.CreateLogDevice(Path.Join(MethodTestDir, "SharedObj.log"), deleteOnClose: true); + objlog = Devices.CreateLogDevice(Path.Join(MethodTestDir, "SharedObj.obj.log"), deleteOnClose: true); + store = new(new() + { + IndexSize = 1L << 13, + LogDevice = log, + ObjectLogDevice = objlog, + MutableFraction = 0.1, + LogMemorySize = 1L << 16, + PageSize = 1L << 13, + CheckpointDir = MethodTestDir + }, StoreFunctions.Create(new TestObjectKey.Comparer(), () => new SharedListSerializer()) + , (allocatorSettings, storeFunctions) => new(allocatorSettings, storeFunctions)); + } + + [TearDown] + public void TearDown() + { + store?.Dispose(); + store = null; + log?.Dispose(); + log = null; + objlog?.Dispose(); + objlog = null; + OnTearDown(); + } + + [Test] + [Category("TsavoriteKV")] + [Category("CheckpointRestore")] + public async Task SupersededObjectIsNotSerializedFromLiveStateAfterCleanup() + { + var key = new TestObjectKey { key = 1 }; + var shared = new List { 1, 2, 3, 4, 5, 6, 7, 8 }; + var original = new SharedListHeapObject(shared); + + using (var session = store.NewSession(new SharedListFunctions())) + _ = session.BasicContext.Upsert(key, original, Empty.Default); + + var sourceAddress = store.Log.TailAddress - 1; + + // Checkpoint C1: pause on entry to WAIT_FLUSH, which is after IN_PROGRESS has been published and the + // fuzzy region has opened, so an RMW in this window is (v+1) against a (v) source and must RCU. + var pause = new PauseAtPhase(Phase.WAIT_FLUSH); + store.stateMachineDriver.UnsafeRegisterCallback(pause); + + Assert.That(store.TryInitiateFullCheckpoint(out _, CheckpointType.Snapshot), Is.True); + Assert.That(pause.Reached.Wait(TimeSpan.FromSeconds(30)), Is.True, "Checkpoint did not reach WAIT_FLUSH"); + + // Real RMW -> CreateNewRecordRMW -> CacheSerializedObjectData on the superseded (v) object. + using (var session = store.NewSession(new SharedListFunctions())) + { + TestObjectInput input = new() { value = 99 }; + TestObjectOutput output = new(); + var status = session.BasicContext.RMW(key, ref input, ref output); + if (status.IsPending) + _ = session.BasicContext.CompletePending(wait: true); + } + + pause.Release.Set(); + await store.CompleteCheckpointAsync().ConfigureAwait(false); + + // The RMW must actually have gone through the CopyUpdate-during-checkpoint path. + Assert.That(original.DoSerializeCount, Is.GreaterThanOrEqualTo(1), + "The RMW did not cache the superseded object's (v) bytes; the test is not exercising the #2101 path"); + var afterCheckpoint = original.DoSerializeCount; + + // Post-checkpoint cleanup, exactly as Garnet's RunPostCheckpointCleanup does. + store.Log.ClearSerializedObjectData(store.Log.BeginAddress, store.Log.TailAddress); + + // The surviving (v+1) record shares this list and keeps mutating it. + shared.Add(1000); + + // Serialize the superseded record again, pausing mid-enumeration so we can mutate concurrently. + // Pre-fix the phase was reset to REST, so this takes the direct path and enumerates the live list. + original.PauseDuringSerialize = new ManualResetEventSlim(false); + Exception serializeFailure = null; + var serializeTask = Task.Run(() => + { + try + { + using var ms = new MemoryStream(); + using var writer = new BinaryWriter(ms); + original.Serialize(writer); + } + catch (Exception ex) + { + serializeFailure = ex; + } + }); + + if (original.ReachedSerialize.Wait(TimeSpan.FromSeconds(2))) + { + // Only reachable if Serialize took the direct path: mutate while the enumerator is live. + shared.Add(2000); + original.PauseDuringSerialize.Set(); + } + await serializeTask.ConfigureAwait(false); + + Assert.Multiple(() => + { + Assert.That(original.DoSerializeCount, Is.EqualTo(afterCheckpoint), + "#2101: the superseded object was re-serialized from its live, shared collection after cleanup"); + Assert.That(serializeFailure, Is.Null, + $"#2101: serializing the superseded object threw {serializeFailure?.GetType().Name}: {serializeFailure?.Message}"); + }); + + _ = sourceAddress; + } + + internal class SharedListFunctions : SessionFunctionsBase + { + public override bool InitialUpdater(ref LogRecord dstLogRecord, in RecordSizeInfo sizeInfo, ref TestObjectInput input, ref TestObjectOutput output, ref RMWInfo rmwInfo) + => dstLogRecord.TrySetValueObject(new SharedListHeapObject([input.value])); + + public override bool InPlaceUpdater(ref LogRecord logRecord, ref TestObjectInput input, ref TestObjectOutput output, ref RMWInfo rmwInfo) + { + ((SharedListHeapObject)logRecord.ValueObject).Items.Add(input.value); + return true; + } + + public override bool CopyUpdater(in TSourceLogRecord srcLogRecord, ref LogRecord dstLogRecord, in RecordSizeInfo sizeInfo, ref TestObjectInput input, ref TestObjectOutput output, ref RMWInfo rmwInfo) + => true; + + public override bool PostCopyUpdater(in TSourceLogRecord srcLogRecord, ref LogRecord dstLogRecord, in RecordSizeInfo sizeInfo, ref TestObjectInput input, ref TestObjectOutput output, ref RMWInfo rmwInfo) + { + ((SharedListHeapObject)dstLogRecord.ValueObject).Items.Add(input.value); + return true; + } + + public override RecordFieldInfo GetRMWModifiedFieldInfo(in TSourceLogRecord srcLogRecord, ref TestObjectInput input) + => new() { KeySize = srcLogRecord.Key.Length, ValueSize = ObjectIdMap.ObjectIdSize, ValueIsObject = true }; + + public override RecordFieldInfo GetRMWInitialFieldInfo(TKey key, ref TestObjectInput input) + => new() { KeySize = key.KeyBytes.Length, ValueSize = ObjectIdMap.ObjectIdSize, ValueIsObject = true }; + + public override RecordFieldInfo GetUpsertFieldInfo(TKey key, IHeapObject value, ref TestObjectInput input) + => new() { KeySize = key.KeyBytes.Length, ValueSize = ObjectIdMap.ObjectIdSize, ValueIsObject = true }; + } + } +} \ No newline at end of file