diff --git a/libs/server/Databases/CheckpointStatus.cs b/libs/server/Databases/CheckpointStatus.cs new file mode 100644 index 00000000000..6fbf1cc55a5 --- /dev/null +++ b/libs/server/Databases/CheckpointStatus.cs @@ -0,0 +1,59 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT license. + +namespace Garnet.server +{ + /// + /// Outcome of a checkpoint request issued against one or more logical databases. + /// + public enum CheckpointStatus + { + /// + /// The checkpoint completed and its data is durable. For a background checkpoint this means the checkpoint + /// was successfully started; its outcome is reported through the database's last save time and status. + /// + Success, + + /// + /// A checkpoint is already in progress for at least one of the requested databases, so no new checkpoint + /// was started. Nothing was written and no existing checkpoint was invalidated. + /// + AlreadyInProgress, + + /// + /// The checkpoint was started but did not complete, so nothing durable was written. + /// + Failed, + } + + /// + /// Outcome of a single database's checkpoint attempt. + /// + internal readonly struct CheckpointResult + { + /// + /// True if the checkpoint completed and its data is durable. Only then may the database's last save time + /// be advanced; advancing it for a failed checkpoint reports data as durable that was never written. + /// + public bool IsSuccessful { get; init; } + + /// + /// Store tail address covered by a full checkpoint, or null for an incremental checkpoint. Null is also + /// returned for a failed checkpoint, which is why it cannot by itself signal failure. + /// + public long? StoreTailAddress { get; init; } + + /// + /// A result denoting a checkpoint that did not complete. + /// + public static CheckpointResult Failed => default; + + /// + /// Creates a result denoting a completed checkpoint. + /// + /// Store tail address covered by a full checkpoint, or null for an incremental checkpoint. + /// A successful result. + public static CheckpointResult Succeeded(long? storeTailAddress) => + new() { IsSuccessful = true, StoreTailAddress = storeTailAddress }; + } +} \ No newline at end of file diff --git a/libs/server/Databases/DatabaseManagerBase.cs b/libs/server/Databases/DatabaseManagerBase.cs index 93bab6e5431..2971dfa7497 100644 --- a/libs/server/Databases/DatabaseManagerBase.cs +++ b/libs/server/Databases/DatabaseManagerBase.cs @@ -38,7 +38,7 @@ internal abstract class DatabaseManagerBase : IDatabaseManager public abstract ValueTask RecoverCheckpointAsync(bool replicaRecover = false, bool recoverFromToken = false, CheckpointMetadata metadata = null); /// - public abstract Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null); + public abstract Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null); /// public abstract Task TakeOnDemandCheckpointAsync(DateTimeOffset entryTime, int dbId = 0); @@ -246,8 +246,13 @@ private static bool AofCanReconstructFromOrigin(GarnetDatabase db, out string st /// Database to checkpoint /// Logger /// Cancellation token - /// Tuple of store tail address and object store tail address - protected async Task TakeCheckpointAsync(GarnetDatabase db, ILogger logger = null, CancellationToken token = default) + /// The checkpoint outcome, including the store tail address covered by a full checkpoint + /// + /// Failures are reported through the returned rather than thrown, so that a + /// background checkpoint cannot tear down the server and so that one database's failure does not abort the + /// bookkeeping of the other databases checkpointed alongside it. + /// + protected async Task TakeCheckpointAsync(GarnetDatabase db, ILogger logger = null, CancellationToken token = default) { try { @@ -258,16 +263,42 @@ private static bool AofCanReconstructFromOrigin(GarnetDatabase db, out string st lastSaveStoreTailAddress - db.LastSaveStoreTailAddress >= StoreWrapper.serverOptions.FullCheckpointLogInterval; var checkpointType = StoreWrapper.serverOptions.UseFoldOverCheckpoints ? CheckpointType.FoldOver : CheckpointType.Snapshot; - await InitiateCheckpointAsync(db, full, checkpointType, logger).ConfigureAwait(false); + if (!await InitiateCheckpointAsync(db, full, checkpointType, logger).ConfigureAwait(false)) + return CheckpointResult.Failed; - return full ? lastSaveStoreTailAddress : null; + return CheckpointResult.Succeeded(full ? lastSaveStoreTailAddress : null); } catch (Exception ex) { - logger?.LogError(ex, "Checkpointing threw exception, DB ID: {id}", db.Id); + // The caller's logger is optional, so fall back to this manager's logger; otherwise a failed + // checkpoint leaves no trace at all. + (logger ?? Logger)?.LogError(ex, "Checkpointing threw exception, DB ID: {id}", db.Id); } - return null; + return CheckpointResult.Failed; + } + + /// + /// Record the outcome of a checkpoint attempt on the specified database + /// + /// Database that was checkpointed + /// Outcome of the checkpoint attempt + /// + /// The last save time is advanced only for a successful checkpoint. Advancing it for a failed one reports + /// data to clients (through LASTSAVE, and to the cluster through on-demand checkpointing) as durable when + /// nothing was written. + /// + protected static void RecordCheckpointOutcome(GarnetDatabase db, CheckpointResult result) + { + db.LastSaveSucceeded = result.IsSuccessful; + + if (!result.IsSuccessful) + return; + + if (result.StoreTailAddress.HasValue) + db.LastSaveStoreTailAddress = result.StoreTailAddress.Value; + + db.LastSaveTime = DateTimeOffset.UtcNow; } /// @@ -559,8 +590,8 @@ private ValueTask CompactionCommitAofAsync(GarnetDatabase db) /// True if full checkpoint should be initiated /// Type of checkpoint /// Logger - /// Task - private async Task InitiateCheckpointAsync(GarnetDatabase db, bool full, CheckpointType checkpointType, + /// True if the checkpoint ran to completion + private async Task InitiateCheckpointAsync(GarnetDatabase db, bool full, CheckpointType checkpointType, ILogger logger = null) { logger?.LogInformation("Initiating checkpoint; full = {full}, type = {checkpointType}, dbId = {dbId}", full, checkpointType, db.Id); @@ -594,6 +625,16 @@ private async Task InitiateCheckpointAsync(GarnetDatabase db, bool full, Checkpo checkpointResult.success = await db.StateMachineDriver.RunAsync(sm).ConfigureAwait(false); + if (!checkpointResult.success) + { + // Another state machine operation (such as an index resize) was already running, so the checkpoint + // never ran. Nothing was written, so the AOF must not be truncated and no checkpoint entry may be + // registered with the cluster - both would discard data this checkpoint does not cover. + (logger ?? Logger)?.LogWarning( + "Checkpoint did not run because another state machine operation is in progress, DB ID: {id}", db.Id); + return false; + } + // If cluster is enabled the replication manager is responsible for truncating AOF if (StoreWrapper.serverOptions.EnableCluster && StoreWrapper.serverOptions.EnableAOF) { @@ -612,6 +653,7 @@ private async Task InitiateCheckpointAsync(GarnetDatabase db, bool full, Checkpo logger ?? Logger); logger?.LogInformation("Completed checkpoint for DB ID: {id}", db.Id); + return true; } internal static void RunPostCheckpointCleanup(Action cleanup, int dbId, ILogger logger) diff --git a/libs/server/Databases/IDatabaseManager.cs b/libs/server/Databases/IDatabaseManager.cs index 31be04cc7fc..a879e0688c0 100644 --- a/libs/server/Databases/IDatabaseManager.cs +++ b/libs/server/Databases/IDatabaseManager.cs @@ -91,8 +91,12 @@ public interface IDatabaseManager : IDisposable /// ID of database to checkpoint, or -1 (default) to checkpoint all active databases /// Cancellation token /// Logger - /// False if another checkpointing process is already in progress - public Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null); + /// + /// if another checkpointing process is already in progress, + /// if a foreground checkpoint did not complete, otherwise + /// . A background checkpoint reports success once it has started. + /// + public Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null); /// /// Take a checkpoint if no checkpoint was taken after the provided time offset diff --git a/libs/server/Databases/MultiDatabaseManager.cs b/libs/server/Databases/MultiDatabaseManager.cs index cf647892100..7169eb59f73 100644 --- a/libs/server/Databases/MultiDatabaseManager.cs +++ b/libs/server/Databases/MultiDatabaseManager.cs @@ -131,6 +131,8 @@ public override async ValueTask RecoverCheckpointAsync(bool replicaRecover = fal if (ex.CandidateTokenCount == 0) { + // As in SingleDatabaseManager, nothing was ever written, so recovery cannot tell that a + // checkpoint the client was told had succeeded is missing; see RecordCheckpointOutcome. Logger?.LogInformation(ex, "No Hybrid Log found for recovery; storeVersion = {storeVersion}; objectStoreVersion = {objectStoreVersion}", storeVersion, objectStoreVersion); @@ -159,16 +161,17 @@ public override async ValueTask RecoverCheckpointAsync(bool replicaRecover = fal } /// - public override Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null) + public override Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null) { // Acquire databasesContentLock (read) so a concurrent swap-db can't move GarnetDatabase // wrappers out from under us mid-checkpoint (which would mis-attribute LASTSAVE to the // swapped DB and let a second BGSAVE race against the in-flight checkpoint). - if (!TryGetDatabasesContentReadLock(token)) return Task.FromResult(false); + if (!TryGetDatabasesContentReadLock(token)) return Task.FromResult(CheckpointStatus.AlreadyInProgress); var multiDbLockHeld = false; int[] pausedDbIds = null; var pausedCount = 0; + var requestedCount = 0; try { @@ -185,12 +188,13 @@ public override Task TakeCheckpointAsync(bool background, int dbId = -1, C if (!multiDbCheckpointingLock.TryWriteLock()) { databasesContentLock.ReadUnlock(); - return Task.FromResult(false); + return Task.FromResult(CheckpointStatus.AlreadyInProgress); } multiDbLockHeld = true; } + requestedCount = activeDbIdsMapSize; pausedDbIds = new int[activeDbIdsMapSize]; var activeDbIdsMapSnapshot = activeDbIds.Map; for (var i = 0; i < activeDbIdsMapSize; i++) @@ -209,11 +213,12 @@ public override Task TakeCheckpointAsync(bool background, int dbId = -1, C if (!TryPauseCheckpoints(dbId)) { databasesContentLock.ReadUnlock(); - return Task.FromResult(false); + return Task.FromResult(CheckpointStatus.AlreadyInProgress); } pausedDbIds = [dbId]; pausedCount = 1; + requestedCount = 1; } } catch @@ -231,10 +236,10 @@ public override Task TakeCheckpointAsync(bool background, int dbId = -1, C throw; } - var checkpointTask = RunPausedCheckpointsAndReleaseLocksAsync(pausedDbIds, pausedCount, multiDbLockHeld, token, logger); + var checkpointTask = RunPausedCheckpointsAndReleaseLocksAsync(pausedDbIds, pausedCount, requestedCount, multiDbLockHeld, token, logger); if (background) - return Task.FromResult(true); + return Task.FromResult(CheckpointStatus.Success); return checkpointTask; } @@ -259,8 +264,8 @@ public override async Task TakeOnDemandCheckpointAsync(DateTimeOffset entryTime, return; // Necessary to take a checkpoint because the latest checkpoint is before entryTime - var storeTailAddress = await TakeCheckpointAsync(db, logger: Logger).ConfigureAwait(false); - UpdateLastSaveData(dbId, storeTailAddress); + var result = await TakeCheckpointAsync(db, logger: Logger).ConfigureAwait(false); + UpdateLastSaveData(dbId, result); } finally { @@ -316,8 +321,8 @@ public override async Task TaskCheckpointBasedOnAofSizeLimitAsync(long aofSizeLi try { - var storeTailAddress = await TakeCheckpointAsync(databasesMapSnapshot[pausedDbId], logger: logger, token: token).ConfigureAwait(false); - UpdateLastSaveData(pausedDbId, storeTailAddress); + var result = await TakeCheckpointAsync(databasesMapSnapshot[pausedDbId], logger: logger, token: token).ConfigureAwait(false); + UpdateLastSaveData(pausedDbId, result); } finally { @@ -1018,8 +1023,15 @@ private void CopyDatabases(IDatabaseManager src, bool enableAof) /// individual one) so a per-DB BGSAVE issued mid-flight during a general BGSAVE reliably /// observes the in-progress checkpoint and fails with "checkpoint already in progress". /// - private async Task RunPausedCheckpointsAndReleaseLocksAsync(int[] pausedDbIds, int pausedCount, - bool multiDbLockHeld, CancellationToken token, ILogger logger) + /// Buffer whose first entries are pause-locked database IDs. + /// Number of databases this request pause-locked and will checkpoint. + /// Number of databases this request was asked to checkpoint, which exceeds + /// when a database was skipped because its checkpoint lock was already held. + /// Whether the caller holds . + /// Cancellation token. + /// Logger. + private async Task RunPausedCheckpointsAndReleaseLocksAsync(int[] pausedDbIds, int pausedCount, + int requestedCount, bool multiDbLockHeld, CancellationToken token, ILogger logger) { // Pre-fill with Task.CompletedTask so the catch path can safely await Task.WhenAll // even if the synchronous task-creation loop below throws partway through. @@ -1027,6 +1039,10 @@ private async Task RunPausedCheckpointsAndReleaseLocksAsync(int[] pausedDb for (var i = 0; i < pausedCount; i++) checkpointTasks[i] = Task.CompletedTask; + // Each checkpoint records its own outcome here rather than through its task's result, so a database + // whose task was never created or which threw stays counted as a failure. + var succeeded = new bool[pausedCount]; + try { // Force async so that the entry point can return synchronously to the caller. @@ -1037,7 +1053,7 @@ private async Task RunPausedCheckpointsAndReleaseLocksAsync(int[] pausedDb try { for (var i = 0; i < pausedCount; i++) - checkpointTasks[i] = TakeOneCheckpointAsync(databaseMapSnapshot[pausedDbIds[i]], pausedDbIds[i]); + checkpointTasks[i] = TakeOneCheckpointAsync(databaseMapSnapshot[pausedDbIds[i]], pausedDbIds[i], i); await Task.WhenAll(checkpointTasks).ConfigureAwait(false); } @@ -1063,28 +1079,48 @@ private async Task RunPausedCheckpointsAndReleaseLocksAsync(int[] pausedDb databasesContentLock.ReadUnlock(); } - return true; + var allSucceeded = true; + for (var i = 0; i < pausedCount; i++) + allSucceeded &= succeeded[i]; + + if (!allSucceeded) + return CheckpointStatus.Failed; + + // A database that could not be pause-locked already had a checkpoint in flight, so this request never + // attempted it and cannot vouch for it. Reporting success would tell a foreground SAVE that every + // requested database is on disk when one of them was skipped, and that skipped checkpoint may still + // fail. The skipped database records its own outcome through its own RecordCheckpointOutcome, so its + // LASTSAVE and rdb_last_bgsave_status stay truthful either way; this only stops the aggregate reply + // from claiming more than the request actually did. + // + // Background requests reply before this runs, so the BGSAVE contract of "skip the busy databases and + // report started" - asserted by MultiDatabaseSaveInProgressTest - is unaffected. + return pausedCount < requestedCount ? CheckpointStatus.AlreadyInProgress : CheckpointStatus.Success; // Local function: take one per-DB checkpoint and update LASTSAVE. Does NOT resume the // per-DB lock — the outer finally above resumes all paused DBs after WhenAll completes. - async Task TakeOneCheckpointAsync(GarnetDatabase db, int dbId) + async Task TakeOneCheckpointAsync(GarnetDatabase db, int dbId, int slot) { - var storeTailAddress = await TakeCheckpointAsync(db, logger: logger, token: token).ConfigureAwait(false); - UpdateLastSaveData(dbId, storeTailAddress); + var result = await TakeCheckpointAsync(db, logger: logger, token: token).ConfigureAwait(false); + UpdateLastSaveData(dbId, result); + succeeded[slot] = result.IsSuccessful; } } - private void UpdateLastSaveData(int dbId, long? storeTailAddress) + /// + /// Resolve the database for the given ID and record its checkpoint outcome + /// + /// ID of the database that was checkpointed + /// Outcome of the checkpoint attempt + /// + /// The database is read from the map here rather than taken from the caller, so that a swap-db that ran while + /// the checkpoint was in flight cannot attribute the outcome to the database that was swapped away. + /// + private void UpdateLastSaveData(int dbId, CheckpointResult result) { var databasesMapSnapshot = databases.Map; - var db = databasesMapSnapshot[dbId]; - db.LastSaveTime = DateTimeOffset.UtcNow; - - if (storeTailAddress.HasValue) - { - db.LastSaveStoreTailAddress = storeTailAddress.Value; - } + RecordCheckpointOutcome(databasesMapSnapshot[dbId], result); } /// diff --git a/libs/server/Databases/SingleDatabaseManager.cs b/libs/server/Databases/SingleDatabaseManager.cs index cc1386b630d..6695c52e478 100644 --- a/libs/server/Databases/SingleDatabaseManager.cs +++ b/libs/server/Databases/SingleDatabaseManager.cs @@ -101,6 +101,10 @@ public override async ValueTask RecoverCheckpointAsync(bool replicaRecover = fal if (ex.CandidateTokenCount == 0) { + // Nothing was ever written, so the server comes up empty and no disk state contradicts that. + // Recovery therefore cannot tell that a checkpoint the client was told had succeeded is missing, + // which is why a failed checkpoint must never be reported as a successful save; see + // RecordCheckpointOutcome. Logger?.LogInformation(ex, "No Hybrid Log found for recovery; storeVersion = {storeVersion};", storeVersion); } else @@ -150,32 +154,28 @@ public override void ResumeCheckpoints(int dbId) } /// - public override async Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null) + public override async Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null) { if (dbId != -1 && dbId != 0) throw new ArgumentOutOfRangeException(nameof(dbId), dbId, "SingleDatabaseManager only supports dbId 0."); // Check if checkpoint already in progress if (!TryPauseCheckpoints(defaultDatabase.Id)) - return false; + return CheckpointStatus.AlreadyInProgress; var checkpointTask = TakeCheckpointHelperAsync(defaultDatabase, logger, token); if (background) - return true; + return CheckpointStatus.Success; - await checkpointTask.ConfigureAwait(false); - return true; + return await checkpointTask.ConfigureAwait(false) ? CheckpointStatus.Success : CheckpointStatus.Failed; - async Task TakeCheckpointHelperAsync(GarnetDatabase defaultDatabase, ILogger logger, CancellationToken token) + async Task TakeCheckpointHelperAsync(GarnetDatabase defaultDatabase, ILogger logger, CancellationToken token) { try { - var storeTailAddress = await TakeCheckpointAsync(defaultDatabase, logger: logger, token: token).ConfigureAwait(false); - - if (storeTailAddress.HasValue) - defaultDatabase.LastSaveStoreTailAddress = storeTailAddress.Value; - - defaultDatabase.LastSaveTime = DateTimeOffset.UtcNow; + var result = await TakeCheckpointAsync(defaultDatabase, logger: logger, token: token).ConfigureAwait(false); + RecordCheckpointOutcome(defaultDatabase, result); + return result.IsSuccessful; } finally { @@ -201,13 +201,7 @@ public override async Task TakeOnDemandCheckpointAsync(DateTimeOffset entryTime, // Necessary to take a checkpoint because the latest checkpoint is before entryTime var result = await TakeCheckpointAsync(defaultDatabase, logger: Logger).ConfigureAwait(false); - - var storeTailAddress = result; - - if (storeTailAddress.HasValue) - defaultDatabase.LastSaveStoreTailAddress = storeTailAddress.Value; - - defaultDatabase.LastSaveTime = DateTimeOffset.UtcNow; + RecordCheckpointOutcome(defaultDatabase, result); } finally { @@ -237,11 +231,8 @@ public override async Task TaskCheckpointBasedOnAofSizeLimitAsync(long aofSizeLi logger?.LogInformation("Enforcing AOF size limit currentAofSize: {aofSize} > AofSizeLimit: {aofSizeLimit}", aofSize, aofSizeLimit); - var storeTailAddress = await TakeCheckpointAsync(defaultDatabase, logger: logger, token: token).ConfigureAwait(false); - if (storeTailAddress.HasValue) - defaultDatabase.LastSaveStoreTailAddress = storeTailAddress.Value; - - defaultDatabase.LastSaveTime = DateTimeOffset.UtcNow; + var result = await TakeCheckpointAsync(defaultDatabase, logger: logger, token: token).ConfigureAwait(false); + RecordCheckpointOutcome(defaultDatabase, result); } finally { diff --git a/libs/server/GarnetDatabase.cs b/libs/server/GarnetDatabase.cs index 2040103313c..bd754ec42a1 100644 --- a/libs/server/GarnetDatabase.cs +++ b/libs/server/GarnetDatabase.cs @@ -64,6 +64,12 @@ public class GarnetDatabase : IDisposable /// public DateTimeOffset LastSaveTime; + /// + /// True if the most recent checkpoint attempt for this database completed, false if it failed. Reported to + /// clients as rdb_last_bgsave_status; a background save's failure cannot be reported in its reply. + /// + public bool LastSaveSucceeded; + /// /// What checkpoint recovery found on disk and what it recovered at startup /// @@ -153,6 +159,7 @@ public GarnetDatabase(int id, GarnetDatabase srcDb, bool enableAof, bool copyLas { LastSaveTime = srcDb.LastSaveTime; LastSaveStoreTailAddress = srcDb.LastSaveStoreTailAddress; + LastSaveSucceeded = srcDb.LastSaveSucceeded; } } @@ -161,6 +168,9 @@ public GarnetDatabase() VersionMap = new WatchVersionMap(DefaultVersionMapSize); LastSaveStoreTailAddress = 0; LastSaveTime = DateTimeOffset.FromUnixTimeSeconds(0); + + // Matches Redis, which reports rdb_last_bgsave_status as ok until a save actually fails. + LastSaveSucceeded = true; } /// diff --git a/libs/server/Metrics/Info/GarnetInfoMetrics.cs b/libs/server/Metrics/Info/GarnetInfoMetrics.cs index dc4add6a9cb..b00610eb18a 100644 --- a/libs/server/Metrics/Info/GarnetInfoMetrics.cs +++ b/libs/server/Metrics/Info/GarnetInfoMetrics.cs @@ -374,7 +374,11 @@ private MetricsItem[] GetDatabasePersistenceStats(StoreWrapper storeWrapper, Gar new($"FlushedUntilAddress", !aofEnabled ? "N/A" : db.AppendOnlyFile.Log.FlushedUntilAddress.ToString()), new($"BeginAddress", !aofEnabled ? "N/A" : db.AppendOnlyFile.Log.BeginAddress.ToString()), new($"TailAddress", !aofEnabled ? "N/A" : db.AppendOnlyFile.Log.TailAddress.ToString()), - new($"SafeAofAddress", !aofEnabled ? "N/A" : storeWrapper.safeAofAddress.ToString()) + new($"SafeAofAddress", !aofEnabled ? "N/A" : storeWrapper.safeAofAddress.ToString()), + + // A background save replies before it runs, so this is the only way a client can observe that one + // failed. Named and valued as in Redis. + new($"rdb_last_bgsave_status", db.LastSaveSucceeded ? "ok" : "err") ]; } @@ -524,8 +528,6 @@ private void GetRespInfo(InfoMetricsType section, int dbId, StoreWrapper storeWr GetSectionRespInfo(header, storeRevivInfo[dbId], sbResponse); return; case InfoMetricsType.PERSISTENCE: - if (!storeWrapper.serverOptions.EnableAOF) - return; PopulatePersistenceInfo(storeWrapper); GetSectionRespInfo(header, persistenceInfo[dbId], sbResponse); return; @@ -604,8 +606,6 @@ private MetricsItem[] GetMetricInternal(InfoMetricsType section, int dbId, Store PopulateStoreRevivInfo(storeWrapper); return storeRevivInfo[dbId]; case InfoMetricsType.PERSISTENCE: - if (!storeWrapper.serverOptions.EnableAOF) - return null; PopulatePersistenceInfo(storeWrapper); return persistenceInfo[dbId]; case InfoMetricsType.CLIENTS: diff --git a/libs/server/Resp/AdminCommands.cs b/libs/server/Resp/AdminCommands.cs index 637d4f027fc..688f6b647be 100644 --- a/libs/server/Resp/AdminCommands.cs +++ b/libs/server/Resp/AdminCommands.cs @@ -978,17 +978,23 @@ private bool NetworkSAVE() var checkpointTask = storeWrapper.TakeCheckpointAsync(false, dbId: dbId, logger: logger); // No choice but to block, we're on the network thread - var success = AsyncUtils.BlockingWait(checkpointTask); + var status = AsyncUtils.BlockingWait(checkpointTask); - if (!success) - { - while (!RespWriteUtils.TryWriteError(CmdStrings.RESP_ERR_CHECKPOINT_ALREADY_IN_PROGRESS, ref dcurr, dend)) - SendAndReset(); - } - else + switch (status) { - while (!RespWriteUtils.TryWriteDirect(CmdStrings.RESP_OK, ref dcurr, dend)) - SendAndReset(); + case CheckpointStatus.AlreadyInProgress: + while (!RespWriteUtils.TryWriteError(CmdStrings.RESP_ERR_CHECKPOINT_ALREADY_IN_PROGRESS, ref dcurr, dend)) + SendAndReset(); + break; + case CheckpointStatus.Failed: + // Nothing durable was written, so replying OK here would report data as saved that was not. + while (!RespWriteUtils.TryWriteError(CmdStrings.RESP_ERR_CHECKPOINT_FAILED, ref dcurr, dend)) + SendAndReset(); + break; + default: + while (!RespWriteUtils.TryWriteDirect(CmdStrings.RESP_OK, ref dcurr, dend)) + SendAndReset(); + break; } return true; @@ -1100,16 +1106,18 @@ private bool NetworkBGSAVE() var checkpointTask = storeWrapper.TakeCheckpointAsync(true, dbId: dbId, logger: logger); // No choice but to block, we're on the network thread - var success = AsyncUtils.BlockingWait(checkpointTask); + var status = AsyncUtils.BlockingWait(checkpointTask); - if (success) + // A background checkpoint replies before it runs, so its failure cannot be reported here; it surfaces + // through LASTSAVE not advancing and through rdb_last_bgsave_status in INFO PERSISTENCE. + if (status == CheckpointStatus.AlreadyInProgress) { - while (!RespWriteUtils.TryWriteSimpleString("Background saving started"u8, ref dcurr, dend)) + while (!RespWriteUtils.TryWriteError(CmdStrings.RESP_ERR_CHECKPOINT_ALREADY_IN_PROGRESS, ref dcurr, dend)) SendAndReset(); } else { - while (!RespWriteUtils.TryWriteError(CmdStrings.RESP_ERR_CHECKPOINT_ALREADY_IN_PROGRESS, ref dcurr, dend)) + while (!RespWriteUtils.TryWriteSimpleString("Background saving started"u8, ref dcurr, dend)) SendAndReset(); } diff --git a/libs/server/Resp/CmdStrings.cs b/libs/server/Resp/CmdStrings.cs index 759747bd824..4b375a64a9c 100644 --- a/libs/server/Resp/CmdStrings.cs +++ b/libs/server/Resp/CmdStrings.cs @@ -316,6 +316,7 @@ static partial class CmdStrings public static ReadOnlySpan RESP_ERR_ZSET_MEMBER => "ERR could not decode requested zset member"u8; public static ReadOnlySpan RESP_ERR_EXPDELSCAN_INVALID => "ERR Cannot execute EXPDELSCAN with background expired key deletion scan enabled"u8; public static ReadOnlySpan RESP_ERR_CHECKPOINT_ALREADY_IN_PROGRESS => "ERR checkpoint already in progress"u8; + public static ReadOnlySpan RESP_ERR_CHECKPOINT_FAILED => "ERR checkpoint failed, check server logs"u8; /// /// Response string templates diff --git a/libs/server/Storage/Functions/GarnetRecordTriggers.cs b/libs/server/Storage/Functions/GarnetRecordTriggers.cs index 2fd9c54ce44..49a95d3e750 100644 --- a/libs/server/Storage/Functions/GarnetRecordTriggers.cs +++ b/libs/server/Storage/Functions/GarnetRecordTriggers.cs @@ -186,6 +186,16 @@ public readonly void OnCheckpoint(CheckpointTrigger trigger, Guid checkpointToke case CheckpointTrigger.CheckpointCompleted: vectorManager?.CheckpointCompleted(); break; + case CheckpointTrigger.CheckpointFailed: + // Release the barrier set at VersionShift, which FlushBegin would have cleared - range index + // operations spin-wait on it, so an aborted checkpoint would block them indefinitely. Idempotent, + // so it is safe when the abort happened after FlushBegin or before the barrier was ever set. + // + // Deliberately does not call vectorManager.CheckpointCompleted(): that reclaims deletions, which + // is only safe once a checkpoint has made them recoverable. A failed checkpoint has not, so they + // must stay queued for the next successful one. + rangeIndexManager?.ClearCheckpointBarrier(); + break; } } diff --git a/libs/server/StoreWrapper.cs b/libs/server/StoreWrapper.cs index 11a1e71fe9b..f85d63a0859 100644 --- a/libs/server/StoreWrapper.cs +++ b/libs/server/StoreWrapper.cs @@ -411,8 +411,12 @@ internal async ValueTask RecoverAsync() /// ID of database to checkpoint, or -1 (default) to checkpoint all active databases /// Cancellation token /// Logger - /// False if another checkpointing process is already in progress - public Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null) + /// + /// if another checkpointing process is already in progress, + /// if a foreground checkpoint did not complete, otherwise + /// + /// + public Task TakeCheckpointAsync(bool background, int dbId = -1, CancellationToken token = default, ILogger logger = null) { if (dbId > 0 && !CheckMultiDatabaseCompatibility()) throw new GarnetException($"Unable to call {nameof(databaseManager.TakeCheckpointAsync)} with DB ID: {dbId}"); diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/HybridLogCheckpointSMTask.cs b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/HybridLogCheckpointSMTask.cs index 4fe9c5f7f7d..1b8820a9fc2 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/HybridLogCheckpointSMTask.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/HybridLogCheckpointSMTask.cs @@ -115,5 +115,28 @@ public virtual void GlobalAfterEnteringState(SystemState next, StateMachineDrive break; } } + + /// + public virtual void OnAbort(StateMachineDriver stateMachineDriver, Exception exception) + { + // Mirrors the Phase.REST handling above, which an aborted state machine never reaches. + + // Lets the application release the barrier it set at VersionShift. Deliberately a distinct trigger from + // CheckpointCompleted: work that is only safe once a checkpoint has made its data recoverable - such as + // reclaiming deletions - must stay pending for the next successful checkpoint. + store.storeFunctions.OnCheckpoint(CheckpointTrigger.CheckpointFailed, guid); + + // Releases any snapshot devices and flush buffers already created, and clears the checkpoint so the next + // one can run. Matches the cleanup CompleteCheckpointAsync performs when it observes a failed checkpoint. + store._hybridLogCheckpoint.Dispose(); + + // Waiters such as ClientSession.WaitForCommitAsync park on store.CheckpointTask, which REST would have + // completed; leaving it pending hangs them forever. Publish the next checkpoint's source before faulting + // the old one, so a continuation that immediately re-reads store.checkpointTcs picks up the source for + // the next checkpoint rather than the one it just watched fail. + var previousTcs = store.checkpointTcs; + store.checkpointTcs = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + _ = previousTcs.TrySetException(exception ?? new TsavoriteException("Checkpoint state machine aborted")); + } } } \ No newline at end of file diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IStateMachineTask.cs b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IStateMachineTask.cs index 68a778374c6..790281990dd 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IStateMachineTask.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IStateMachineTask.cs @@ -1,6 +1,8 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT license. +using System; + namespace Tsavorite.core { /// @@ -21,5 +23,20 @@ public interface IStateMachineTask /// /// public void GlobalAfterEnteringState(SystemState nextState, StateMachineDriver stateMachineDriver); + + /// + /// Called when the state machine aborts before reaching , so that the task can + /// release whatever it set up in the phases it did enter. + /// + /// + /// The exception that aborted the state machine. + /// + /// The REST phase is where a task normally releases its resources and completes whatever the rest of the + /// system is waiting on, and an aborted state machine never enters it. A task that leaves state behind here + /// makes every subsequent run of the same state machine fail, and a task that leaves a waiter uncompleted + /// hangs it, so the failure of one operation becomes permanent. Defaults to doing nothing for tasks that hold + /// no such state. + /// + public void OnAbort(StateMachineDriver stateMachineDriver, Exception exception) { } } } \ No newline at end of file diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IndexCheckpointSMTask.cs b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IndexCheckpointSMTask.cs index ef5f1bc877f..f23efeebb8a 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IndexCheckpointSMTask.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/IndexCheckpointSMTask.cs @@ -62,5 +62,14 @@ public void GlobalBeforeEnteringState(SystemState next, StateMachineDriver state public void GlobalAfterEnteringState(SystemState next, StateMachineDriver stateMachineDriver) { } + + /// + public void OnAbort(StateMachineDriver stateMachineDriver, Exception exception) + { + // Mirrors the Phase.REST handling above, which an aborted state machine never reaches. Leaving + // _indexCheckpoint set would make the PREPARE phase of every later checkpoint fail its IsDefault check, + // so one failed checkpoint would stop the store from ever checkpointing again. + store._indexCheckpoint.Reset(); + } } } \ No newline at end of file diff --git a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineBase.cs b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineBase.cs index f031e40f8ac..f09b3a0b50e 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineBase.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineBase.cs @@ -1,6 +1,8 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT license. +using System; + namespace Tsavorite.core { /// @@ -37,5 +39,12 @@ public void GlobalAfterEnteringState(SystemState next, StateMachineDriver stateM foreach (var task in tasks) task.GlobalAfterEnteringState(next, stateMachineDriver); } + + /// + public void OnAbort(StateMachineDriver stateMachineDriver, Exception exception) + { + foreach (var task in tasks) + task.OnAbort(stateMachineDriver, exception); + } } } \ No newline at end of file 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 5c9fba08bae..4b42be1d90c 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineDriver.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/Checkpointing/StateMachineDriver.cs @@ -377,6 +377,7 @@ async Task RunStateMachine(CancellationToken token = default) catch (Exception e) { FastForwardStateMachineToRest(); + ReleaseAbortedStateMachineResources(e); logger?.LogError(e, "Exception in state machine"); ex = e; throw; @@ -401,6 +402,25 @@ async Task RunStateMachine(CancellationToken token = default) } } + /// + /// Lets the state machine release resources it acquired in the phases it entered before aborting, and + /// complete anything the rest of the system is waiting on, since the REST phase that normally does both is + /// never reached. + /// + /// The exception that aborted the state machine. + void ReleaseAbortedStateMachineResources(Exception exception) + { + try + { + stateMachine.OnAbort(this, exception); + } + catch (Exception e) + { + // Must not replace the exception that aborted the state machine, which is the actionable one. + logger?.LogError(e, "Exception while releasing the resources of an aborted state machine"); + } + } + void FastForwardStateMachineToRest() { // Move system state to the next REST phase diff --git a/libs/storage/Tsavorite/cs/src/core/Index/StoreFunctions/CheckpointTrigger.cs b/libs/storage/Tsavorite/cs/src/core/Index/StoreFunctions/CheckpointTrigger.cs index 1f9099ef86b..9f9f1e5bcfb 100644 --- a/libs/storage/Tsavorite/cs/src/core/Index/StoreFunctions/CheckpointTrigger.cs +++ b/libs/storage/Tsavorite/cs/src/core/Index/StoreFunctions/CheckpointTrigger.cs @@ -26,6 +26,16 @@ public enum CheckpointTrigger /// REST phase, after the checkpoint is fully persisted. The application /// should clean up outdated external checkpoint artifacts. /// - CheckpointCompleted + CheckpointCompleted, + + /// + /// The checkpoint state machine aborted before reaching , so nothing was + /// persisted. The application should release the barrier set during without doing + /// any of the work that is only safe once a checkpoint has made its data recoverable. + /// + /// + /// May be raised even when was never reached, so handling must be idempotent. + /// + CheckpointFailed } } \ No newline at end of file diff --git a/test/standalone/Garnet.test.scripting/MultiDatabaseTests.cs b/test/standalone/Garnet.test.scripting/MultiDatabaseTests.cs index dde3fe48475..689000e098a 100644 --- a/test/standalone/Garnet.test.scripting/MultiDatabaseTests.cs +++ b/test/standalone/Garnet.test.scripting/MultiDatabaseTests.cs @@ -1461,6 +1461,66 @@ public void MultiDatabaseGeneralSaveBlocksGeneralSaveTest() ClassicAssert.Greater(lastsave, lastsaveBaseline, "LASTSAVE did not advance within timeout"); } + [Test] + public void MultiDatabaseForegroundSaveSkippingBusyDatabaseDoesNotReportSuccess() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db0 = redis.GetDatabase(0); + var db1 = redis.GetDatabase(1); + + db0.StringSet("k", "v"); + db1.StringSet("k", "v"); + + var storeWrapper = server.Provider.StoreWrapper; + + // Hold DB 1's per-DB checkpoint lock, which is exactly the state a per-DB BGSAVE on DB 1 + // leaves the server in while it is running. TakeCheckpointAsync skips a database it cannot + // pause-lock, so a general SAVE here checkpoints DB 0 only. Taking the lock directly makes + // the skip deterministic rather than depending on a real checkpoint still being in flight. + ClassicAssert.IsTrue(storeWrapper.TryPauseCheckpoints(1), "Could not pause checkpoints on DB 1"); + + try + { + // A foreground SAVE consumes the aggregate outcome, so it must not answer +OK after + // skipping DB 1 - that would tell the client every database is on disk when one was + // never attempted. + var ex = Assert.Throws(() => db0.Execute("SAVE")); + ClassicAssert.AreEqual( + Encoding.ASCII.GetString(CmdStrings.RESP_ERR_CHECKPOINT_ALREADY_IN_PROGRESS), + ex.Message); + + // BGSAVE replies before the aggregate is computed, so its contract of "skip the busy + // databases and report started" is unchanged. + var res = db0.Execute("BGSAVE"); + ClassicAssert.AreEqual("Background saving started", res.ToString()); + } + finally + { + storeWrapper.ResumeCheckpoints(1); + } + + // Once no database is busy, a general SAVE succeeds again - the error above was the skip, + // not a permanent refusal. The BGSAVE above may still be finishing, so poll. + var deadline = DateTime.UtcNow.AddSeconds(30); + string saveResult; + do + { + try + { + saveResult = db0.Execute("SAVE").ToString(); + break; + } + catch (RedisServerException) + { + saveResult = null; + Thread.Sleep(10); + } + } + while (DateTime.UtcNow < deadline); + + ClassicAssert.AreEqual("OK", saveResult, "SAVE did not succeed within timeout once all DBs were free"); + } + [Test] [TestCase(false)] [TestCase(true)] diff --git a/test/standalone/Garnet.test/CheckpointFailureTests.cs b/test/standalone/Garnet.test/CheckpointFailureTests.cs new file mode 100644 index 00000000000..33f995aa28c --- /dev/null +++ b/test/standalone/Garnet.test/CheckpointFailureTests.cs @@ -0,0 +1,443 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT license. + +using System; +using System.IO; +using System.Linq; +using System.Threading; +using Garnet.server; +using NUnit.Framework; +using NUnit.Framework.Legacy; +using StackExchange.Redis; +using Tsavorite.core; + +namespace Garnet.test +{ + /// + /// A checkpoint that fails must never be reported to clients as a successful save. LASTSAVE must not advance, + /// SAVE must return an error, and INFO PERSISTENCE must report rdb_last_bgsave_status:err. Otherwise a + /// client following the documented "BGSAVE then poll LASTSAVE" pattern treats data as durable that was never + /// written, and the loss only surfaces as an empty store after the next restart. + /// + /// + /// The failure is injected through rather than by making the + /// filesystem itself fail, so the tests do not depend on path length, file permissions or platform. + /// + [TestFixture] + public class CheckpointFailureTests : TestBase + { + const string StatusPrefix = "rdb_last_bgsave_status:"; + const string TestKey = "CheckpointFailureTestKey"; + const string TestValue = "CheckpointFailureTestValue"; + + static readonly long EpochTicks = DateTimeOffset.FromUnixTimeSeconds(0).Ticks; + static readonly TimeSpan CheckpointTimeout = TimeSpan.FromSeconds(30); + + GarnetServer server; + GarnetServerOptions options; + FailingCheckpointDeviceFactoryCreator deviceFactoryCreator; + + [SetUp] + public void Setup() + { + TestUtils.DeleteDirectory(TestUtils.MethodTestDir, wait: true); + server = CreateServer(tryRecover: false); + server.Start(); + } + + [TearDown] + public void TearDown() + { + server?.Dispose(); + TestUtils.OnTearDown(); + } + + GarnetServer CreateServer(bool tryRecover) => CreateServer(tryRecover, enableAof: false); + + GarnetServer CreateServer(bool tryRecover, bool enableAof) + { + // Built through GetGarnetServerOptions rather than CreateGarnetServer because only the options object + // exposes DeviceFactoryCreator, which is how the checkpoint failure is injected. + options = TestUtils.GetGarnetServerOptions( + checkpointDir: TestUtils.MethodTestDir, + logDir: TestUtils.MethodTestDir, + endpoint: TestUtils.EndPoint, + enableCluster: false, + enableAOF: enableAof, + tryRecover: tryRecover); + + deviceFactoryCreator = new FailingCheckpointDeviceFactoryCreator(options.StoreCheckpointBaseDirectory); + options.DeviceFactoryCreator = deviceFactoryCreator; + + return new GarnetServer(options); + } + + [Test] + public void BackgroundSaveFailureDoesNotAdvanceLastSave() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db = redis.GetDatabase(0); + var redisServer = redis.GetServer(TestUtils.EndPoint); + + db.StringSet(TestKey, TestValue); + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE should start at the epoch"); + + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.All; + redisServer.Save(SaveType.BackgroundSave); + + // BGSAVE replies before the checkpoint runs, so wait for the outcome to be recorded rather than sleeping + // for an arbitrary interval. + WaitForLastSaveStatus(db, "err"); + + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks, + "LASTSAVE advanced for a checkpoint that failed, so a client polling it would treat unwritten data as durable"); + ClassicAssert.AreEqual(0, CountCheckpointFiles(dbId: 0), "A failed checkpoint should not leave checkpoint files behind"); + } + + [Test] + public void SaveFailureReturnsErrorToClient() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db = redis.GetDatabase(0); + var redisServer = redis.GetServer(TestUtils.EndPoint); + + db.StringSet(TestKey, TestValue); + + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.All; + var ex = Assert.Throws(() => db.Execute("SAVE")); + ClassicAssert.IsTrue(ex.Message.StartsWith("ERR checkpoint failed", StringComparison.Ordinal), + $"SAVE should report the failure to the client, but replied: {ex.Message}"); + + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE advanced for a failed SAVE"); + ClassicAssert.AreEqual("err", GetLastSaveStatus(db)); + ClassicAssert.AreEqual(0, CountCheckpointFiles(dbId: 0)); + } + + [Test] + public void SaveSucceedsAfterFailureAndDataSurvivesRestart() + { + using (var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true))) + { + var db = redis.GetDatabase(0); + var redisServer = redis.GetServer(TestUtils.EndPoint); + + db.StringSet(TestKey, TestValue); + + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.All; + _ = Assert.Throws(() => db.Execute("SAVE")); + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks); + + // A failed checkpoint must leave the server able to take the next one. + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.None; + ClassicAssert.AreEqual("OK", db.Execute("SAVE").ToString()); + + ClassicAssert.AreNotEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE should advance for a successful save"); + ClassicAssert.AreEqual("ok", GetLastSaveStatus(db)); + ClassicAssert.Greater(CountCheckpointFiles(dbId: 0), 0); + } + + server.Dispose(false); + server = CreateServer(tryRecover: true); + server.Start(); + + using (var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true))) + { + var db = redis.GetDatabase(0); + ClassicAssert.AreEqual(TestValue, db.StringGet(TestKey).ToString(), + "The successful save should have been recoverable"); + } + } + + [Test] + public void BackgroundSaveFailureDoesNotAdvanceLastSaveForAnyDatabase() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db0 = redis.GetDatabase(0); + + // Touching a non-zero database promotes the server to the multi-database manager, which records the + // last save time through its own code path. + var db1 = redis.GetDatabase(1); + + db0.StringSet(TestKey, TestValue); + db1.StringSet(TestKey, TestValue); + + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.All; + ClassicAssert.AreEqual("Background saving started", db0.Execute("BGSAVE").ToString()); + + WaitForLastSaveStatus(db0, "err"); + WaitForLastSaveStatus(db1, "err"); + + ClassicAssert.AreEqual(0, (long)db0.Execute("LASTSAVE"), "LASTSAVE advanced for DB 0 after a failed checkpoint"); + ClassicAssert.AreEqual(0, (long)db1.Execute("LASTSAVE"), "LASTSAVE advanced for DB 1 after a failed checkpoint"); + ClassicAssert.AreEqual(0, CountCheckpointFiles(dbId: 0)); + ClassicAssert.AreEqual(0, CountCheckpointFiles(dbId: 1)); + } + + [Test] + public void SnapshotDeviceFailureDoesNotAdvanceLastSaveAndLeavesStoreCheckpointable() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db = redis.GetDatabase(0); + var redisServer = redis.GetServer(TestUtils.EndPoint); + + db.StringSet(TestKey, TestValue); + + // The snapshot devices are requested at WAIT_FLUSH, once both the index checkpoint and the hybrid-log + // checkpoint have been initialized, so this aborts with state from both to release. The all-devices mode + // instead fails in the index checkpoint, before the hybrid-log checkpoint exists. + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.SnapshotLogDevicesOnly; + + var ex = Assert.Throws(() => db.Execute("SAVE")); + ClassicAssert.IsTrue(ex.Message.StartsWith("ERR checkpoint failed", StringComparison.Ordinal), + $"SAVE should report the failure to the client, but replied: {ex.Message}"); + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE advanced for a failed snapshot"); + ClassicAssert.AreEqual("err", GetLastSaveStatus(db)); + + // Assert the absence of leaks directly, not just that the next checkpoint works: a leaked device would + // still let the next checkpoint pass while leaking a file handle on every failure. + AssertNoCheckpointStateLeaked(); + + // Only possible if the abort released the hybrid-log checkpoint; otherwise the next checkpoint fails too. + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.None; + ClassicAssert.AreEqual("OK", db.Execute("SAVE").ToString()); + + ClassicAssert.AreNotEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE should advance for a successful save"); + ClassicAssert.AreEqual("ok", GetLastSaveStatus(db)); + ClassicAssert.Greater(CountCheckpointFiles(dbId: 0), 0); + } + + [Test] + public void CheckpointBlockedByAnotherStateMachineFailsWithoutTruncatingAof() + { + // The AOF is what makes the truncation this branch prevents observable. + server.Dispose(); + server = CreateServer(tryRecover: false, enableAof: true); + server.Start(); + + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db = redis.GetDatabase(0); + var redisServer = redis.GetServer(TestUtils.EndPoint); + + for (var i = 0; i < 64; i++) + db.StringSet($"{TestKey}{i}", TestValue); + + var beginAddressBeforeSave = GetPersistenceField(db, "BeginAddress:"); + + // Occupy the database's state machine driver, so the checkpoint's RunAsync returns false instead of + // throwing. No checkpoint is taken in that case, so the AOF must not be truncated and the save must not + // be reported as successful. + var driver = server.Provider.StoreWrapper.DefaultDatabase.StateMachineDriver; + using var release = new ManualResetEventSlim(false); + var blocking = new BlockingStateMachine(release); + ClassicAssert.IsTrue(driver.Register(blocking), "Could not occupy the state machine driver"); + + try + { + ClassicAssert.IsTrue(blocking.Entered.Wait(CheckpointTimeout), "Blocking state machine never started"); + + var ex = Assert.Throws(() => db.Execute("SAVE")); + ClassicAssert.IsTrue(ex.Message.StartsWith("ERR checkpoint failed", StringComparison.Ordinal), + $"SAVE should report the failure to the client, but replied: {ex.Message}"); + + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks, + "LASTSAVE advanced for a checkpoint that never ran"); + ClassicAssert.AreEqual("err", GetLastSaveStatus(db)); + ClassicAssert.AreEqual(beginAddressBeforeSave, GetPersistenceField(db, "BeginAddress:"), + "The AOF was truncated for a checkpoint that never ran"); + ClassicAssert.AreEqual(0, CountCheckpointFiles(dbId: 0)); + } + finally + { + release.Set(); + } + + ClassicAssert.IsTrue(blocking.Completed.Wait(CheckpointTimeout), "Blocking state machine never finished"); + } + + [Test] + public void FailedCheckpointFaultsCheckpointWaitersAndPublishesANewTask() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db = redis.GetDatabase(0); + + db.StringSet(TestKey, TestValue); + + var store = server.Provider.StoreWrapper.store; + var taskBeforeSave = store.CheckpointTask; + ClassicAssert.IsFalse(taskBeforeSave.IsCompleted, "The checkpoint task should start out pending"); + + deviceFactoryCreator.FailureMode = CheckpointDeviceFailure.SnapshotLogDevicesOnly; + _ = Assert.Throws(() => db.Execute("SAVE")); + + // Waiters such as ClientSession.WaitForCommitAsync park on this task, which only the REST phase would + // have completed. Leaving it pending after an abort hangs them forever. + ClassicAssert.IsTrue(taskBeforeSave.IsCompleted, "The checkpoint task was left pending after a failed checkpoint"); + ClassicAssert.IsTrue(taskBeforeSave.IsFaulted, "The checkpoint task should be faulted for a failed checkpoint"); + _ = taskBeforeSave.Exception; + + // The replacement must already be published, so a waiter that re-reads the task after being woken picks + // up the source for the next checkpoint rather than the one it just watched fail. + var taskAfterSave = store.CheckpointTask; + ClassicAssert.AreNotSame(taskBeforeSave, taskAfterSave, "A new checkpoint task should have been published"); + ClassicAssert.IsFalse(taskAfterSave.IsCompleted, "The replacement checkpoint task should be pending"); + } + + [Test] + public void CheckpointAbortedAtVersionShiftLeaksNothingAndLeavesStoreCheckpointable() + { + using var redis = ConnectionMultiplexer.Connect(TestUtils.GetConfig(allowAdmin: true)); + var db = redis.GetDatabase(0); + var redisServer = redis.GetServer(TestUtils.EndPoint); + + db.StringSet(TestKey, TestValue); + + // A device failure can only abort where a device is created, which is either before the version shift + // (index device, in PREPARE) or after FlushBegin has already cleared the range-index barrier (snapshot + // devices, in WAIT_FLUSH). Aborting at IN_PROGRESS is therefore the only way to reach the window where + // the barrier is set but not yet cleared - the window CheckpointTrigger.CheckpointFailed exists for. + var driver = server.Provider.StoreWrapper.DefaultDatabase.StateMachineDriver; + var failOnce = new ThrowOnceAtPhase(Phase.IN_PROGRESS); + driver.UnsafeRegisterCallback(failOnce); + + var ex = Assert.Throws(() => db.Execute("SAVE")); + ClassicAssert.IsTrue(ex.Message.StartsWith("ERR checkpoint failed", StringComparison.Ordinal), + $"SAVE should report the failure to the client, but replied: {ex.Message}"); + ClassicAssert.IsTrue(failOnce.Fired, "The abort was never injected"); + + ClassicAssert.AreEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE advanced for an aborted checkpoint"); + ClassicAssert.AreEqual("err", GetLastSaveStatus(db)); + AssertNoCheckpointStateLeaked(); + + // The callback only throws once, so the next checkpoint exercises recovery from the abort. + ClassicAssert.AreEqual("OK", db.Execute("SAVE").ToString()); + ClassicAssert.AreNotEqual(EpochTicks, redisServer.LastSave().Ticks, "LASTSAVE should advance for a successful save"); + ClassicAssert.AreEqual("ok", GetLastSaveStatus(db)); + ClassicAssert.Greater(CountCheckpointFiles(dbId: 0), 0); + } + + /// + /// Asserts that an aborted checkpoint released every device and buffer it had taken. Checking the fields + /// directly catches a leak that a "the next checkpoint still works" assertion would miss. + /// + void AssertNoCheckpointStateLeaked() + { + var store = server.Provider.StoreWrapper.store; + + Assert.Multiple(() => + { + Assert.That(store._hybridLogCheckpoint.snapshotFileDevice, Is.Null, "Snapshot log device was leaked"); + Assert.That(store._hybridLogCheckpoint.snapshotFileObjectLogDevice, Is.Null, "Snapshot object log device was leaked"); + Assert.That(store._hybridLogCheckpoint.objectLogFlushBuffers, Is.Null, "Object log flush buffers were leaked"); + Assert.That(store._indexCheckpoint.IsDefault, Is.True, "Index checkpoint was not reset"); + Assert.That(store._indexCheckpoint.main_ht_device, Is.Null, "Index hash table device was leaked"); + }); + } + + /// + /// Number of files written under the given database's log checkpoint directory. + /// + int CountCheckpointFiles(int dbId) + { + var namingScheme = new DefaultCheckpointNamingScheme(options.GetStoreCheckpointDirectory(dbId)); + var checkpointDir = Path.Combine(namingScheme.BaseName, namingScheme.LogCheckpointBasePath); + + return Directory.Exists(checkpointDir) + ? Directory.GetFiles(checkpointDir, "*", SearchOption.AllDirectories).Length + : 0; + } + + static string GetLastSaveStatus(IDatabase db) => GetPersistenceField(db, StatusPrefix); + + /// + /// Value of a single field:value line in the INFO PERSISTENCE section. + /// + static string GetPersistenceField(IDatabase db, string prefix) + { + var info = db.Execute("INFO", "PERSISTENCE").ToString(); + var line = info.Split("\r\n").FirstOrDefault(x => x.StartsWith(prefix, StringComparison.Ordinal)); + ClassicAssert.IsNotNull(line, $"INFO PERSISTENCE did not report {prefix}"); + + return line[prefix.Length..]; + } + + static void WaitForLastSaveStatus(IDatabase db, string expected) + { + var deadline = DateTime.UtcNow + CheckpointTimeout; + while (GetLastSaveStatus(db) != expected && DateTime.UtcNow < deadline) + Thread.Sleep(10); + + ClassicAssert.AreEqual(expected, GetLastSaveStatus(db), + $"{StatusPrefix} did not become '{expected}' within {CheckpointTimeout.TotalSeconds} seconds"); + } + + /// + /// Occupies a until released, so that a checkpoint issued meanwhile finds the + /// driver already running a state machine and does not run at all. + /// + private sealed class BlockingStateMachine : IStateMachine + { + readonly ManualResetEventSlim release; + + /// Set once the state machine is occupying the driver. + internal readonly ManualResetEventSlim Entered = new(false); + + /// Set once the state machine has been released and is returning to REST. + internal readonly ManualResetEventSlim Completed = new(false); + + internal BlockingStateMachine(ManualResetEventSlim release) => this.release = release; + + /// + public SystemState NextState(SystemState start) + { + var next = SystemState.Copy(ref start); + next.Phase = start.Phase == Phase.REST ? Phase.PREPARE_GROW : Phase.REST; + return next; + } + + /// + public void GlobalBeforeEnteringState(SystemState nextState, StateMachineDriver stateMachineDriver) + { + // Blocks before the driver takes epoch protection, so parking here holds no epoch. + if (nextState.Phase == Phase.PREPARE_GROW) + { + Entered.Set(); + release.Wait(); + } + else if (nextState.Phase == Phase.REST) + { + Completed.Set(); + } + } + + /// + public void GlobalAfterEnteringState(SystemState nextState, StateMachineDriver stateMachineDriver) { } + } + + /// + /// Aborts the state machine once, on first entry to a chosen phase. Device-agnostic, so it reaches abort + /// points that no device failure can - notably , which is after the version + /// shift has published its lifecycle notification but before any snapshot device exists. + /// + /// + /// The driver has no way to unregister a callback, so this stays attached for the store's lifetime and must + /// throw only once; later checkpoints have to be able to succeed. + /// + private sealed class ThrowOnceAtPhase : IStateMachineCallback + { + readonly Phase phase; + int fired; + + internal ThrowOnceAtPhase(Phase phase) => this.phase = phase; + + /// True once the abort has been injected. + internal bool Fired => Volatile.Read(ref fired) != 0; + + /// + public void BeforeEnteringState(SystemState nextState) + { + if (nextState.Phase == phase && Interlocked.CompareExchange(ref fired, 1, 0) == 0) + throw new IOException($"Simulated checkpoint abort at {phase}"); + } + } + } +} \ No newline at end of file diff --git a/test/standalone/Garnet.test/FailingCheckpointDeviceFactoryCreator.cs b/test/standalone/Garnet.test/FailingCheckpointDeviceFactoryCreator.cs new file mode 100644 index 00000000000..a31d0a60063 --- /dev/null +++ b/test/standalone/Garnet.test/FailingCheckpointDeviceFactoryCreator.cs @@ -0,0 +1,129 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT license. + +using System; +using System.Collections.Generic; +using System.IO; +using Tsavorite.core; + +namespace Garnet.test +{ + /// + /// Which store checkpoint devices a should fail to create. + /// + internal enum CheckpointDeviceFailure + { + /// Every checkpoint device is created normally. + None, + + /// + /// Every checkpoint device creation throws. A full checkpoint therefore fails in the index checkpoint task, + /// before the hybrid-log checkpoint has been initialized. + /// + All, + + /// + /// Only the hybrid-log snapshot devices throw. Those are requested at Phase.WAIT_FLUSH, after the index + /// checkpoint and the hybrid-log checkpoint have both been initialized, so the abort has state from both to + /// release. This is the shape of the real failure that motivated this work, where the device layer rejected an + /// over-long snapshot.obj.dat path. + /// + SnapshotLogDevicesOnly, + } + + /// + /// Device factory creator that makes store checkpoint device creation fail on demand while leaving every other + /// device - hybrid log, storage tier and AOF - working normally. + /// + /// + /// Injected by assigning it to GarnetServerOptions.DeviceFactoryCreator, which GarnetServer.CreateStore + /// hands to the store's checkpoint manager. DeviceLogCommitCheckpointManager does not guard its + /// deviceFactory.Get calls, so a throw from the returned factory propagates out of the checkpoint state + /// machine the same way a real device-creation failure does, with no dependence on path length, file permissions + /// or platform. Only devices below the store checkpoint base directory are affected, so the server keeps serving + /// while its checkpoints fail. + /// + internal sealed class FailingCheckpointDeviceFactoryCreator : INamedDeviceFactoryCreator + { + readonly INamedDeviceFactoryCreator inner = new LocalStorageNamedDeviceFactoryCreator(); + readonly string storeCheckpointBaseDirectory; + + /// + /// Which checkpoint devices currently fail to be created. Read on every device request, so a failure can be + /// armed and disarmed against a running server. + /// + internal volatile CheckpointDeviceFailure FailureMode = CheckpointDeviceFailure.None; + + // Taken from the naming scheme rather than hard-coded, so renaming a checkpoint file cannot silently turn + // SnapshotLogDevicesOnly into a mode that fails nothing. + static readonly string LogSnapshotFileName = new DefaultCheckpointNamingScheme(string.Empty).LogSnapshot(Guid.Empty).fileName; + static readonly string ObjectLogSnapshotFileName = new DefaultCheckpointNamingScheme(string.Empty).ObjectLogSnapshot(Guid.Empty).fileName; + + /// + /// Creates a creator that fails checkpoint devices under the given directory. + /// + /// Store checkpoint base directory, as given by + /// GarnetServerOptions.StoreCheckpointBaseDirectory. Every database's checkpoint directory lives + /// beneath it, so a prefix match covers multi-database servers too. + internal FailingCheckpointDeviceFactoryCreator(string storeCheckpointBaseDirectory) + { + this.storeCheckpointBaseDirectory = storeCheckpointBaseDirectory; + } + + /// + public INamedDeviceFactory Create(string baseName) + { + var factory = inner.Create(baseName); + return baseName is not null && baseName.StartsWith(storeCheckpointBaseDirectory, StringComparison.Ordinal) + ? new FailingNamedDeviceFactory(factory, this) + : factory; + } + + /// + /// True if a device request should throw under the current failure mode. + /// + private bool ShouldFail(FileDescriptor fileInfo) => FailureMode switch + { + CheckpointDeviceFailure.All => true, + CheckpointDeviceFailure.SnapshotLogDevicesOnly => + string.Equals(fileInfo.fileName, LogSnapshotFileName, StringComparison.Ordinal) || + string.Equals(fileInfo.fileName, ObjectLogSnapshotFileName, StringComparison.Ordinal), + _ => false, + }; + + /// + /// Delegating device factory that throws instead of creating a device while its owner is armed. + /// + private sealed class FailingNamedDeviceFactory : INamedDeviceFactory + { + readonly INamedDeviceFactory underlying; + readonly FailingCheckpointDeviceFactoryCreator owner; + + internal FailingNamedDeviceFactory(INamedDeviceFactory underlying, FailingCheckpointDeviceFactoryCreator owner) + { + this.underlying = underlying; + this.owner = owner; + } + + /// + public IDevice Get(FileDescriptor fileInfo) + { + if (owner.ShouldFail(fileInfo)) + { + // Thrown before the underlying device is created, so a failed checkpoint leaves no open handle + // behind for the test's directory cleanup to trip over. + throw new IOException( + $"Simulated checkpoint device failure for {Path.Combine(fileInfo.directoryName ?? string.Empty, fileInfo.fileName ?? string.Empty)}"); + } + + return underlying.Get(fileInfo); + } + + /// + public void Delete(FileDescriptor fileInfo) => underlying.Delete(fileInfo); + + /// + public IEnumerable ListContents(string path) => underlying.ListContents(path); + } + } +} \ No newline at end of file diff --git a/test/standalone/Garnet.test/RespAdminCommandsTests.cs b/test/standalone/Garnet.test/RespAdminCommandsTests.cs index 28acc8ee546..3cd421ba11e 100644 --- a/test/standalone/Garnet.test/RespAdminCommandsTests.cs +++ b/test/standalone/Garnet.test/RespAdminCommandsTests.cs @@ -187,7 +187,7 @@ public async Task OverflowKeyCheckpointTest() }); await writerTask.ConfigureAwait(false); - Assert.That(await checkpointTask.ConfigureAwait(false), Is.True); + Assert.That(await checkpointTask.ConfigureAwait(false), Is.EqualTo(CheckpointStatus.Success)); // Cleanup previously failed while remapping an overflow key, which discarded this full-checkpoint tail. var firstCheckpointTail = server.Provider.StoreWrapper.DefaultDatabase.LastSaveStoreTailAddress; @@ -196,7 +196,7 @@ public async Task OverflowKeyCheckpointTest() // With no further writes, log growth is below FullCheckpointLogInterval, so the next checkpoint is // incremental and leaves the last full-checkpoint tail unchanged. - Assert.That(await server.Provider.StoreWrapper.TakeCheckpointAsync(background: false).ConfigureAwait(false), Is.True); + Assert.That(await server.Provider.StoreWrapper.TakeCheckpointAsync(background: false).ConfigureAwait(false), Is.EqualTo(CheckpointStatus.Success)); Assert.That(server.Provider.StoreWrapper.DefaultDatabase.LastSaveStoreTailAddress, Is.EqualTo(firstCheckpointTail), "An incremental checkpoint must not replace the preceding full-checkpoint tail"); } diff --git a/website/docs/commands/checkpoint.md b/website/docs/commands/checkpoint.md index c4b66c02e23..c1d30fa93c7 100644 --- a/website/docs/commands/checkpoint.md +++ b/website/docs/commands/checkpoint.md @@ -14,12 +14,17 @@ BGSAVE [SCHEDULE] [DBID] Save all databases inside the Garnet instance in the background. If a DB ID is specified, save save only that specific database. +The reply is sent before the checkpoint runs, so it does not report whether the checkpoint succeeded. A background +save that fails leaves `LASTSAVE` unchanged and sets `rdb_last_bgsave_status` to `err` in the `PERSISTENCE` section of +[INFO](server.md#info), with details in the server log. + #### Resp Reply One of the following: * Simple string reply: Background saving started. * Simple string reply: Background saving scheduled. +* Error reply: checkpoint already in progress. --- @@ -35,7 +40,10 @@ The SAVE commands performs a synchronous save of the dataset producing a point i #### Resp Reply -Simple string reply: OK. +One of the following: + +* Simple string reply: OK. +* Error reply, if the checkpoint did not complete. Nothing was written, so `LASTSAVE` is left unchanged. --- ### LASTSAVE @@ -47,6 +55,9 @@ LASTSAVE [DBID] Return the UNIX TIME of the last DB save executed with success for the current database or, if a DB ID is specified, the last DB save executed with success for the specified database. +Only a checkpoint that completed advances this timestamp, so it can be polled after `BGSAVE` to confirm that the data +was actually persisted. + #### Resp Reply Integer reply: UNIX TIME of the last DB save executed with success. diff --git a/website/docs/commands/server.md b/website/docs/commands/server.md index 5521ce39b1c..d748a5b4f27 100644 --- a/website/docs/commands/server.md +++ b/website/docs/commands/server.md @@ -229,7 +229,7 @@ Sections: * `STORE`: Per-database store details, including the current and last-checkpointed version, system state, hash-index bucket counts and sizes, and the hybrid-log and read-cache page/memory/heap sizes together with the key log addresses (`Log.BeginAddress`, `Log.HeadAddress`, `Log.SafeReadOnlyAddress`, `Log.FlushedUntilAddress`, `Log.TailAddress`, and the corresponding `ReadCache.*` addresses). * `STOREHASHTABLE`: Per-database dump of the hash-table bucket distribution (how records are spread across the index bucket chains); useful for diagnosing index sizing. * `STOREREVIV`: Per-database revivification (deleted-record free list) statistics. -* `PERSISTENCE`: Checkpoint and append-only-file persistence information. +* `PERSISTENCE`: Checkpoint and append-only-file persistence information. `rdb_last_bgsave_status` reports whether the last checkpoint attempt for the database succeeded (`ok`) or failed (`err`); the append-only-file addresses are reported as `N/A` when the append-only file is disabled. * `CLIENTS`: Connected-client statistics. * `KEYSPACE`: Per-database key counts. Requested explicitly only, since it requires a full log scan. * `MODULES`: Loaded module information.