From 436605d1ae739c506b26885953fb153b4c331c0c Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Wed, 30 Sep 2026 16:47:36 +0800 Subject: [PATCH] Ensure archive groups only act on a fixed number of rows --- .../Recoverability/MessageArchiver.cs | 97 ++++- .../ArchiveOperationTotalSyncTests.cs | 225 ++++++++++++ .../EFCore/ArchiveGroupBudgetTests.cs | 339 ++++++++++++++++++ .../ArchiveGroupArrivalsTests.cs | 158 ++++++++ .../ArchiveProgressReconciliationTests.cs | 144 ++++++++ .../Archiving/InMemoryArchive.cs | 12 + .../Archiving/InMemoryUnarchive.cs | 12 + 7 files changed, 970 insertions(+), 17 deletions(-) create mode 100644 src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveOperationTotalSyncTests.cs create mode 100644 src/ServiceControl.Persistence.Tests/EFCore/ArchiveGroupBudgetTests.cs create mode 100644 src/ServiceControl.Persistence.Tests/Recoverability/ArchiveGroupArrivalsTests.cs create mode 100644 src/ServiceControl.Persistence.Tests/Recoverability/ArchiveProgressReconciliationTests.cs diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs index 90a73be8ca..491be0569c 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs @@ -60,22 +60,40 @@ public async Task ArchiveAllInGroup(string groupId, AuditUser? initiatedBy = nul // ── Start in-memory tracking ── await archivingManager.StartArchiving(operationEntity, cancellationToken); + // ── Batch loop ── - string[] batchIds; - do + // Bound the live query by the original batch budget; later arrivals remain for another operation. + while (operationEntity.CurrentBatch < operationEntity.NumberOfBatches + && operationEntity.NumberOfMessagesProcessed < operationEntity.TotalNumberOfMessages) { + string[] batchIds; + int affectedRows; + await using var batchScope = scopeFactory.CreateAsyncScope(); var batchDbContext = batchScope.ServiceProvider.GetRequiredService(); operationEntity = await batchDbContext.ArchiveOperations.FindAsync([groupId, ArchiveType.FailureGroup, ArchiveOperationType.Archive], cancellationToken) ?? throw new InvalidOperationException($"No in progress Archive Operation found for {groupId}"); - batchIds = await UpdateGroupStatusAsync(batchDbContext, groupId, FailedMessageStatus.Unresolved, FailedMessageStatus.Archived, batchSize, cancellationToken); - await archivingManager.BatchArchived(groupId, ArchiveType.FailureGroup, batchIds.Length, cancellationToken); + (batchIds, affectedRows) = await UpdateGroupStatusAsync(batchDbContext, groupId, FailedMessageStatus.Unresolved, FailedMessageStatus.Archived, batchSize, cancellationToken); + + if (batchIds.Length == 0) + { + // No group members left in the source status; the plan cannot make progress. + break; + } + + await archivingManager.BatchArchived(groupId, ArchiveType.FailureGroup, affectedRows, cancellationToken); // Update progress tracking operationEntity.CurrentBatch++; - operationEntity.NumberOfMessagesProcessed += batchIds.Length; + operationEntity.NumberOfMessagesProcessed += affectedRows; + + if (operationEntity.NumberOfMessagesProcessed > operationEntity.TotalNumberOfMessages) + { + operationEntity.TotalNumberOfMessages = operationEntity.NumberOfMessagesProcessed; + } + await batchDbContext.SaveChangesAsync(cancellationToken); // Raise batch domain event @@ -85,9 +103,22 @@ public async Task ArchiveAllInGroup(string groupId, AuditUser? initiatedBy = nul AuditArchivedMessages(MessageActionKind.Archive, Permissions.ErrorRecoverabilityGroupsArchive, auditUser, auditOperationId, batchIds); logger.LogInformation("Archiving of {MessageCount} messages from group {GroupId} completed", batchIds.Length, groupId); - } while (batchIds.Length >= batchSize); + + if (batchIds.Length < batchSize) + { + // Partial fetched batch: fewer source-status members remain than a full batch. + break; + } + } // ── Finalize ── + if (operationEntity.NumberOfMessagesProcessed > operationEntity.TotalNumberOfMessages) + { + // Only reachable for a resumed row persisted by a version that did not reconcile the + // two counts and whose plan was already complete, so no batch ran to widen it. + operationEntity.TotalNumberOfMessages = operationEntity.NumberOfMessagesProcessed; + } + logger.LogInformation("Archiving of group {GroupId} is complete", groupId); await archivingManager.ArchiveOperationFinalizing(groupId, ArchiveType.FailureGroup, cancellationToken); await archivingManager.ArchiveOperationCompleted(groupId, ArchiveType.FailureGroup, cancellationToken); @@ -139,21 +170,39 @@ public async Task UnarchiveAllInGroup(string groupId, AuditUser? initiatedBy = n } await unarchivingManager.StartUnarchiving(operationEntity, cancellationToken); - string[] batchIds; - do + + // ── Batch loop ── + // Bound the live query by the original batch budget; later arrivals remain for another operation. + while (operationEntity.CurrentBatch < operationEntity.NumberOfBatches + && operationEntity.NumberOfMessagesProcessed < operationEntity.TotalNumberOfMessages) { + string[] batchIds; + int affectedRows; + await using var batchScope = scopeFactory.CreateAsyncScope(); var batchDbContext = batchScope.ServiceProvider.GetRequiredService(); operationEntity = await batchDbContext.ArchiveOperations.FindAsync([groupId, ArchiveType.FailureGroup, ArchiveOperationType.UnArchive], cancellationToken) ?? throw new InvalidOperationException($"No in progress Unarchive Operation found for {groupId}"); - batchIds = await UpdateGroupStatusAsync(batchDbContext, groupId, FailedMessageStatus.Archived, FailedMessageStatus.Unresolved, batchSize, cancellationToken); + (batchIds, affectedRows) = await UpdateGroupStatusAsync(batchDbContext, groupId, FailedMessageStatus.Archived, FailedMessageStatus.Unresolved, batchSize, cancellationToken); - await unarchivingManager.BatchUnarchived(groupId, ArchiveType.FailureGroup, batchIds.Length, cancellationToken); + if (batchIds.Length == 0) + { + // No group members left in the source status; the plan cannot make progress. + break; + } + + await unarchivingManager.BatchUnarchived(groupId, ArchiveType.FailureGroup, affectedRows, cancellationToken); // Update progress tracking operationEntity.CurrentBatch++; - operationEntity.NumberOfMessagesProcessed += batchIds.Length; + operationEntity.NumberOfMessagesProcessed += affectedRows; + + if (operationEntity.NumberOfMessagesProcessed > operationEntity.TotalNumberOfMessages) + { + operationEntity.TotalNumberOfMessages = operationEntity.NumberOfMessagesProcessed; + } + await batchDbContext.SaveChangesAsync(cancellationToken); // Raise batch domain event @@ -163,9 +212,22 @@ public async Task UnarchiveAllInGroup(string groupId, AuditUser? initiatedBy = n AuditArchivedMessages(MessageActionKind.Unarchive, Permissions.ErrorRecoverabilityGroupsUnarchive, auditUser, auditOperationId, batchIds); logger.LogInformation("Unarchiving of {MessageCount} messages from group {GroupId} completed", batchIds.Length, groupId); - } while (batchIds.Length >= batchSize); + + if (batchIds.Length < batchSize) + { + // Partial fetched batch: fewer source-status members remain than a full batch. + break; + } + } // ── Finalize ── + if (operationEntity.NumberOfMessagesProcessed > operationEntity.TotalNumberOfMessages) + { + // Only reachable for a resumed row persisted by a version that did not reconcile the + // two counts and whose plan was already complete, so no batch ran to widen it. + operationEntity.TotalNumberOfMessages = operationEntity.NumberOfMessagesProcessed; + } + logger.LogInformation("Unarchiving of group {GroupId} is complete", groupId); await unarchivingManager.UnarchiveOperationFinalizing(groupId, ArchiveType.FailureGroup, cancellationToken); await unarchivingManager.UnarchiveOperationCompleted(groupId, ArchiveType.FailureGroup, cancellationToken); @@ -278,26 +340,27 @@ void AuditArchivedMessages(MessageActionKind kind, string permission, AuditUser } } - async Task UpdateGroupStatusAsync(ServiceControlDbContext dbContext, string groupId, FailedMessageStatus fromStatus, FailedMessageStatus toStatus, int batchSize, CancellationToken cancellationToken) + async Task<(string[] BatchIds, int AffectedRows)> UpdateGroupStatusAsync(ServiceControlDbContext dbContext, string groupId, FailedMessageStatus fromStatus, FailedMessageStatus toStatus, int batchSize, CancellationToken cancellationToken) { var batchIds = await GetNextBatch(dbContext, groupId, fromStatus, batchSize) .Select(x => x.UniqueMessageId) .ToListAsync(cancellationToken); + var affectedRows = 0; if (batchIds.Count > 0) { var now = timeProvider.GetUtcNow().UtcDateTime; - // Bulk status change with re-asserted status filter - await dbContext.FailedMessages - .Where(fm => batchIds.Contains(fm.UniqueMessageId)) + // A fetched row may have changed status before the update. + affectedRows = await dbContext.FailedMessages + .Where(fm => batchIds.Contains(fm.UniqueMessageId) && fm.Status == fromStatus) .ExecuteUpdateAsync(s => s .SetProperty(fm => fm.Status, toStatus) .SetProperty(fm => fm.StatusChangedAt, now) .SetProperty(fm => fm.LastModified, now), cancellationToken); } - return batchIds.Select(id => id.ToString()).ToArray(); + return (batchIds.Select(id => id.ToString()).ToArray(), affectedRows); } static async Task<(int count, string groupName)> GetGroupDetails( diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveOperationTotalSyncTests.cs b/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveOperationTotalSyncTests.cs new file mode 100644 index 0000000000..27f8e9e333 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveOperationTotalSyncTests.cs @@ -0,0 +1,225 @@ +namespace ServiceControl.Persistence.Tests.RavenDB.Archiving; + +using System; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Infrastructure.DomainEvents; +using ServiceControl.MessageFailures; +using ServiceControl.Recoverability; + +/// +/// RavenDB iterates a fixed set of pre-built batches while the operation's total comes from a +/// separate group-count index, so the batches can carry more message ids than the planned total. +/// The persisted operation document must adopt the actually processed count when batches overrun +/// the total, so interim progress never reports a negative remaining count and the widened total +/// is what a resumed operation and the completion event see. +/// +[TestFixture] +class ArchiveOperationTotalSyncTests : RavenPersistenceTestBase +{ + readonly ProbingDomainEvents events = new(); + + public ArchiveOperationTotalSyncTests() => + RegisterServices = services => services.AddSingleton(events); + + [Test] + public async Task Archive_operation_widens_the_persisted_total_when_batches_overrun_it() + { + const string groupId = "TestGroup"; + + using (var session = DocumentStore.OpenAsyncSession()) + { + foreach (var id in new[] { "A", "B", "C" }) + { + await session.StoreAsync(new FailedMessage + { + Id = "FailedMessages/" + id, + UniqueMessageId = id, + Status = FailedMessageStatus.Unresolved + }); + } + + // The count index reported 1 message but the batch stream supplies 3 ids. + await session.StoreAsync(new ArchiveBatch + { + Id = ArchiveBatch.MakeId(groupId, ArchiveType.FailureGroup, 0), + DocumentIds = ["FailedMessages/A", "FailedMessages/B", "FailedMessages/C"] + }); + + await session.StoreAsync(new ArchiveOperation + { + Id = ArchiveOperation.MakeId(groupId, ArchiveType.FailureGroup), + RequestId = groupId, + ArchiveType = ArchiveType.FailureGroup, + TotalNumberOfMessages = 1, + NumberOfMessagesArchived = 0, + Started = DateTime.UtcNow, + GroupName = "Test Group", + NumberOfBatches = 1, + CurrentBatch = 0 + }); + + await session.SaveChangesAsync(); + } + + ArchiveOperation storedAfterBatch = null; + events.OnRaised = async domainEvent => + { + if (domainEvent is FailedMessageGroupBatchArchived && storedAfterBatch == null) + { + using var session = DocumentStore.OpenAsyncSession(); + storedAfterBatch = await session.LoadAsync(ArchiveOperation.MakeId(groupId, ArchiveType.FailureGroup)); + } + }; + + await ArchiveMessages.ArchiveAllInGroup(groupId); + + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0), "remaining must never go negative"); + Assert.That(progress.Last().TotalNumberOfMessages, Is.EqualTo(3), "the planned total must adopt the processed count"); + Assert.That(progress.Last().NumberOfMessagesArchived, Is.EqualTo(3)); + Assert.That(storedAfterBatch, Is.Not.Null, "the operation document should still exist after the batch"); + Assert.That(storedAfterBatch.TotalNumberOfMessages, Is.EqualTo(3), "the persisted total must be synchronized, so a resume starts from the widened total"); + Assert.That(storedAfterBatch.NumberOfMessagesArchived, Is.EqualTo(3)); + Assert.That(completed.MessagesCount, Is.EqualTo(3), "the completion event reports the widened total"); + } + } + + [Test] + public async Task Unarchive_operation_widens_the_persisted_total_when_batches_overrun_it() + { + const string groupId = "TestGroup"; + + using (var session = DocumentStore.OpenAsyncSession()) + { + foreach (var id in new[] { "A", "B", "C" }) + { + await session.StoreAsync(new FailedMessage + { + Id = "FailedMessages/" + id, + UniqueMessageId = id, + Status = FailedMessageStatus.Archived + }); + } + + await session.StoreAsync(new UnarchiveBatch + { + Id = UnarchiveBatch.MakeId(groupId, ArchiveType.FailureGroup, 0), + DocumentIds = ["FailedMessages/A", "FailedMessages/B", "FailedMessages/C"] + }); + + await session.StoreAsync(new UnarchiveOperation + { + Id = UnarchiveOperation.MakeId(groupId, ArchiveType.FailureGroup), + RequestId = groupId, + ArchiveType = ArchiveType.FailureGroup, + TotalNumberOfMessages = 1, + NumberOfMessagesUnarchived = 0, + Started = DateTime.UtcNow, + GroupName = "Test Group", + NumberOfBatches = 1, + CurrentBatch = 0 + }); + + await session.SaveChangesAsync(); + } + + UnarchiveOperation storedAfterBatch = null; + events.OnRaised = async domainEvent => + { + if (domainEvent is FailedMessageGroupBatchUnarchived && storedAfterBatch == null) + { + using var session = DocumentStore.OpenAsyncSession(); + storedAfterBatch = await session.LoadAsync(UnarchiveOperation.MakeId(groupId, ArchiveType.FailureGroup)); + } + }; + + await ArchiveMessages.UnarchiveAllInGroup(groupId); + + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0), "remaining must never go negative"); + Assert.That(progress.Last().TotalNumberOfMessages, Is.EqualTo(3), "the planned total must adopt the processed count"); + Assert.That(progress.Last().NumberOfMessagesUnarchived, Is.EqualTo(3)); + Assert.That(storedAfterBatch, Is.Not.Null, "the operation document should still exist after the batch"); + Assert.That(storedAfterBatch.TotalNumberOfMessages, Is.EqualTo(3), "the persisted total must be synchronized, so a resume starts from the widened total"); + Assert.That(storedAfterBatch.NumberOfMessagesUnarchived, Is.EqualTo(3)); + Assert.That(completed.MessagesCount, Is.EqualTo(3), "the completion event reports the widened total"); + } + } + + [Test] + public async Task Archive_operation_resumed_with_a_processed_count_above_the_total_reports_non_negative_progress() + { + const string groupId = "TestGroup"; + + using (var session = DocumentStore.OpenAsyncSession()) + { + await session.StoreAsync(new FailedMessage + { + Id = "FailedMessages/A", + UniqueMessageId = "A", + Status = FailedMessageStatus.Unresolved + }); + + await session.StoreAsync(new ArchiveBatch + { + Id = ArchiveBatch.MakeId(groupId, ArchiveType.FailureGroup, 0), + DocumentIds = ["FailedMessages/A"] + }); + + // Persisted by a version that did not reconcile the plan: 2 processed of 1 planned. + await session.StoreAsync(new ArchiveOperation + { + Id = ArchiveOperation.MakeId(groupId, ArchiveType.FailureGroup), + RequestId = groupId, + ArchiveType = ArchiveType.FailureGroup, + TotalNumberOfMessages = 1, + NumberOfMessagesArchived = 2, + Started = DateTime.UtcNow, + GroupName = "Test Group", + NumberOfBatches = 1, + CurrentBatch = 0 + }); + + await session.SaveChangesAsync(); + } + + await ArchiveMessages.ArchiveAllInGroup(groupId); + + var starting = events.Raised.OfType().Single(); + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(starting.Progress.MessagesRemaining, Is.EqualTo(0), "a resumed operation must not report a negative remaining count"); + Assert.That(starting.Progress.TotalNumberOfMessages, Is.EqualTo(2), "the persisted processed count wins over the stale plan"); + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0)); + Assert.That(completed.MessagesCount, Is.EqualTo(3), "completion reports the processed count of the resumed run"); + } + } + + class ProbingDomainEvents : IDomainEvents + { + public Func OnRaised { get; set; } = _ => Task.CompletedTask; + + public System.Collections.Generic.List Raised { get; } = []; + + public async Task Raise(T domainEvent, CancellationToken cancellationToken = default) where T : IDomainEvent + { + Raised.Add(domainEvent); + await OnRaised(domainEvent); + } + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests/EFCore/ArchiveGroupBudgetTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/ArchiveGroupBudgetTests.cs new file mode 100644 index 0000000000..bc7f1c1afe --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/ArchiveGroupBudgetTests.cs @@ -0,0 +1,339 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.Collections.Generic; +using System.Data.Common; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Diagnostics; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Infrastructure.DomainEvents; +using ServiceControl.MessageFailures; +using ServiceControl.Persistence.EFCore.DbContexts; +using ServiceControl.Persistence.EFCore.Entities; +using ServiceControl.Recoverability; + +/// +/// EF Core archives a failure group with a batch loop over live group membership, so batch counts +/// can overrun the totals planned when the operation started. These tests pin the invariants of +/// the bounded loop: the persisted total is widened when the processed count overruns it (so the +/// reported remaining count never goes negative, and a resume starts from the widened total), only +/// rows actually affected by the re-asserted status update are counted as processed, and a resumed +/// operation whose plan is already complete does not process a new batch. +/// +[TestFixture] +class ArchiveGroupBudgetTests : ErrorIngestionTestBase +{ + static readonly DateTime Noon = new(2026, 7, 22, 12, 0, 0, DateTimeKind.Utc); + + readonly HookingDomainEvents events = new(); + readonly FlipOneUnresolvedRowBeforeStatusUpdate flipOneRow = new(); + + public ArchiveGroupBudgetTests() + { + var registerServices = RegisterServices; + RegisterServices = services => + { + registerServices(services); + services.AddSingleton(events); + PersistenceTestsContext.InterceptDatabaseCommands(services, flipOneRow); + }; + } + + [TearDown] + public void DisarmInterceptor() => flipOneRow.Disarm(); + + [Test] + [CancelAfter(180_000)] + public async Task A_batch_overrunning_the_planned_total_widens_the_persisted_total() + { + var group = NewGroup(); + await Insert(group, 1500, FailedMessageStatus.Unresolved); + + var snapshots = new List(); + events.OnRaised = async domainEvent => + { + if (domainEvent is FailedMessageGroupBatchArchived) + { + if (snapshots.Count == 0) + { + // 700 failures join the group after the plan (1500) was captured + await Insert(group, 700, FailedMessageStatus.Unresolved); + } + + snapshots.Add(await Query(db => db.ArchiveOperations.AsNoTracking().SingleAsync())); + } + }; + + await ArchiveMessages.ArchiveAllInGroup(group.Id, cancellationToken: TestContext.CurrentContext.CancellationToken); + + await CompleteDatabaseOperation(); + + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(snapshots, Has.Count.EqualTo(2), "the planned batch budget (ceil(1500/1000)) caps the loop at two batches"); + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0), "remaining must never go negative"); + Assert.That(progress.Last().TotalNumberOfMessages, Is.EqualTo(2000), "the reported total adopts the processed count"); + Assert.That(progress.Last().NumberOfMessagesArchived, Is.EqualTo(2000)); + Assert.That(snapshots.Last().TotalNumberOfMessages, Is.EqualTo(2000), "the persisted total must adopt the processed count, so a resume starts from it"); + Assert.That(snapshots.Last().NumberOfMessagesProcessed, Is.EqualTo(2000)); + Assert.That(completed.MessagesCount, Is.EqualTo(2000), "the completion event reports the widened total"); + } + + using (Assert.EnterMultipleScope()) + { + var leftover = await GroupsStore.GetUnresolvedGroup(group.Id, null, null); + Assert.That(leftover.Results?.Count, Is.EqualTo(200), "members that joined late are left for a later operation"); + Assert.That(await Query(db => db.ArchiveOperations.CountAsync()), Is.Zero, "the operation row is removed on completion"); + } + } + + [Test] + [CancelAfter(180_000)] + public async Task An_unarchive_batch_overrunning_the_planned_total_widens_the_persisted_total() + { + var group = NewGroup(); + await Insert(group, 1500, FailedMessageStatus.Archived); + + var snapshots = new List(); + events.OnRaised = async domainEvent => + { + if (domainEvent is FailedMessageGroupBatchUnarchived) + { + if (snapshots.Count == 0) + { + // 700 failures join the group after the plan (1500) was captured + await Insert(group, 700, FailedMessageStatus.Archived); + } + + snapshots.Add(await Query(db => db.ArchiveOperations.AsNoTracking().SingleAsync())); + } + }; + + await ArchiveMessages.UnarchiveAllInGroup(group.Id, cancellationToken: TestContext.CurrentContext.CancellationToken); + + await CompleteDatabaseOperation(); + + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(snapshots, Has.Count.EqualTo(2), "the planned batch budget (ceil(1500/1000)) caps the loop at two batches"); + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0), "remaining must never go negative"); + Assert.That(progress.Last().TotalNumberOfMessages, Is.EqualTo(2000), "the reported total adopts the processed count"); + Assert.That(progress.Last().NumberOfMessagesUnarchived, Is.EqualTo(2000)); + Assert.That(snapshots.Last().TotalNumberOfMessages, Is.EqualTo(2000), "the persisted total must adopt the processed count, so a resume starts from it"); + Assert.That(snapshots.Last().NumberOfMessagesProcessed, Is.EqualTo(2000)); + Assert.That(completed.MessagesCount, Is.EqualTo(2000), "the completion event reports the widened total"); + } + + using (Assert.EnterMultipleScope()) + { + var leftover = await GroupsStore.GetArchivedGroup(group.Id, null, null); + Assert.That(leftover.Results?.Count, Is.EqualTo(200), "members that joined late are left for a later operation"); + Assert.That(await Query(db => db.ArchiveOperations.CountAsync()), Is.Zero, "the operation row is removed on completion"); + } + } + + [Test] + [CancelAfter(180_000)] + public async Task Only_rows_still_in_the_source_status_are_counted_as_processed() + { + var group = NewGroup(); + await Insert(group, 3, FailedMessageStatus.Unresolved); + + var snapshots = new List(); + events.OnRaised = async domainEvent => + { + if (domainEvent is FailedMessageGroupBatchArchived) + { + snapshots.Add(await Query(db => db.ArchiveOperations.AsNoTracking().SingleAsync())); + } + }; + + // Between the fetch of the batch and the bulk status update, a concurrent change moves one + // group member out of the source status. The re-asserted update must skip that row and the + // progress counters must only claim the rows the update actually affected. + flipOneRow.Arm(async () => + { + await using var scope = ServiceProvider.CreateAsyncScope(); + var dbContext = scope.ServiceProvider.GetRequiredService(); + var victim = await dbContext.FailedMessages + .Where(fm => fm.Status == FailedMessageStatus.Unresolved + && dbContext.FailedMessageGroups.Any(g => g.GroupId == group.Id && g.FailedMessageUniqueId == fm.UniqueMessageId)) + .OrderBy(fm => fm.UniqueMessageId) + .FirstAsync(); + await dbContext.FailedMessages + .Where(fm => fm.UniqueMessageId == victim.UniqueMessageId) + .ExecuteUpdateAsync(s => s.SetProperty(fm => fm.Status, FailedMessageStatus.Resolved)); + }); + + await ArchiveMessages.ArchiveAllInGroup(group.Id, cancellationToken: TestContext.CurrentContext.CancellationToken); + + Assert.That(flipOneRow.WasTriggered, Is.True, "the status update must pass through the race interceptor"); + var batchEvent = events.Raised.OfType().Single(); + var progress = events.Raised.OfType().Select(e => e.Progress).Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(batchEvent.FailedMessagesIds, Has.Length.EqualTo(3), "the fetch selected all three group members"); + Assert.That(progress.NumberOfMessagesArchived, Is.EqualTo(2), "only the rows the status update affected are counted as archived"); + Assert.That(progress.MessagesRemaining, Is.EqualTo(1), "the concurrently changed row is not claimed as processed"); + Assert.That(progress.Percentage, Is.EqualTo(Math.Round(2 / 3.0, 2))); + Assert.That(progress.TotalNumberOfMessages, Is.EqualTo(3), "the plan is not widened when nothing overran it"); + Assert.That(snapshots.Single().NumberOfMessagesProcessed, Is.EqualTo(2), "the persisted checkpoint counts affected rows only"); + } + + using (Assert.EnterMultipleScope()) + { + Assert.That(await Query(db => db.FailedMessages.CountAsync(fm => fm.Status == FailedMessageStatus.Archived)), Is.EqualTo(2)); + Assert.That(await Query(db => db.FailedMessages.CountAsync(fm => fm.Status == FailedMessageStatus.Resolved)), Is.EqualTo(1), "the concurrently changed row keeps its new status"); + } + } + + [Test] + [CancelAfter(180_000)] + public async Task A_resumed_operation_whose_plan_is_already_complete_does_not_process_a_new_batch() + { + var group = NewGroup(); + + // A crash between the last batch and finalization leaves the plan complete but unfinalized; + // the row was persisted by a version that did not reconcile the plan with what it processed + await Store(new ArchiveOperationEntity + { + RequestId = group.Id, + GroupName = group.Title, + ArchiveType = ArchiveType.FailureGroup, + OperationType = ArchiveOperationType.Archive, + TotalNumberOfMessages = 2, + NumberOfMessagesProcessed = 3, + NumberOfBatches = 1, + CurrentBatch = 1, + Started = Now + }); + + await Insert(group, 2, FailedMessageStatus.Unresolved); // arrived after the plan was completed + + await ArchiveMessages.ArchiveAllInGroup(group.Id, cancellationToken: TestContext.CurrentContext.CancellationToken); + + var starting = events.Raised.OfType().Single(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(events.Raised.OfType(), Is.Empty, "no new batch may be processed once the planned total has been reached"); + Assert.That(starting.Progress.MessagesRemaining, Is.EqualTo(0), "a resumed operation must not report a negative remaining count"); + Assert.That(starting.Progress.TotalNumberOfMessages, Is.EqualTo(3), "the persisted processed count wins over the stale plan"); + Assert.That(completed.MessagesCount, Is.EqualTo(3), "the completion event reports the reconciled total"); + } + + using (Assert.EnterMultipleScope()) + { + var leftover = await GroupsStore.GetUnresolvedGroup(group.Id, null, null); + Assert.That(leftover.Results?.Count, Is.EqualTo(2), "the late arrivals are left for a later operation"); + Assert.That(await Query(db => db.ArchiveOperations.CountAsync()), Is.Zero, "the resumed operation row is removed on completion"); + } + } + + static FailedMessage.FailureGroup NewGroup() => + new() { Id = Guid.NewGuid().ToString(), Title = "OrderPlaced", Type = "Message Type" }; + + static IngestedFailure InGroup(FailedMessage.FailureGroup group) => + new() + { + Groups = [group], + AttemptedAt = Noon, + TimeOfFailure = Noon, + TimeSent = Noon.AddMinutes(-1) + }; + + async Task Insert(FailedMessage.FailureGroup group, int count, FailedMessageStatus status) + { + var failures = Enumerable.Range(0, count).Select(_ => InGroup(group)).Select(f => f.ToFailedMessage(status)).ToArray(); + + foreach (var message in failures) + { + message.Id = PersistenceTestsContext.GenerateFailedMessageRecordId(message.UniqueMessageId); + } + + await PersistenceTestsContext.InsertFailedMessages(failures); + await CompleteDatabaseOperation(); + } + + /// + /// Records every domain event and offers a hook that runs synchronously when each is raised. + /// + sealed class HookingDomainEvents : IDomainEvents + { + public Func OnRaised { get; set; } = _ => Task.CompletedTask; + + public List Raised { get; } = []; + + public async Task Raise(T domainEvent, CancellationToken cancellationToken = default) where T : IDomainEvent + { + cancellationToken.ThrowIfCancellationRequested(); + + Raised.Add(domainEvent); + await OnRaised(domainEvent); + } + } + + /// + /// Arms a one-shot change that runs between the archiver's batch fetch and its bulk status + /// update, simulating a concurrent status change on a row the fetch already selected. EF + /// renders that update as the only UPDATE on FailedMessages that sets StatusChangedAt. + /// + sealed class FlipOneUnresolvedRowBeforeStatusUpdate : DbCommandInterceptor + { + Func armed; + public bool WasTriggered { get; private set; } + + public void Arm(Func flip) => Interlocked.Exchange(ref armed, flip); + + public void Disarm() => Interlocked.Exchange(ref armed, null); + + public override async ValueTask> NonQueryExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) + { + await FlipBeforeUpdate(command); + return await base.NonQueryExecutingAsync(command, eventData, result, cancellationToken); + } + + public override async ValueTask> ReaderExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) + { + await FlipBeforeUpdate(command); + return await base.ReaderExecutingAsync(command, eventData, result, cancellationToken); + } + + public override async ValueTask> ScalarExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) + { + await FlipBeforeUpdate(command); + return await base.ScalarExecutingAsync(command, eventData, result, cancellationToken); + } + + async Task FlipBeforeUpdate(DbCommand command) + { + if (IsFailedMessagesStatusUpdate(command) && Interlocked.Exchange(ref armed, null) is { } flip) + { + await flip(); + WasTriggered = true; + } + } + + // Seed inserts and the operation-row bookkeeping (INSERT/UPDATE/DELETE on ArchiveOperations) + // do not match; the command shape is checked before the one-shot arm is consumed. + static bool IsFailedMessagesStatusUpdate(DbCommand command) + { + var sql = command.CommandText.Replace("_", ""); // PostgreSQL uses snake_case names. + return sql.Contains("UPDATE", StringComparison.OrdinalIgnoreCase) + && sql.Contains("FailedMessages", StringComparison.OrdinalIgnoreCase) + && sql.Contains("StatusChangedAt", StringComparison.OrdinalIgnoreCase); + } + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/ArchiveGroupArrivalsTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/ArchiveGroupArrivalsTests.cs new file mode 100644 index 0000000000..25c523a21d --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/Recoverability/ArchiveGroupArrivalsTests.cs @@ -0,0 +1,158 @@ +namespace ServiceControl.Persistence.Tests.Recoverability; + +using System; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Infrastructure.DomainEvents; +using ServiceControl.MessageFailures; +using ServiceControl.Recoverability; + +/// +/// Failure-group archive/unarchive must stay within the budget planned when the operation started +/// even if new failures keep joining the group while it runs, and its reported progress must never +/// go negative. Each batch completion injects new failures into the group, simulating live +/// arrivals: an unbounded batch loop (or a remaining count computed from a stale plan) makes these +/// assertions fail. Members that join late are accepted as leftovers for a later operation. +/// +[TestFixture] +class ArchiveGroupArrivalsTests : PersistenceTestBase +{ + const int BatchSize = 1000; // must match the persisters' archive batch size + const int MaxArrivalInjections = 10; // bounds the run when the loop is unbounded + + static readonly DateTime Noon = new(2026, 7, 22, 12, 0, 0, DateTimeKind.Utc); + + readonly ArrivalsDomainEvents events = new(); + + public ArchiveGroupArrivalsTests() => + RegisterServices = services => services.AddSingleton(events); + + [Test] + [CancelAfter(180_000)] + public async Task Archive_group_stays_within_its_planned_budget_when_failures_arrive_mid_operation(CancellationToken cancellationToken = default) + { + var group = NewGroup(); + await Insert(group, 2 * BatchSize, FailedMessageStatus.Unresolved); + + events.OnRaised = domainEvent => + { + if (domainEvent is FailedMessageGroupBatchArchived && events.ArrivalInjections < MaxArrivalInjections) + { + events.ArrivalInjections++; + return Insert(group, BatchSize, FailedMessageStatus.Unresolved); + } + + return Task.CompletedTask; + }; + + await ArchiveMessages.ArchiveAllInGroup(group.Id, cancellationToken: cancellationToken); + + await CompleteDatabaseOperation(); + + var batchEvents = events.Raised.OfType().ToArray(); + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(batchEvents, Has.Length.EqualTo(2), "the planned batch count (2 x 1000) is a hard cap"); + Assert.That(batchEvents.Select(e => e.FailedMessagesIds.Length), Is.EqualTo(new[] { BatchSize, BatchSize })); + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0), "remaining must never go negative"); + Assert.That(progress, Has.All.Matches(p => p.TotalNumberOfMessages >= p.NumberOfMessagesArchived)); + Assert.That(completed.MessagesCount, Is.EqualTo(2 * BatchSize), "the completed event reports the planned total"); + } + + var leftover = await GroupsStore.GetUnresolvedGroup(group.Id, null, null, cancellationToken); + Assert.That(leftover.Results?.Count, Is.EqualTo(2 * BatchSize), + "members that joined late are left for a later operation, and the loop must not run past its budget"); + } + + [Test] + [CancelAfter(180_000)] + public async Task Unarchive_group_stays_within_its_planned_budget_when_failures_arrive_mid_operation(CancellationToken cancellationToken = default) + { + var group = NewGroup(); + await Insert(group, 2 * BatchSize, FailedMessageStatus.Archived); + + events.OnRaised = domainEvent => + { + if (domainEvent is FailedMessageGroupBatchUnarchived && events.ArrivalInjections < MaxArrivalInjections) + { + events.ArrivalInjections++; + return Insert(group, BatchSize, FailedMessageStatus.Archived); + } + + return Task.CompletedTask; + }; + + await ArchiveMessages.UnarchiveAllInGroup(group.Id, cancellationToken: cancellationToken); + + await CompleteDatabaseOperation(); + + var batchEvents = events.Raised.OfType().ToArray(); + var progress = events.Raised.OfType().Select(e => e.Progress).ToArray(); + var completed = events.Raised.OfType().Single(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(batchEvents, Has.Length.EqualTo(2), "the planned batch count (2 x 1000) is a hard cap"); + Assert.That(batchEvents.Select(e => e.FailedMessagesIds.Length), Is.EqualTo(new[] { BatchSize, BatchSize })); + Assert.That(progress, Has.All.Matches(p => p.MessagesRemaining >= 0), "remaining must never go negative"); + Assert.That(progress, Has.All.Matches(p => p.TotalNumberOfMessages >= p.NumberOfMessagesUnarchived)); + Assert.That(completed.MessagesCount, Is.EqualTo(2 * BatchSize), "the completed event reports the planned total"); + } + + var leftover = await GroupsStore.GetArchivedGroup(group.Id, null, null, cancellationToken); + Assert.That(leftover.Results?.Count, Is.EqualTo(2 * BatchSize), + "members that joined late are left for a later operation, and the loop must not run past its budget"); + } + + static FailedMessage.FailureGroup NewGroup() => + new() { Id = Guid.NewGuid().ToString(), Title = "OrderPlaced", Type = "Message Type" }; + + static IngestedFailure InGroup(FailedMessage.FailureGroup group) => + new() + { + Groups = [group], + AttemptedAt = Noon, + TimeOfFailure = Noon, + TimeSent = Noon.AddMinutes(-1) + }; + + async Task Insert(FailedMessage.FailureGroup group, int count, FailedMessageStatus status) + { + var failures = Enumerable.Range(0, count).Select(_ => InGroup(group)).Select(f => f.ToFailedMessage(status)).ToArray(); + + foreach (var message in failures) + { + message.Id = PersistenceTestsContext.GenerateFailedMessageRecordId(message.UniqueMessageId); + } + + await PersistenceTestsContext.InsertFailedMessages(failures); + await CompleteDatabaseOperation(); + } + + /// + /// Records every domain event and offers a hook that runs when each is raised. Cancellation is + /// honoured on the way in so a runaway loop is cut short once the test timeout fires. + /// + sealed class ArrivalsDomainEvents : IDomainEvents + { + public Func OnRaised { get; set; } = _ => Task.CompletedTask; + + public System.Collections.Generic.List Raised { get; } = []; + + public int ArrivalInjections { get; set; } + + public async Task Raise(T domainEvent, CancellationToken cancellationToken = default) where T : IDomainEvent + { + cancellationToken.ThrowIfCancellationRequested(); + + Raised.Add(domainEvent); + await OnRaised(domainEvent); + } + } +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/ArchiveProgressReconciliationTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/ArchiveProgressReconciliationTests.cs new file mode 100644 index 0000000000..6ebdd558dc --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/Recoverability/ArchiveProgressReconciliationTests.cs @@ -0,0 +1,144 @@ +namespace ServiceControl.Persistence.Tests.Recoverability; + +using System.Linq; +using System.Threading.Tasks; +using Microsoft.Extensions.Time.Testing; +using NUnit.Framework; +using ServiceControl.Recoverability; + +/// +/// Regression tests for the archive/unarchive progress counters shared by every persister: +/// when the batches processed for an operation overran the total planned when it started (live +/// group membership grew mid-operation, or the persisted plan undercounted), the total is widened +/// to the processed count so the reported remaining count never goes negative, both while batches +/// are being counted and when an operation is resumed from persisted state. +/// +[TestFixture] +class ArchiveProgressReconciliationTests +{ + readonly FakeTimeProvider timeProvider = new(); + + [Test] + public async Task Archive_batches_overrunning_the_total_widen_it_so_remaining_never_goes_negative() + { + var archive = NewArchive(total: 1000); + + await archive.Start(); + await archive.BatchArchived(1000); + await archive.BatchArchived(500); // 500 messages joined the group after the plan was made + + var progress = archive.GetProgress(); + using (Assert.EnterMultipleScope()) + { + Assert.That(progress.TotalNumberOfMessages, Is.EqualTo(1500), "the total must adopt the processed count"); + Assert.That(progress.NumberOfMessagesArchived, Is.EqualTo(1500)); + Assert.That(progress.MessagesRemaining, Is.EqualTo(0), "remaining must not go negative"); + Assert.That(progress.Percentage, Is.EqualTo(1.0), "the progress fraction must not exceed 1"); + } + } + + [Test] + public async Task Archive_progress_events_report_widened_totals() + { + var events = new FakeDomainEvents(); + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, events, timeProvider) { TotalNumberOfMessages = 1000 }; + + await archive.Start(); + await archive.BatchArchived(1500); + + var batchCompleted = events.RaisedEvents.OfType().Single(); + using (Assert.EnterMultipleScope()) + { + Assert.That(batchCompleted.Progress.TotalNumberOfMessages, Is.EqualTo(1500)); + Assert.That(batchCompleted.Progress.NumberOfMessagesArchived, Is.EqualTo(1500)); + Assert.That(batchCompleted.Progress.MessagesRemaining, Is.EqualTo(0)); + Assert.That(batchCompleted.Progress.Percentage, Is.EqualTo(1.0)); + } + } + + [Test] + public async Task An_archive_operation_resumed_with_a_processed_count_above_the_total_reconciles_on_start() + { + var events = new FakeDomainEvents(); + var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, events, timeProvider) + { + TotalNumberOfMessages = 1000, + NumberOfMessagesArchived = 1500 // persisted by a plan that undercounted + }; + + await archive.Start(); + + var starting = events.RaisedEvents.OfType().Single(); + using (Assert.EnterMultipleScope()) + { + Assert.That(starting.Progress.TotalNumberOfMessages, Is.EqualTo(1500), "the persisted processed count wins over the stale plan"); + Assert.That(starting.Progress.NumberOfMessagesArchived, Is.EqualTo(1500)); + Assert.That(starting.Progress.MessagesRemaining, Is.EqualTo(0), "a resumed operation must not report a negative remaining count"); + } + } + + [Test] + public async Task Unarchive_batches_overrunning_the_total_widen_it_so_remaining_never_goes_negative() + { + var unarchive = NewUnarchive(total: 1000); + + await unarchive.Start(); + await unarchive.BatchUnarchived(1000); + await unarchive.BatchUnarchived(500); + + var progress = unarchive.GetProgress(); + using (Assert.EnterMultipleScope()) + { + Assert.That(progress.TotalNumberOfMessages, Is.EqualTo(1500), "the total must adopt the processed count"); + Assert.That(progress.NumberOfMessagesUnarchived, Is.EqualTo(1500)); + Assert.That(progress.MessagesRemaining, Is.EqualTo(0), "remaining must not go negative"); + Assert.That(progress.Percentage, Is.EqualTo(1.0), "the progress fraction must not exceed 1"); + } + } + + [Test] + public async Task Unarchive_progress_events_report_widened_totals() + { + var events = new FakeDomainEvents(); + var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, events, timeProvider) { TotalNumberOfMessages = 1000 }; + + await unarchive.Start(); + await unarchive.BatchUnarchived(1500); + + var batchCompleted = events.RaisedEvents.OfType().Single(); + using (Assert.EnterMultipleScope()) + { + Assert.That(batchCompleted.Progress.TotalNumberOfMessages, Is.EqualTo(1500)); + Assert.That(batchCompleted.Progress.NumberOfMessagesUnarchived, Is.EqualTo(1500)); + Assert.That(batchCompleted.Progress.MessagesRemaining, Is.EqualTo(0)); + Assert.That(batchCompleted.Progress.Percentage, Is.EqualTo(1.0)); + } + } + + [Test] + public async Task An_unarchive_operation_resumed_with_a_processed_count_above_the_total_reconciles_on_start() + { + var events = new FakeDomainEvents(); + var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, events, timeProvider) + { + TotalNumberOfMessages = 1000, + NumberOfMessagesUnarchived = 1500 + }; + + await unarchive.Start(); + + var starting = events.RaisedEvents.OfType().Single(); + using (Assert.EnterMultipleScope()) + { + Assert.That(starting.Progress.TotalNumberOfMessages, Is.EqualTo(1500)); + Assert.That(starting.Progress.NumberOfMessagesUnarchived, Is.EqualTo(1500)); + Assert.That(starting.Progress.MessagesRemaining, Is.EqualTo(0)); + } + } + + InMemoryArchive NewArchive(int total) => + new("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), timeProvider) { TotalNumberOfMessages = total }; + + InMemoryUnarchive NewUnarchive(int total) => + new("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), timeProvider) { TotalNumberOfMessages = total }; +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs index bef83dd441..d7b0a2828f 100644 --- a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs @@ -46,6 +46,12 @@ public ArchiveProgress GetProgress() public Task Start(CancellationToken cancellationToken = default) { + if (NumberOfMessagesArchived > TotalNumberOfMessages) + { + // Reconcile progress from an older persisted operation. + TotalNumberOfMessages = NumberOfMessagesArchived; + } + ArchiveState = ArchiveState.ArchiveStarted; CompletionTime = null; operationMetrics?.Started(); @@ -63,6 +69,12 @@ public Task BatchArchived(int numberOfMessagesArchivedInBatch, CancellationToken { ArchiveState = ArchiveState.ArchiveProgressing; NumberOfMessagesArchived += numberOfMessagesArchivedInBatch; + + if (NumberOfMessagesArchived > TotalNumberOfMessages) + { + TotalNumberOfMessages = NumberOfMessagesArchived; + } + CurrentBatch++; Last = timeProvider.GetUtcNow().UtcDateTime; operationMetrics?.BatchCompleted(numberOfMessagesArchivedInBatch); diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs index 375f98f38d..923d8b20a5 100644 --- a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs @@ -46,6 +46,12 @@ public UnarchiveProgress GetProgress() public Task Start(CancellationToken cancellationToken = default) { + if (NumberOfMessagesUnarchived > TotalNumberOfMessages) + { + // Reconcile progress from an older persisted operation. + TotalNumberOfMessages = NumberOfMessagesUnarchived; + } + ArchiveState = ArchiveState.ArchiveStarted; CompletionTime = null; operationMetrics?.Started(); @@ -63,6 +69,12 @@ public Task BatchUnarchived(int numberOfMessagesUnarchivedInBatch, CancellationT { ArchiveState = ArchiveState.ArchiveProgressing; NumberOfMessagesUnarchived += numberOfMessagesUnarchivedInBatch; + + if (NumberOfMessagesUnarchived > TotalNumberOfMessages) + { + TotalNumberOfMessages = NumberOfMessagesUnarchived; + } + CurrentBatch++; Last = timeProvider.GetUtcNow().UtcDateTime; operationMetrics?.BatchCompleted(numberOfMessagesUnarchivedInBatch);