Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<ServiceControlDbContext>();

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
Expand All @@ -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);
Expand Down Expand Up @@ -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<ServiceControlDbContext>();
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
Expand All @@ -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);
Expand Down Expand Up @@ -278,26 +340,27 @@ void AuditArchivedMessages(MessageActionKind kind, string permission, AuditUser
}
}

async Task<string[]> 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(
Expand Down
Original file line number Diff line number Diff line change
@@ -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;

/// <summary>
/// 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.
/// </summary>
[TestFixture]
class ArchiveOperationTotalSyncTests : RavenPersistenceTestBase
{
readonly ProbingDomainEvents events = new();

public ArchiveOperationTotalSyncTests() =>
RegisterServices = services => services.AddSingleton<IDomainEvents>(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>(ArchiveOperation.MakeId(groupId, ArchiveType.FailureGroup));
}
};

await ArchiveMessages.ArchiveAllInGroup(groupId);

var progress = events.Raised.OfType<ArchiveOperationBatchCompleted>().Select(e => e.Progress).ToArray();
var completed = events.Raised.OfType<FailedMessageGroupArchived>().Single();

using (Assert.EnterMultipleScope())
{
Assert.That(progress, Has.All.Matches<ArchiveProgress>(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>(UnarchiveOperation.MakeId(groupId, ArchiveType.FailureGroup));
}
};

await ArchiveMessages.UnarchiveAllInGroup(groupId);

var progress = events.Raised.OfType<UnarchiveOperationBatchCompleted>().Select(e => e.Progress).ToArray();
var completed = events.Raised.OfType<FailedMessageGroupUnarchived>().Single();

using (Assert.EnterMultipleScope())
{
Assert.That(progress, Has.All.Matches<UnarchiveProgress>(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<ArchiveOperationStarting>().Single();
var progress = events.Raised.OfType<ArchiveOperationBatchCompleted>().Select(e => e.Progress).ToArray();
var completed = events.Raised.OfType<FailedMessageGroupArchived>().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<ArchiveProgress>(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<object, Task> OnRaised { get; set; } = _ => Task.CompletedTask;

public System.Collections.Generic.List<object> Raised { get; } = [];

public async Task Raise<T>(T domainEvent, CancellationToken cancellationToken = default) where T : IDomainEvent
{
Raised.Add(domainEvent);
await OnRaised(domainEvent);
}
}
}
Loading
Loading