From 785985a69f368a3315892e93dbc121f9145b4552 Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Thu, 13 Aug 2026 13:23:01 +0800 Subject: [PATCH 1/3] Convert EditFailedMessageManager to an atomic write operation --- docs/coding-and-design-guidelines.md | 17 ++ .../EditFailedMessagesDataStore.cs | 95 ++++++++- .../EditFailedMessagesManager.cs | 96 --------- .../Editing/EditFailedMessageManager.cs | 59 ------ .../Editing/EditFailedMessagesDataStore.cs | 63 +++++- .../RavenTransactionalDataStore.cs | 18 -- .../Editing/EditAcquisitionTests.cs | 51 +++++ ...erTests.cs => AtomicEditRetentionTests.cs} | 56 ++---- .../EditFailedMessagesDataStoreTests.cs | 182 ++++++++++++++++++ .../PersistenceTestBase.cs | 7 + .../Recoverability/EditMessageTests.cs | 31 ++- .../IDataSessionManager.cs | 11 -- .../IEditFailedMessagesDataStore.cs | 48 ++++- .../IEditFailedMessagesManager.cs | 15 -- .../EditFailedMessagesControllerAuditTests.cs | 20 +- .../Api/EditFailedMessagesController.cs | 5 +- .../Recoverability/Editing/EditHandler.cs | 49 ++--- 17 files changed, 495 insertions(+), 328 deletions(-) delete mode 100644 src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesManager.cs delete mode 100644 src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessageManager.cs delete mode 100644 src/ServiceControl.Persistence.RavenDB/Transactions/RavenTransactionalDataStore.cs create mode 100644 src/ServiceControl.Persistence.Tests.RavenDB/Editing/EditAcquisitionTests.cs rename src/ServiceControl.Persistence.Tests/EFCore/{EditFailedMessagesManagerTests.cs => AtomicEditRetentionTests.cs} (57%) create mode 100644 src/ServiceControl.Persistence.Tests/EditFailedMessagesDataStoreTests.cs delete mode 100644 src/ServiceControl.Persistence/IDataSessionManager.cs delete mode 100644 src/ServiceControl.Persistence/IEditFailedMessagesManager.cs diff --git a/docs/coding-and-design-guidelines.md b/docs/coding-and-design-guidelines.md index c58fb6654f..e3cecf4d1d 100644 --- a/docs/coding-and-design-guidelines.md +++ b/docs/coding-and-design-guidelines.md @@ -40,6 +40,23 @@ There are a few things that are still registered using convention. Note that the Additionally, because NServiceBus does a type-scan at startup it will automatically register any implementations of `Feature` and `IHandleMessage<>`. We have chosen to leave this alone as we would be fighting with NServiceBus in order to turn this off. +## Prefer explicit persistence operations + +Use a direct data-store method when all inputs for a persistence operation fit in one method call. The method should own and dispose its EF Core scope/context or RavenDB session, accept a `CancellationToken`, and commit before returning. Returned entities are detached snapshots; callers must not be required to mutate tracked entities as an implicit persistence command. + +Atomic operations that can encounter concurrency conflicts should document their provider guarantees and translate expected provider exceptions into explicit domain outcomes. For example, a unique-key or optimistic-concurrency conflict should not escape when contention is part of the operation's normal contract. + +Reserve a specialized unit of work for cases where a caller genuinely composes several writes into one atomic batch. New unit-of-work APIs should consistently provide: + +- an `I...UnitOfWorkFactory`; +- a `StartNew` factory method; +- a `Complete(CancellationToken)` commit method; +- `IAsyncDisposable` lifetime ownership; +- explicit operation-recording methods rather than mutation of tracked return values; +- documented commit, abandon, repeated-completion, and concurrency behavior. + +Do not introduce generic `IDataSessionManager`-style abstractions or persistence managers with hidden call-order protocols. During review, prefer one explicit store operation unless caller-composed atomicity requires a unit of work. + ## Avoid property injection Although the Autofac container can be configured to allow property injection, we prefer to avoid it. There is no way to specify property injection using the Microsoft DI abstractions, and the default `IServiceProvider` implementation does not support it. Where possible, use constructor injection instead. diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs index d902b300e1..b12f9ac729 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs @@ -1,14 +1,97 @@ namespace ServiceControl.Persistence.EFCore.Implementation; -using DbContexts; +using Entities; +using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; +using ServiceControl.MessageFailures; -public class EditFailedMessagesDataStore(IServiceScopeFactory scopeFactory, TimeProvider timeProvider) : IEditFailedMessagesDataStore +public class EditFailedMessagesDataStore(IServiceScopeFactory scopeFactory, TimeProvider timeProvider) + : DataStoreBase(scopeFactory), IEditFailedMessagesDataStore { - public Task CreateEditFailedMessageManager(CancellationToken cancellationToken = default) + public Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default) { - var scope = scopeFactory.CreateAsyncScope(); - return Task.FromResult( - new EditFailedMessagesManager(scope, scope.ServiceProvider.GetRequiredService(), timeProvider)); + if (!Guid.TryParse(failedMessageId, out var uniqueMessageId)) + { + return Task.FromResult(null); + } + + return ExecuteWithDbContext((dbContext, token) => dbContext.FailedMessageEdits + .AsNoTracking() + .Where(edit => edit.UniqueMessageId == uniqueMessageId) + .Select(edit => edit.EditId) + .SingleOrDefaultAsync(token), cancellationToken); + } + + public Task TryBeginEdit(string failedMessageId, string editingMessageId, CancellationToken cancellationToken = default) + { + if (!Guid.TryParse(failedMessageId, out var uniqueMessageId)) + { + return Task.FromResult(new BeginEditResult(BeginEditOutcome.MessageNotFound)); + } + + return ExecuteWithDbContext(async (dbContext, ct) => + { + var entity = await dbContext.FailedMessages + .SingleOrDefaultAsync(message => message.UniqueMessageId == uniqueMessageId, ct); + + if (entity is null) + { + return new BeginEditResult(BeginEditOutcome.MessageNotFound); + } + + if (entity.Status != FailedMessageStatus.Unresolved) + { + return new BeginEditResult(BeginEditOutcome.MessageNotUnresolved); + } + + var editId = await dbContext.FailedMessageEdits + .Where(edit => edit.UniqueMessageId == uniqueMessageId) + .Select(edit => edit.EditId) + .SingleOrDefaultAsync(ct); + + if (editId is null) + { + editId = editingMessageId; + dbContext.FailedMessageEdits.Add(new FailedMessageEditEntity { UniqueMessageId = uniqueMessageId, EditId = editingMessageId }); + + var now = timeProvider.GetUtcNow().UtcDateTime; + entity.Status = FailedMessageStatus.Resolved; + entity.StatusChangedAt = now; + entity.LastModified = now; + + try + { + await dbContext.SaveChangesAsync(ct); + } + catch (DbUpdateException exception) when (dbContext.IsDuplicateKeyException(exception)) + { + dbContext.ChangeTracker.Clear(); + + var winningEditId = await dbContext.FailedMessageEdits + .AsNoTracking() + .Where(edit => edit.UniqueMessageId == uniqueMessageId) + .Select(edit => edit.EditId) + .SingleOrDefaultAsync(ct); + + if (winningEditId is null) + { + throw; + } + + entity = await dbContext.FailedMessages + .AsNoTracking() + .SingleAsync(message => message.UniqueMessageId == uniqueMessageId, ct); + + editId = winningEditId; + } + } + + if (editId == editingMessageId) + { + return new BeginEditResult(BeginEditOutcome.Acquired, entity.ToFailedMessage([]), editId); + } + + return new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: editId); + }, cancellationToken); } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesManager.cs b/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesManager.cs deleted file mode 100644 index fb1eafc8e3..0000000000 --- a/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesManager.cs +++ /dev/null @@ -1,96 +0,0 @@ -namespace ServiceControl.Persistence.EFCore.Implementation; - -using Microsoft.EntityFrameworkCore; -using ServiceControl.MessageFailures; -using ServiceControl.Persistence.EFCore.DbContexts; -using ServiceControl.Persistence.EFCore.Entities; - -public class EditFailedMessagesManager(IAsyncDisposable scope, ServiceControlDbContext dbContext, TimeProvider timeProvider) : IEditFailedMessagesManager -{ - FailedMessage? failedMessage; // cached after GetFailedMessage - FailedMessageEditEntity? editEntity; // tracked after GetCurrentEditingRequestId / SetCurrentEditingRequestId - - public async Task GetFailedMessage(string failedMessageId, CancellationToken cancellationToken = default) - { - if (!Guid.TryParse(failedMessageId, out var uniqueMessageId)) - { - return null; - } - - var entity = await dbContext.FailedMessages - .AsNoTracking() - .SingleOrDefaultAsync(message => message.UniqueMessageId == uniqueMessageId, cancellationToken); - - if (entity == null) - { - return null!; - } - - // Reuse the same mapping helper as FailedMessageQueryDataStore so the manager and the - // query store don't diverge. The edit manager passes an empty group list (it does not - // need failure groups); ToFailedMessage tolerates an empty collection. - var result = entity.ToFailedMessage([]); - failedMessage = result; - return result; - } - - public async Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default) - { - if (!Guid.TryParse(failedMessageId, out var uniqueMessageId)) - { - return null!; - } - - editEntity = await dbContext.FailedMessageEdits - .AsNoTracking() - .SingleOrDefaultAsync(e => e.UniqueMessageId == uniqueMessageId, cancellationToken); - - return editEntity?.EditId; - } - - public Task SetCurrentEditingRequestId(string editingMessageId, CancellationToken cancellationToken = default) - { - if (failedMessage == null) - { - throw new InvalidOperationException("No failed message loaded"); - } - - editEntity = new FailedMessageEditEntity - { - UniqueMessageId = Guid.Parse(failedMessage.UniqueMessageId), - EditId = editingMessageId - }; - dbContext.FailedMessageEdits.Add(editEntity); - - return Task.CompletedTask; - } - - public async Task SetFailedMessageAsResolved(CancellationToken cancellationToken = default) - { - if (failedMessage == null) - { - throw new InvalidOperationException("No failed message loaded"); - } - - var message = failedMessage; - - // Critical: must update the tracked entity (not just the in-memory FailedMessage) and - // MUST set StatusChangedAt + LastModified so the retention sweeper's filter - // (StatusChangedAt < cutoff AND Status in Resolved/Archived) works correctly. Leaving - // StatusChangedAt at its previous value could cause the just-resolved message to be - // swept immediately. - var now = timeProvider.GetUtcNow().UtcDateTime; - var entity = await dbContext.FailedMessages - .SingleOrDefaultAsync(m => m.UniqueMessageId == Guid.Parse(message.UniqueMessageId), cancellationToken) - ?? throw new InvalidOperationException("Failed message entity not found"); - - entity.Status = FailedMessageStatus.Resolved; - entity.StatusChangedAt = now; - entity.LastModified = now; - message.Status = FailedMessageStatus.Resolved; - } - - public Task SaveChanges(CancellationToken cancellationToken = default) => dbContext.SaveChangesAsync(cancellationToken); - - public ValueTask DisposeAsync() => scope.DisposeAsync(); -} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessageManager.cs b/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessageManager.cs deleted file mode 100644 index eb6b2587a2..0000000000 --- a/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessageManager.cs +++ /dev/null @@ -1,59 +0,0 @@ -namespace ServiceControl.Persistence.RavenDB -{ - using System; - using System.Threading; - using System.Threading.Tasks; - using Raven.Client.Documents.Session; - using ServiceControl.MessageFailures; - using ServiceControl.Persistence.Recoverability.Editing; - - class EditFailedMessageManager : AbstractSessionManager, IEditFailedMessagesManager - { - readonly IAsyncDocumentSession session; - readonly ExpirationManager expirationManager; - FailedMessage failedMessage; - - public EditFailedMessageManager(IAsyncDocumentSession session, ExpirationManager expirationManager) - : base(session) - { - this.session = session; - this.expirationManager = expirationManager; - } - - public async Task GetFailedMessage(string failedMessageId, CancellationToken cancellationToken = default) - { - failedMessage = await session.LoadAsync(FailedMessageIdGenerator.MakeDocumentId(failedMessageId), cancellationToken); - return failedMessage; - } - - public async Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default) - { - var edit = await session.LoadAsync(FailedMessageEdit.MakeDocumentId(failedMessageId), cancellationToken); - return edit?.EditId; - } - - public Task SetCurrentEditingRequestId(string editingMessageId, CancellationToken cancellationToken = default) - { - if (failedMessage == null) - { - throw new InvalidOperationException("No failed message loaded"); - } - return session.StoreAsync(new FailedMessageEdit - { - Id = FailedMessageEdit.MakeDocumentId(failedMessage.UniqueMessageId), - FailedMessageId = failedMessage.Id, - EditId = editingMessageId - }, cancellationToken); - } - - public Task SetFailedMessageAsResolved(CancellationToken cancellationToken = default) - { - // Instance is tracked by the document session - failedMessage.Status = FailedMessageStatus.Resolved; - - expirationManager.EnableExpiration(session, failedMessage); - - return Task.CompletedTask; - } - } -} diff --git a/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs b/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs index 87b54cbc18..50b8a16320 100644 --- a/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs @@ -2,11 +2,68 @@ namespace ServiceControl.Persistence.RavenDB.Editing { using System.Threading; using System.Threading.Tasks; + using Raven.Client.Exceptions; + using ServiceControl.MessageFailures; + using ServiceControl.Persistence.Recoverability.Editing; class EditFailedMessagesDataStore(IRavenSessionProvider sessionProvider, ExpirationManager expirationManager) : IEditFailedMessagesDataStore { - public async Task CreateEditFailedMessageManager(CancellationToken cancellationToken = default) => - // the edit failed message manager manages the lifetime of the session - new EditFailedMessageManager(await sessionProvider.OpenSession(cancellationToken: cancellationToken), expirationManager); + public async Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default) + { + using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken); + var edit = await session.LoadAsync(FailedMessageEdit.MakeDocumentId(failedMessageId), cancellationToken); + return edit?.EditId; + } + + public async Task TryBeginEdit(string failedMessageId, string editingMessageId, CancellationToken cancellationToken = default) + { + using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken); + session.Advanced.UseOptimisticConcurrency = true; + + var failedMessage = await session.LoadAsync(FailedMessageIdGenerator.MakeDocumentId(failedMessageId), cancellationToken); + if (failedMessage is null) + { + return new BeginEditResult(BeginEditOutcome.MessageNotFound); + } + if (failedMessage.Status != FailedMessageStatus.Unresolved) + { + return new BeginEditResult(BeginEditOutcome.MessageNotUnresolved); + } + + var editDocumentId = FailedMessageEdit.MakeDocumentId(failedMessageId); + var existingEdit = await session.LoadAsync(editDocumentId, cancellationToken); + if (existingEdit is null) + { + await session.StoreAsync(new FailedMessageEdit { Id = editDocumentId, FailedMessageId = failedMessage.Id, EditId = editingMessageId }, cancellationToken); + + failedMessage.Status = FailedMessageStatus.Resolved; + expirationManager.EnableExpiration(session, failedMessage); + + try + { + await session.SaveChangesAsync(cancellationToken); + return new(BeginEditOutcome.Acquired, failedMessage); + } + catch (ConcurrencyException) + { + // One bounded reload is sufficient: Raven reports the conflict only after the + // competing atomic batch has won, so the persisted claim identifies the outcome. + existingEdit = await ReloadConflictResult(failedMessageId, cancellationToken); + } + } + + if (existingEdit.EditId == editingMessageId) + { + return new BeginEditResult(BeginEditOutcome.Acquired, failedMessage, existingEdit.EditId); + } + return new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: existingEdit.EditId); + } + + async Task ReloadConflictResult(string failedMessageId, CancellationToken cancellationToken) + { + using var reloadSession = await sessionProvider.OpenSession(cancellationToken: cancellationToken); + return await reloadSession.LoadAsync(FailedMessageEdit.MakeDocumentId(failedMessageId), cancellationToken); + } + } } diff --git a/src/ServiceControl.Persistence.RavenDB/Transactions/RavenTransactionalDataStore.cs b/src/ServiceControl.Persistence.RavenDB/Transactions/RavenTransactionalDataStore.cs deleted file mode 100644 index c7610f5120..0000000000 --- a/src/ServiceControl.Persistence.RavenDB/Transactions/RavenTransactionalDataStore.cs +++ /dev/null @@ -1,18 +0,0 @@ -namespace ServiceControl.Persistence.RavenDB -{ - using System.Threading; - using System.Threading.Tasks; - using Raven.Client.Documents.Session; - - abstract class AbstractSessionManager(IAsyncDocumentSession session) : IDataSessionManager - { - protected IAsyncDocumentSession Session { get; } = session; - - public Task SaveChanges(CancellationToken cancellationToken = default) => Session.SaveChangesAsync(cancellationToken); - public ValueTask DisposeAsync() - { - Session.Dispose(); - return ValueTask.CompletedTask; - } - } -} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/Editing/EditAcquisitionTests.cs b/src/ServiceControl.Persistence.Tests.RavenDB/Editing/EditAcquisitionTests.cs new file mode 100644 index 0000000000..856e7b377e --- /dev/null +++ b/src/ServiceControl.Persistence.Tests.RavenDB/Editing/EditAcquisitionTests.cs @@ -0,0 +1,51 @@ +namespace ServiceControl.Persistence.Tests.RavenDB.Editing; + +using System; +using System.Threading.Tasks; +using NUnit.Framework; +using Raven.Client; +using ServiceControl.Contracts.Operations; +using ServiceControl.MessageFailures; + +class EditAcquisitionTests : PersistenceTestBase +{ + [Test] + public async Task TryBeginEdit_applies_the_failed_message_expiration() + { + var failedMessageId = Guid.NewGuid().ToString(); + var attemptedAt = DateTime.UtcNow; + var failedMessage = await SeedFailedMessage(new FailedMessage + { + UniqueMessageId = failedMessageId, + Status = FailedMessageStatus.Unresolved, + ProcessingAttempts = + [ + new FailedMessage.ProcessingAttempt + { + AttemptedAt = attemptedAt, + MessageId = Guid.NewGuid().ToString(), + Body = "body", + Headers = [], + FailureDetails = new FailureDetails + { + AddressOfFailingEndpoint = "Shipping", + TimeOfFailure = attemptedAt + } + } + ] + }); + + var result = await EditFailedMessagesStore.TryBeginEdit(failedMessageId, Guid.NewGuid().ToString()); + + using var session = PersistenceTestsContext.DocumentStore.OpenAsyncSession(); + var persisted = await session.LoadAsync(failedMessage.Id); + var metadata = session.Advanced.GetMetadataFor(persisted); + + using (Assert.EnterMultipleScope()) + { + Assert.That(result.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); + Assert.That(persisted.Status, Is.EqualTo(FailedMessageStatus.Resolved)); + Assert.That(metadata.ContainsKey(Constants.Documents.Metadata.Expires), Is.True); + } + } +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/EditFailedMessagesManagerTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/AtomicEditRetentionTests.cs similarity index 57% rename from src/ServiceControl.Persistence.Tests/EFCore/EditFailedMessagesManagerTests.cs rename to src/ServiceControl.Persistence.Tests/EFCore/AtomicEditRetentionTests.cs index 52b9ed1d6c..e1f08e8d5a 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/EditFailedMessagesManagerTests.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/AtomicEditRetentionTests.cs @@ -10,20 +10,18 @@ namespace ServiceControl.Persistence.Tests; using ServiceControl.Persistence.EFCore.Entities; /// -/// Focused EFCore-only tests for . These run only against -/// the SQL Server and PostgreSQL persisters (the RavenDB test project excludes the EFCore test -/// folder). They lock in the / -/// requirement that the retention sweeper relies -/// on, and confirm the edit-id round-trips through the database across manager scopes. +/// EF Core-specific retention tests for atomic edit acquisition. These run only against the SQL +/// Server and PostgreSQL persisters and lock in the +/// and behavior used by the retention sweeper. /// -class EditFailedMessagesManagerTests : ErrorIngestionTestBase +class AtomicEditRetentionTests : ErrorIngestionTestBase { IEditFailedMessagesDataStore EditStore => ServiceProvider.GetRequiredService(); DateTime Now => PersistenceTestsContext.FakeTime.GetUtcNow().UtcDateTime; [Test] - public async Task Resolving_via_edit_manager_stamps_StatusChangedAt_and_LastModified() + public async Task TryBeginEdit_stamps_StatusChangedAt_and_LastModified() { // Seed an Unresolved message with an ancient timestamp so that a stale // StatusChangedAt would make it immediately sweepable. @@ -51,18 +49,8 @@ await Store(new FailedMessageEntity var resolveTime = Now; - await using (var manager = await EditStore.CreateEditFailedMessageManager()) - { - var failedMessage = await manager.GetFailedMessage(failedMessageId); - Assert.That(failedMessage, Is.Not.Null); - Assert.That(failedMessage.Status, Is.EqualTo(FailedMessageStatus.Unresolved)); - - Assert.That(await manager.GetCurrentEditingRequestId(failedMessageId), Is.Null); - - await manager.SetCurrentEditingRequestId(editId); - await manager.SetFailedMessageAsResolved(); - await manager.SaveChanges(); - } + var result = await EditStore.TryBeginEdit(failedMessageId, editId); + Assert.That(result.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); var entity = await GetFailedMessage(uniqueMessageId); @@ -74,10 +62,7 @@ await Store(new FailedMessageEntity Assert.That(entity.LastModified, Is.EqualTo(resolveTime), "LastModified must be stamped on resolve"); } - // The edit id round-trips through the database across a brand new manager scope (i.e. it is - // persisted, not held in memory). - await using var assertionManager = await EditStore.CreateEditFailedMessageManager(); - Assert.That(await assertionManager.GetCurrentEditingRequestId(failedMessageId), Is.EqualTo(editId)); + Assert.That(await EditStore.GetCurrentEditingRequestId(failedMessageId), Is.EqualTo(editId)); } [Test] @@ -108,15 +93,10 @@ await Store(new FailedMessageEntity var failedMessageId = uniqueMessageId.ToString(); - // Resolving via the edit manager re-stamps StatusChangedAt to "now", moving the message + // Atomic edit acquisition re-stamps StatusChangedAt to "now", moving the message // back inside the retention window. - await using (var manager = await EditStore.CreateEditFailedMessageManager()) - { - await manager.GetFailedMessage(failedMessageId); - await manager.SetCurrentEditingRequestId(Guid.NewGuid().ToString()); - await manager.SetFailedMessageAsResolved(); - await manager.SaveChanges(); - } + var result = await EditStore.TryBeginEdit(failedMessageId, Guid.NewGuid().ToString()); + Assert.That(result.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); Assert.That(await FindFailedMessage(uniqueMessageId), Is.Not.Null, "the just-resolved message must not be swept immediately"); @@ -129,18 +109,4 @@ await Store(new FailedMessageEntity Assert.That(await FindFailedMessage(uniqueMessageId), Is.Null, "the resolved-via-edit message must be swept after the retention period"); } - [Test] - public async Task GetFailedMessage_returns_null_when_not_found() - { - await using var manager = await EditStore.CreateEditFailedMessageManager(); - Assert.That(await manager.GetFailedMessage(Guid.NewGuid().ToString()), Is.Null); - } - - [Test] - public async Task SetCurrentEditingRequestId_and_SetFailedMessageAsResolved_throw_when_no_message_loaded() - { - await using var manager = await EditStore.CreateEditFailedMessageManager(); - Assert.ThrowsAsync(() => manager.SetCurrentEditingRequestId(Guid.NewGuid().ToString())); - Assert.ThrowsAsync(() => manager.SetFailedMessageAsResolved()); - } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.Tests/EditFailedMessagesDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/EditFailedMessagesDataStoreTests.cs new file mode 100644 index 0000000000..73e00c4c9e --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EditFailedMessagesDataStoreTests.cs @@ -0,0 +1,182 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Contracts.Operations; +using MessageFailures; +using NUnit.Framework; + +class EditFailedMessagesDataStoreTests : PersistenceTestBase +{ + [Test] + public async Task TryBeginEdit_acquires_an_unresolved_message() + { + var failedMessage = await CreateUnresolvedFailedMessage(); + var editId = Guid.NewGuid().ToString(); + + var result = await EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, editId, TestContext.CurrentContext.CancellationToken); + + var persisted = await FailedMessageQueryStore.GetFailedMessage(failedMessage.UniqueMessageId, TestContext.CurrentContext.CancellationToken); + using (Assert.EnterMultipleScope()) + { + Assert.That(result.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); + Assert.That(result.FailedMessage, Is.Not.Null); + Assert.That(result.FailedMessage!.UniqueMessageId, Is.EqualTo(failedMessage.UniqueMessageId)); + Assert.That(result.FailedMessage.ProcessingAttempts.Single().MessageId, Is.EqualTo(failedMessage.ProcessingAttempts.Single().MessageId)); + Assert.That(result.ExistingEditId, Is.Null); + Assert.That(await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessage.UniqueMessageId), Is.EqualTo(editId)); + Assert.That(persisted!.Status, Is.EqualTo(FailedMessageStatus.Resolved)); + } + } + + [Test] + [Repeat(5)] + public async Task TryBeginEdit_is_idempotent_for_the_same_edit_id() + { + var failedMessage = await CreateUnresolvedFailedMessage(); + var editId = Guid.NewGuid().ToString(); + + var first = await EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, editId); + var second = await EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, editId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(first.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); + Assert.That(second.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); + Assert.That(second.FailedMessage, Is.Not.Null, "the dispatch snapshot is required when an edit message is retried"); + Assert.That(second.FailedMessage!.ProcessingAttempts.Single().MessageId, Is.EqualTo(failedMessage.ProcessingAttempts.Single().MessageId)); + Assert.That(second.ExistingEditId, Is.EqualTo(editId)); + Assert.That(await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessage.UniqueMessageId), Is.EqualTo(editId)); + } + } + + [Test] + public async Task TryBeginEdit_reports_the_existing_edit_id_without_mutating_the_message() + { + var failedMessage = await CreateUnresolvedFailedMessage(); + var winningEditId = Guid.NewGuid().ToString(); + var losingEditId = Guid.NewGuid().ToString(); + await EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, winningEditId); + + var result = await EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, losingEditId); + var persisted = await FailedMessageQueryStore.GetFailedMessage(failedMessage.UniqueMessageId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(result.Outcome, Is.EqualTo(BeginEditOutcome.AcquiredByAnotherEdit)); + Assert.That(result.FailedMessage, Is.Null); + Assert.That(result.ExistingEditId, Is.EqualTo(winningEditId)); + Assert.That(await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessage.UniqueMessageId), Is.EqualTo(winningEditId)); + Assert.That(persisted!.Status, Is.EqualTo(FailedMessageStatus.Resolved)); + } + } + + [Test] + [Repeat(5)] + public async Task Concurrent_different_edit_ids_produce_exactly_one_winner() + { + var failedMessage = await CreateUnresolvedFailedMessage(); + var editIds = new[] { Guid.NewGuid().ToString(), Guid.NewGuid().ToString() }; + + var results = await Task.WhenAll(editIds.Select(editId => + EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, editId, TestContext.CurrentContext.CancellationToken))); + + var acquired = results.Single(result => result.Outcome == BeginEditOutcome.Acquired); + var loser = results.Single(result => result.Outcome == BeginEditOutcome.AcquiredByAnotherEdit); + var persistedEditId = await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessage.UniqueMessageId); + var persisted = await FailedMessageQueryStore.GetFailedMessage(failedMessage.UniqueMessageId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(acquired.FailedMessage, Is.Not.Null); + Assert.That(loser.FailedMessage, Is.Null); + Assert.That(loser.ExistingEditId, Is.EqualTo(persistedEditId)); + Assert.That(editIds, Does.Contain(persistedEditId)); + Assert.That(persisted!.Status, Is.EqualTo(FailedMessageStatus.Resolved)); + } + } + + [Test] + public async Task TryBeginEdit_returns_MessageNotFound_for_a_missing_message() + { + var failedMessageId = Guid.NewGuid().ToString(); + + var result = await EditFailedMessagesStore.TryBeginEdit(failedMessageId, Guid.NewGuid().ToString()); + + using (Assert.EnterMultipleScope()) + { + Assert.That(result, Is.EqualTo(new BeginEditResult(BeginEditOutcome.MessageNotFound))); + Assert.That(await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessageId), Is.Null); + } + } + + [TestCase(FailedMessageStatus.Resolved)] + [TestCase(FailedMessageStatus.RetryIssued)] + [TestCase(FailedMessageStatus.Archived)] + public async Task TryBeginEdit_returns_MessageNotUnresolved_without_creating_a_claim(FailedMessageStatus status) + { + var failedMessage = await CreateFailedMessage(status); + + var result = await EditFailedMessagesStore.TryBeginEdit(failedMessage.UniqueMessageId, Guid.NewGuid().ToString()); + var persisted = await FailedMessageQueryStore.GetFailedMessage(failedMessage.UniqueMessageId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(result, Is.EqualTo(new BeginEditResult(BeginEditOutcome.MessageNotUnresolved))); + Assert.That(await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessage.UniqueMessageId), Is.Null); + Assert.That(persisted!.Status, Is.EqualTo(status)); + } + } + + [Test] + public void GetCurrentEditingRequestId_propagates_cancellation() + { + using var cancellation = new CancellationTokenSource(); + cancellation.Cancel(); + + Assert.That( + async () => await EditFailedMessagesStore.GetCurrentEditingRequestId(Guid.NewGuid().ToString(), cancellation.Token), + Throws.InstanceOf()); + } + + [Test] + public void TryBeginEdit_propagates_cancellation() + { + using var cancellation = new CancellationTokenSource(); + cancellation.Cancel(); + + Assert.That( + async () => await EditFailedMessagesStore.TryBeginEdit(Guid.NewGuid().ToString(), Guid.NewGuid().ToString(), cancellation.Token), + Throws.InstanceOf()); + } + + Task CreateUnresolvedFailedMessage() => CreateFailedMessage(FailedMessageStatus.Unresolved); + + async Task CreateFailedMessage(FailedMessageStatus status) + { + var failedMessageId = Guid.NewGuid().ToString(); + var attemptedAt = DateTime.UtcNow; + return await SeedFailedMessage(new FailedMessage + { + UniqueMessageId = failedMessageId, + Status = status, + ProcessingAttempts = + [ + new FailedMessage.ProcessingAttempt + { + AttemptedAt = attemptedAt, + MessageId = Guid.NewGuid().ToString(), + Body = "body", + Headers = [], + FailureDetails = new FailureDetails + { + AddressOfFailingEndpoint = "Shipping", + TimeOfFailure = attemptedAt + } + } + ] + }); + } +} diff --git a/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs b/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs index 0e684356a8..d32c6db6a2 100644 --- a/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs +++ b/src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs @@ -76,6 +76,13 @@ public async Task TearDown() protected Task CompleteDatabaseOperation() => PersistenceTestsContext.CompleteDatabaseOperation(); + protected async Task SeedFailedMessage(FailedMessage failedMessage) + { + failedMessage.Id = PersistenceTestsContext.GenerateFailedMessageRecordId(failedMessage.UniqueMessageId); + await PersistenceTestsContext.InsertFailedMessages(failedMessage); + return failedMessage; + } + protected static async Task WaitUntil(Func> conditionChecker, string condition, TimeSpan timeout = default) { timeout = timeout == default ? TimeSpan.FromSeconds(10) : timeout; diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs index afcb089bb4..9e8cc10d7f 100644 --- a/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs +++ b/src/ServiceControl.Persistence.Tests/Recoverability/EditMessageTests.cs @@ -57,8 +57,7 @@ public async Task Should_discard_edit_if_edited_message_not_unresolved(FailedMes var failedMessage = await FailedMessageQueryStore.GetFailedMessage(failedMessageId); - var editFailedMessagesManager = await EditFailedMessagesStore.CreateEditFailedMessageManager(); - var editOperation = await editFailedMessagesManager.GetCurrentEditingRequestId(failedMessageId); + var editOperation = await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessageId); using (Assert.EnterMultipleScope()) { @@ -76,28 +75,21 @@ public async Task Should_discard_edit_when_different_edit_already_exists() _ = await CreateAndStoreFailedMessage(failedMessageId); - await using (var editFailedMessagesManager = await EditFailedMessagesStore.CreateEditFailedMessageManager()) - { - _ = await editFailedMessagesManager.GetFailedMessage(failedMessageId); - await editFailedMessagesManager.SetCurrentEditingRequestId(previousEdit); - await editFailedMessagesManager.SaveChanges(); - } + var previousAcquisition = await EditFailedMessagesStore.TryBeginEdit(failedMessageId, previousEdit); + Assert.That(previousAcquisition.Outcome, Is.EqualTo(BeginEditOutcome.Acquired)); var message = CreateEditMessage(failedMessageId); // Act await handler.Handle(message, new TestableMessageHandlerContext()); - await using (var editFailedMessagesManagerAssert = await EditFailedMessagesStore.CreateEditFailedMessageManager()) + var failedMessage = await FailedMessageQueryStore.GetFailedMessage(failedMessageId); + var editId = await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessageId); + + using (Assert.EnterMultipleScope()) { - var failedMessage = await editFailedMessagesManagerAssert.GetFailedMessage(failedMessageId); - var editId = await editFailedMessagesManagerAssert.GetCurrentEditingRequestId(failedMessageId); - - using (Assert.EnterMultipleScope()) - { - Assert.That(editId, Is.EqualTo(previousEdit)); - Assert.That(failedMessage.Status, Is.EqualTo(FailedMessageStatus.Unresolved)); - } + Assert.That(editId, Is.EqualTo(previousEdit)); + Assert.That(failedMessage.Status, Is.EqualTo(FailedMessageStatus.Resolved)); } Assert.That(dispatcher.DispatchedMessages, Is.Empty); @@ -125,11 +117,10 @@ public async Task Should_dispatch_edited_message_when_first_edit() Assert.That(dispatchedMessage.Item1.Message.Headers["someKey"], Is.EqualTo("someValue")); } - await using var x = await EditFailedMessagesStore.CreateEditFailedMessageManager(); - var failedMessage2 = await x.GetFailedMessage(failedMessage.UniqueMessageId); + var failedMessage2 = await FailedMessageQueryStore.GetFailedMessage(failedMessage.UniqueMessageId); Assert.That(failedMessage2, Is.Not.Null, "Edited failed message"); - var editId = await x.GetCurrentEditingRequestId(failedMessage2.UniqueMessageId); + var editId = await EditFailedMessagesStore.GetCurrentEditingRequestId(failedMessage2.UniqueMessageId); using (Assert.EnterMultipleScope()) { diff --git a/src/ServiceControl.Persistence/IDataSessionManager.cs b/src/ServiceControl.Persistence/IDataSessionManager.cs deleted file mode 100644 index 72a91d99b3..0000000000 --- a/src/ServiceControl.Persistence/IDataSessionManager.cs +++ /dev/null @@ -1,11 +0,0 @@ -namespace ServiceControl.Persistence -{ - using System; - using System.Threading; - using System.Threading.Tasks; - - public interface IDataSessionManager : IAsyncDisposable - { - Task SaveChanges(CancellationToken cancellationToken = default); - } -} diff --git a/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs b/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs index 9c8293d8fb..11acb41999 100644 --- a/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs +++ b/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs @@ -1,10 +1,44 @@ -namespace ServiceControl.Persistence +#nullable enable +namespace ServiceControl.Persistence; + +using System.Threading; +using System.Threading.Tasks; +using ServiceControl.MessageFailures; + +public interface IEditFailedMessagesDataStore { - using System.Threading; - using System.Threading.Tasks; + /// + /// Gets the edit request that currently owns the failed message, if any. + /// This query is advisory; callers must use to acquire an edit safely. + /// + Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default); + + /// + /// Atomically claims an unresolved failed message and marks the original failure resolved. + /// An acquired result is committed before this method returns. Non-acquired results do not + /// modify the failed message. + /// + /// + /// returns the failed-message snapshot and no existing edit ID. + /// returns the failed-message snapshot and the supplied edit ID. + /// returns the winning edit ID and no failed-message snapshot. + /// and return neither optional value. + /// + Task TryBeginEdit(string failedMessageId, string editingMessageId, CancellationToken cancellationToken = default); +} - public interface IEditFailedMessagesDataStore - { - Task CreateEditFailedMessageManager(CancellationToken cancellationToken = default); - } +public enum BeginEditOutcome +{ + Acquired, + AcquiredByAnotherEdit, + MessageNotFound, + MessageNotUnresolved } + +/// +/// The result of attempting to acquire a failed message for editing. +/// +/// The acquisition outcome. +/// The snapshot used to dispatch an acquired or idempotently reacquired edit. +/// The existing claim for an idempotent retry or competing edit. +public sealed record BeginEditResult(BeginEditOutcome Outcome, FailedMessage? FailedMessage = null, string? ExistingEditId = null); diff --git a/src/ServiceControl.Persistence/IEditFailedMessagesManager.cs b/src/ServiceControl.Persistence/IEditFailedMessagesManager.cs deleted file mode 100644 index aa99151262..0000000000 --- a/src/ServiceControl.Persistence/IEditFailedMessagesManager.cs +++ /dev/null @@ -1,15 +0,0 @@ -#nullable enable -namespace ServiceControl.Persistence -{ - using System.Threading; - using System.Threading.Tasks; - using ServiceControl.MessageFailures; - - public interface IEditFailedMessagesManager : IDataSessionManager - { - Task GetFailedMessage(string failedMessageId, CancellationToken cancellationToken = default); - Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default); - Task SetCurrentEditingRequestId(string editingMessageId, CancellationToken cancellationToken = default); - Task SetFailedMessageAsResolved(CancellationToken cancellationToken = default); - } -} diff --git a/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs b/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs index 434a1b83c1..85db4e3aff 100644 --- a/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs +++ b/src/ServiceControl.UnitTests/MessageFailures/EditFailedMessagesControllerAuditTests.cs @@ -45,25 +45,17 @@ public async Task Edit_emits_single_operation() Assert.That(op.Count, Is.EqualTo(1)); } - sealed class FakeEditFailedMessagesManager : IEditFailedMessagesManager + sealed class StubErrorMessageDataStore : IFailedMessageQueryDataStore, IEditFailedMessagesDataStore { + public FailedMessage? ErrorByResult { get; set; } public string? CurrentEditingRequestId { get; set; } - public Task SaveChanges(CancellationToken cancellationToken = default) => Task.CompletedTask; - public Task GetFailedMessage(string failedMessageId, CancellationToken cancellationToken = default) => Task.FromResult(null); - public Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default) => Task.FromResult(CurrentEditingRequestId); - public Task SetCurrentEditingRequestId(string editingMessageId, CancellationToken cancellationToken = default) => Task.CompletedTask; - public Task SetFailedMessageAsResolved(CancellationToken cancellationToken = default) => Task.CompletedTask; - - public ValueTask DisposeAsync() => ValueTask.CompletedTask; - } + public Task GetCurrentEditingRequestId(string failedMessageId, CancellationToken cancellationToken = default) => + Task.FromResult(CurrentEditingRequestId); - sealed class StubErrorMessageDataStore : IFailedMessageQueryDataStore, IEditFailedMessagesDataStore - { - public FailedMessage? ErrorByResult { get; set; } - public FakeEditFailedMessagesManager EditManager { get; } = new(); + public Task TryBeginEdit(string failedMessageId, string editingMessageId, CancellationToken cancellationToken = default) => + Task.FromResult(new BeginEditResult(BeginEditOutcome.MessageNotFound)); - public Task CreateEditFailedMessageManager(CancellationToken cancellationToken = default) => Task.FromResult(EditManager); public Task GetFailedMessage(string failedMessageId, CancellationToken cancellationToken = default) => Task.FromResult(ErrorByResult); public Task GetFailedMessagesByIds(Guid[] ids, CancellationToken cancellationToken = default) => throw new NotImplementedException(); diff --git a/src/ServiceControl/MessageFailures/Api/EditFailedMessagesController.cs b/src/ServiceControl/MessageFailures/Api/EditFailedMessagesController.cs index 4b8913c345..fb794d2ca6 100644 --- a/src/ServiceControl/MessageFailures/Api/EditFailedMessagesController.cs +++ b/src/ServiceControl/MessageFailures/Api/EditFailedMessagesController.cs @@ -44,9 +44,8 @@ public async Task> Edit(string failedMessageId, return NotFound(); } - //HINT: This validation is the first one because we want to minimize the chance of two users concurrently execute an edit-retry. - var editManager = await editStore.CreateEditFailedMessageManager(cancellationToken); - var editId = await editManager.GetCurrentEditingRequestId(failedMessageId, cancellationToken); + // This early query is advisory only. The handler atomically acquires the edit before dispatch. + var editId = await editStore.GetCurrentEditingRequestId(failedMessageId, cancellationToken); if (editId != null) { logger.LogWarning("Cannot edit message {FailedMessageId} because it has already been edited", failedMessageId); diff --git a/src/ServiceControl/Recoverability/Editing/EditHandler.cs b/src/ServiceControl/Recoverability/Editing/EditHandler.cs index 9b7f68dc0b..67d2ddba55 100644 --- a/src/ServiceControl/Recoverability/Editing/EditHandler.cs +++ b/src/ServiceControl/Recoverability/Editing/EditHandler.cs @@ -22,43 +22,30 @@ class EditHandler(IEditFailedMessagesDataStore store, IMessageRedirectsDataStore { public async Task Handle(EditAndSend message, IMessageHandlerContext context) { - FailedMessage failedMessage; - string editId; - await using (var session = await store.CreateEditFailedMessageManager(context.CancellationToken)) - { - failedMessage = await session.GetFailedMessage(message.FailedMessageId, context.CancellationToken); + var beginEdit = await store.TryBeginEdit(message.FailedMessageId, context.MessageId, context.CancellationToken); - if (failedMessage == null) - { + switch (beginEdit.Outcome) + { + case BeginEditOutcome.MessageNotFound: logger.LogWarning("Discarding edit {MessageId} because no message failure for id {FailedMessageId} has been found", context.MessageId, message.FailedMessageId); return; - } - - editId = await session.GetCurrentEditingRequestId(message.FailedMessageId, context.CancellationToken); - if (editId == null) - { - if (failedMessage.Status != FailedMessageStatus.Unresolved) - { - logger.LogWarning("Discarding edit {MessageId} because message failure {FailedMessageId} doesn't have state 'Unresolved'", context.MessageId, message.FailedMessageId); - return; - } - - // create a retries document to prevent concurrent edits - await session.SetCurrentEditingRequestId(context.MessageId, context.CancellationToken); - } - else if (editId != context.MessageId) - { - logger.LogWarning("Discarding edit & retry request because the failed message id {FailedMessageId} has already been edited by Message ID {EditedMessageId}", message.FailedMessageId, editId); + case BeginEditOutcome.MessageNotUnresolved: + logger.LogWarning("Discarding edit {MessageId} because message failure {FailedMessageId} doesn't have state 'Unresolved'", context.MessageId, message.FailedMessageId); return; - } - - // the original failure is marked as resolved as any failures of the edited message are treated as a new message failure. - await session.SetFailedMessageAsResolved(context.CancellationToken); - - - await session.SaveChanges(context.CancellationToken); + case BeginEditOutcome.AcquiredByAnotherEdit: + logger.LogWarning("Discarding edit & retry request because the failed message id {FailedMessageId} has already been edited by Message ID {EditedMessageId}", message.FailedMessageId, beginEdit.ExistingEditId); + return; + case BeginEditOutcome.Acquired: + break; + default: + throw new InvalidOperationException($"Unknown begin-edit outcome: {beginEdit.Outcome}"); } + // The store commits resolution of the original failure before returning. Any failure + // of the edited message is therefore treated as a new failure, preserving the existing + // resolve-before-dispatch behavior. + var failedMessage = beginEdit.FailedMessage!; + var redirects = await redirectsStore.GetRedirects(context.CancellationToken); var attempt = failedMessage.ProcessingAttempts.Last(); From 1a50a930cbc06affbddcab4928913161b78cdd67 Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Thu, 13 Aug 2026 16:21:54 +0800 Subject: [PATCH 2/3] Give test back a more meaningful name --- ...ionTests.cs => EditFailedMessagesDataStoreRetentionTests.cs} | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) rename src/ServiceControl.Persistence.Tests/EFCore/{AtomicEditRetentionTests.cs => EditFailedMessagesDataStoreRetentionTests.cs} (98%) diff --git a/src/ServiceControl.Persistence.Tests/EFCore/AtomicEditRetentionTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/EditFailedMessagesDataStoreRetentionTests.cs similarity index 98% rename from src/ServiceControl.Persistence.Tests/EFCore/AtomicEditRetentionTests.cs rename to src/ServiceControl.Persistence.Tests/EFCore/EditFailedMessagesDataStoreRetentionTests.cs index e1f08e8d5a..8ec279cb72 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/AtomicEditRetentionTests.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/EditFailedMessagesDataStoreRetentionTests.cs @@ -14,7 +14,7 @@ namespace ServiceControl.Persistence.Tests; /// Server and PostgreSQL persisters and lock in the /// and behavior used by the retention sweeper. /// -class AtomicEditRetentionTests : ErrorIngestionTestBase +class EditFailedMessagesDataStoreRetentionTests : ErrorIngestionTestBase { IEditFailedMessagesDataStore EditStore => ServiceProvider.GetRequiredService(); From 2817b83e4171f03f6c28222dbedeea28b7498dbb Mon Sep 17 00:00:00 2001 From: Rhys Bevilaqua Date: Fri, 14 Aug 2026 16:14:40 +0800 Subject: [PATCH 3/3] Fix concurrency logic --- .../EditFailedMessagesDataStore.cs | 78 ++++++++++--------- .../Editing/EditFailedMessagesDataStore.cs | 46 +++++------ .../IEditFailedMessagesDataStore.cs | 3 +- 3 files changed, 65 insertions(+), 62 deletions(-) diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs index b12f9ac729..b77aa9630f 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/EditFailedMessagesDataStore.cs @@ -39,59 +39,61 @@ public Task TryBeginEdit(string failedMessageId, string editing return new BeginEditResult(BeginEditOutcome.MessageNotFound); } + var existingEditId = await dbContext.FailedMessageEdits + .Where(edit => edit.UniqueMessageId == uniqueMessageId) + .Select(edit => edit.EditId) + .SingleOrDefaultAsync(ct); + + if (existingEditId is not null) + { + return existingEditId == editingMessageId + ? new BeginEditResult(BeginEditOutcome.Acquired, entity.ToFailedMessage([]), existingEditId) + : new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: existingEditId); + } + if (entity.Status != FailedMessageStatus.Unresolved) { return new BeginEditResult(BeginEditOutcome.MessageNotUnresolved); } - var editId = await dbContext.FailedMessageEdits - .Where(edit => edit.UniqueMessageId == uniqueMessageId) - .Select(edit => edit.EditId) - .SingleOrDefaultAsync(ct); + dbContext.FailedMessageEdits.Add(new FailedMessageEditEntity { UniqueMessageId = uniqueMessageId, EditId = editingMessageId }); + + var now = timeProvider.GetUtcNow().UtcDateTime; + entity.Status = FailedMessageStatus.Resolved; + entity.StatusChangedAt = now; + entity.LastModified = now; - if (editId is null) + try { - editId = editingMessageId; - dbContext.FailedMessageEdits.Add(new FailedMessageEditEntity { UniqueMessageId = uniqueMessageId, EditId = editingMessageId }); + await dbContext.SaveChangesAsync(ct); + return new BeginEditResult(BeginEditOutcome.Acquired, entity.ToFailedMessage([])); + } + catch (DbUpdateException exception) when (dbContext.IsDuplicateKeyException(exception)) + { + dbContext.ChangeTracker.Clear(); - var now = timeProvider.GetUtcNow().UtcDateTime; - entity.Status = FailedMessageStatus.Resolved; - entity.StatusChangedAt = now; - entity.LastModified = now; + var winningEditId = await dbContext.FailedMessageEdits + .AsNoTracking() + .Where(edit => edit.UniqueMessageId == uniqueMessageId) + .Select(edit => edit.EditId) + .SingleOrDefaultAsync(ct); - try + if (winningEditId is null) { - await dbContext.SaveChangesAsync(ct); + throw; } - catch (DbUpdateException exception) when (dbContext.IsDuplicateKeyException(exception)) - { - dbContext.ChangeTracker.Clear(); - - var winningEditId = await dbContext.FailedMessageEdits - .AsNoTracking() - .Where(edit => edit.UniqueMessageId == uniqueMessageId) - .Select(edit => edit.EditId) - .SingleOrDefaultAsync(ct); - - if (winningEditId is null) - { - throw; - } - - entity = await dbContext.FailedMessages - .AsNoTracking() - .SingleAsync(message => message.UniqueMessageId == uniqueMessageId, ct); - editId = winningEditId; + if (winningEditId != editingMessageId) + { + return new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: winningEditId); } - } - if (editId == editingMessageId) - { - return new BeginEditResult(BeginEditOutcome.Acquired, entity.ToFailedMessage([]), editId); - } + entity = await dbContext.FailedMessages + .AsNoTracking() + .SingleAsync(message => message.UniqueMessageId == uniqueMessageId, ct); - return new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: editId); + return new BeginEditResult(BeginEditOutcome.Acquired, entity.ToFailedMessage([]), winningEditId); + } }, cancellationToken); } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs b/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs index 50b8a16320..55cd75225d 100644 --- a/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/Editing/EditFailedMessagesDataStore.cs @@ -25,38 +25,40 @@ public async Task TryBeginEdit(string failedMessageId, string e { return new BeginEditResult(BeginEditOutcome.MessageNotFound); } + var editDocumentId = FailedMessageEdit.MakeDocumentId(failedMessageId); + var existingEdit = await session.LoadAsync(editDocumentId, cancellationToken); + if (existingEdit is not null) + { + return existingEdit.EditId == editingMessageId + ? new BeginEditResult(BeginEditOutcome.Acquired, failedMessage, existingEdit.EditId) + : new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: existingEdit.EditId); + } + if (failedMessage.Status != FailedMessageStatus.Unresolved) { return new BeginEditResult(BeginEditOutcome.MessageNotUnresolved); } - var editDocumentId = FailedMessageEdit.MakeDocumentId(failedMessageId); - var existingEdit = await session.LoadAsync(editDocumentId, cancellationToken); - if (existingEdit is null) - { - await session.StoreAsync(new FailedMessageEdit { Id = editDocumentId, FailedMessageId = failedMessage.Id, EditId = editingMessageId }, cancellationToken); + await session.StoreAsync(new FailedMessageEdit { Id = editDocumentId, FailedMessageId = failedMessage.Id, EditId = editingMessageId }, cancellationToken); - failedMessage.Status = FailedMessageStatus.Resolved; - expirationManager.EnableExpiration(session, failedMessage); + failedMessage.Status = FailedMessageStatus.Resolved; + expirationManager.EnableExpiration(session, failedMessage); - try - { - await session.SaveChangesAsync(cancellationToken); - return new(BeginEditOutcome.Acquired, failedMessage); - } - catch (ConcurrencyException) - { - // One bounded reload is sufficient: Raven reports the conflict only after the - // competing atomic batch has won, so the persisted claim identifies the outcome. - existingEdit = await ReloadConflictResult(failedMessageId, cancellationToken); - } + try + { + await session.SaveChangesAsync(cancellationToken); + return new(BeginEditOutcome.Acquired, failedMessage); } - - if (existingEdit.EditId == editingMessageId) + catch (ConcurrencyException) { - return new BeginEditResult(BeginEditOutcome.Acquired, failedMessage, existingEdit.EditId); + // One bounded reload is sufficient: Raven reports the conflict only after the + // competing atomic batch has won, so the persisted claim identifies the outcome. + existingEdit = await ReloadConflictResult(failedMessageId, cancellationToken); } - return new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: existingEdit.EditId); + + return existingEdit.EditId == editingMessageId + ? new BeginEditResult(BeginEditOutcome.Acquired, failedMessage, existingEdit.EditId) + : new BeginEditResult(BeginEditOutcome.AcquiredByAnotherEdit, ExistingEditId: existingEdit.EditId); } async Task ReloadConflictResult(string failedMessageId, CancellationToken cancellationToken) diff --git a/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs b/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs index 11acb41999..252ec4cbb7 100644 --- a/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs +++ b/src/ServiceControl.Persistence/IEditFailedMessagesDataStore.cs @@ -19,8 +19,7 @@ public interface IEditFailedMessagesDataStore /// modify the failed message. /// /// - /// returns the failed-message snapshot and no existing edit ID. - /// returns the failed-message snapshot and the supplied edit ID. + /// returns the failed-message snapshot. A new acquisition has no existing edit ID; an idempotent retry returns the supplied edit ID. /// returns the winning edit ID and no failed-message snapshot. /// and return neither optional value. ///