diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/GroupsDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/GroupsDataStore.cs index 56760ece42..b967aed312 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/GroupsDataStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/GroupsDataStore.cs @@ -12,7 +12,7 @@ namespace ServiceControl.Persistence.EFCore.Implementation; public class GroupsDataStore(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), IGroupsDataStore { - public Task> GetFailureGroupsByClassifier(string classifier, string classifierFilter) => + public Task> GetUnresolvedGroupsByClassifier(string classifier, string classifierFilter) => ExecuteWithDbContext(dbContext => { var groups = ByClassifier(dbContext, classifier); @@ -25,30 +25,15 @@ public Task> GetFailureGroupsByClassifier(string classif return MostRecent(groups.AggregateGroups(WithStatus(dbContext, FailedMessageStatus.Unresolved))); }); - public Task> GetArchivedFailureGroupsByClassifier(string classifier) => + public Task> GetArchivedGroupsByClassifier(string classifier) => ExecuteWithDbContext(dbContext => MostRecent( ByClassifier(dbContext, classifier).AggregateGroups(WithStatus(dbContext, FailedMessageStatus.Archived)))); - // Implemented once retry batches are persisted, together with IRetryDocumentDataStore. - public Task GetCurrentForwardingBatch() => - throw new NotImplementedException(); - - public Task>> GetGroup(string groupId, string status, string modified) => - ExecuteWithDbContext(async dbContext => - { - var groups = await ById(dbContext, groupId, FailedMessageStatus.Unresolved, status, modified).ToListAsync(); + public Task> GetUnresolvedGroup(string groupId, string status, string modified) => + ExecuteWithDbContext(dbContext => SingleGroup(dbContext, groupId, FailedMessageStatus.Unresolved, status, modified)); - return new QueryResult>(groups, groups.ToQueryStatsInfo()); - }); - - public Task> GetFailureGroupView(string groupId, string status, string modified) => - ExecuteWithDbContext(async dbContext => - { - var groups = await ById(dbContext, groupId, FailedMessageStatus.Archived, status, modified).ToListAsync(); - - // A missing group is reported as a null result, the same as the RavenDB persister does. - return new QueryResult(groups.FirstOrDefault()!, groups.ToQueryStatsInfo()); - }); + public Task> GetArchivedGroup(string groupId, string status, string modified) => + ExecuteWithDbContext(dbContext => SingleGroup(dbContext, groupId, FailedMessageStatus.Archived, status, modified)); public Task>> GetGroupErrors(string groupId, string status, string modified, SortInfo sortInfo, PagingInfo pagingInfo) => ExecuteWithDbContext(dbContext => InGroup(dbContext, groupId, status, modified).ToPagedResult(pagingInfo, sortInfo)); @@ -67,18 +52,18 @@ static IQueryable ByClassifier(ServiceControlDbContext .AsNoTracking() .Where(group => group.Type == classifier); - /// - /// The status a group is read at, before the caller's own status and modified filters narrow it - /// further. RavenDB reads open groups out of an unresolved-only index and archived groups out of - /// an archived-only one, which is what stands in for here. - /// - static IQueryable ById(ServiceControlDbContext dbContext, string groupId, FailedMessageStatus baseline, string status, string modified) => - dbContext.FailedMessageGroups + static async Task> SingleGroup(ServiceControlDbContext dbContext, string groupId, FailedMessageStatus baseline, string status, string modified) + { + var groups = await dbContext.FailedMessageGroups .AsNoTracking() .Where(group => group.GroupId == groupId) .AggregateGroups(WithStatus(dbContext, baseline) .FilterByStatus(status) - .FilterByLastModifiedRange(modified)); + .FilterByLastModifiedRange(modified)) + .ToListAsync(); + + return new QueryResult(groups.FirstOrDefault()!, groups.ToQueryStatsInfo()); + } static IQueryable WithStatus(ServiceControlDbContext dbContext, FailedMessageStatus status) => dbContext.FailedMessages diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/RetryDocumentDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/RetryDocumentDataStore.cs index f29758b8a5..af757abef0 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/RetryDocumentDataStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/RetryDocumentDataStore.cs @@ -36,6 +36,6 @@ public Task GetBatchesForFailedQueueAddress(DateTime cutoff, string failedQueueA public Task GetBatchesForFailureGroup(string groupId, string groupTitle, string groupType, DateTime cutoff, Func callback) => throw new NotImplementedException(); - public Task QueryFailureGroupViewOnGroupId(string groupId) => + public Task GetCurrentForwardingBatch() => throw new NotImplementedException(); } diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/GroupsDataStore.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/GroupsDataStore.cs index d517b37a83..d12249702e 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/GroupsDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/GroupsDataStore.cs @@ -14,7 +14,7 @@ namespace ServiceControl.Persistence.RavenDB.Recoverability class GroupsDataStore(IRavenSessionProvider sessionProvider) : IGroupsDataStore { - public async Task> GetFailureGroupsByClassifier(string classifier, string classifierFilter) + public async Task> GetUnresolvedGroupsByClassifier(string classifier, string classifierFilter) { using var session = await sessionProvider.OpenSession(); var query = Queryable.Where(session.Query(), v => v.Type == classifier); @@ -40,7 +40,7 @@ public async Task> GetFailureGroupsByClassifier(string c return groups; } - public async Task> GetArchivedFailureGroupsByClassifier(string classifier) + public async Task> GetArchivedGroupsByClassifier(string classifier) { using var session = await sessionProvider.OpenSession(); var groups = session @@ -55,30 +55,21 @@ public async Task> GetArchivedFailureGroupsByClassifier( return results; } - public async Task GetCurrentForwardingBatch() + public async Task> GetUnresolvedGroup(string groupId, string status, string modified) { using var session = await sessionProvider.OpenSession(); - var nowForwarding = await session.Include(r => r.RetryBatchId) - .LoadAsync(RetryDocumentDataStore.NowForwardingDocumentId); - - return nowForwarding == null ? null : await session.LoadAsync(nowForwarding.RetryBatchId); - } - - public async Task>> GetGroup(string groupId, string status, string modified) - { - using var session = await sessionProvider.OpenSession(); - var queryResult = await session.Advanced + var document = await session.Advanced .AsyncDocumentQuery() .Statistics(out var stats) .WhereEquals(group => group.Id, groupId) .FilterByStatusWhere(status) .FilterByLastModifiedRange(modified) - .ToListAsync(); + .FirstOrDefaultAsync(); - return queryResult.ToQueryResult(stats); + return new QueryResult(document, stats.ToQueryStatsInfo()); } - public async Task> GetFailureGroupView(string groupId, string status, string modified) + public async Task> GetArchivedGroup(string groupId, string status, string modified) { using var session = await sessionProvider.OpenSession(); var document = await session.Advanced diff --git a/src/ServiceControl.Persistence.RavenDB/RetryDocumentDataStore.cs b/src/ServiceControl.Persistence.RavenDB/RetryDocumentDataStore.cs index b2606ece2a..3abc04edda 100644 --- a/src/ServiceControl.Persistence.RavenDB/RetryDocumentDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/RetryDocumentDataStore.cs @@ -200,12 +200,13 @@ public async Task GetBatchesForFailureGroup(string groupId, string groupTitle, s } } - public async Task QueryFailureGroupViewOnGroupId(string groupId) + public async Task GetCurrentForwardingBatch() { using var session = await sessionProvider.OpenSession(); - var group = await session.Query() - .FirstOrDefaultAsync(x => x.Id == groupId); - return group; + var nowForwarding = await session.Include(r => r.RetryBatchId) + .LoadAsync(NowForwardingDocumentId); + + return nowForwarding == null ? null : await session.LoadAsync(nowForwarding.RetryBatchId); } public static string MakeDocumentId(string messageUniqueId) => "RetryBatches/" + messageUniqueId; diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/GroupsDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/GroupsDataStoreTests.cs index 4f8654f193..d3a42859f7 100644 --- a/src/ServiceControl.Persistence.Tests/Recoverability/GroupsDataStoreTests.cs +++ b/src/ServiceControl.Persistence.Tests/Recoverability/GroupsDataStoreTests.cs @@ -25,7 +25,7 @@ await Insert( InGroup(group, failedAt: Noon), InGroup(group, failedAt: Noon.AddHours(2))); - var view = (await GroupsStore.GetFailureGroupsByClassifier(Classifier, null)).Single(); + var view = (await GroupsStore.GetUnresolvedGroupsByClassifier(Classifier, null)).Single(); using (Assert.EnterMultipleScope()) { @@ -46,7 +46,7 @@ public async Task Returns_only_the_groups_of_the_requested_classifier() await Insert(InGroup(requested), InGroup(other)); - var groups = await GroupsStore.GetFailureGroupsByClassifier(Classifier, null); + var groups = await GroupsStore.GetUnresolvedGroupsByClassifier(Classifier, null); Assert.That(groups.Select(group => group.Id), Is.EqualTo(new[] { requested.Id })); } @@ -59,7 +59,7 @@ public async Task Narrows_the_groups_to_the_classifier_filter() await Insert(InGroup(matching), InGroup(other)); - var groups = await GroupsStore.GetFailureGroupsByClassifier(Classifier, "OrderPlaced"); + var groups = await GroupsStore.GetUnresolvedGroupsByClassifier(Classifier, "OrderPlaced"); Assert.That(groups.Select(group => group.Id), Is.EqualTo(new[] { matching.Id })); } @@ -74,7 +74,7 @@ await Insert( InGroup(group).ToFailedMessage(FailedMessageStatus.Archived), InGroup(group).ToFailedMessage(FailedMessageStatus.Resolved)); - var view = (await GroupsStore.GetFailureGroupsByClassifier(Classifier, null)).Single(); + var view = (await GroupsStore.GetUnresolvedGroupsByClassifier(Classifier, null)).Single(); Assert.That(view.Count, Is.EqualTo(1)); } @@ -88,7 +88,7 @@ await Insert( InGroup(group).ToFailedMessage(), InGroup(group).ToFailedMessage(FailedMessageStatus.Archived)); - var view = (await GroupsStore.GetArchivedFailureGroupsByClassifier(Classifier)).Single(); + var view = (await GroupsStore.GetArchivedGroupsByClassifier(Classifier)).Single(); using (Assert.EnterMultipleScope()) { @@ -109,7 +109,7 @@ await Insert( InGroup(newest, failedAt: Noon.AddHours(4)), InGroup(middle, failedAt: Noon.AddHours(2))); - var groups = await GroupsStore.GetFailureGroupsByClassifier(Classifier, null); + var groups = await GroupsStore.GetUnresolvedGroupsByClassifier(Classifier, null); Assert.That(groups.Select(group => group.Title), Is.EqualTo(new[] { "Newest", "Middle", "Oldest" })); } @@ -121,9 +121,9 @@ public async Task Returns_a_single_group_by_id() await Insert(InGroup(requested), InGroup(NewGroup("OrderCancelled"))); - var result = await GroupsStore.GetGroup(requested.Id, null, null); + var result = await GroupsStore.GetUnresolvedGroup(requested.Id, null, null); - var view = result.Results.Single(); + var view = result.Results; using (Assert.EnterMultipleScope()) { @@ -137,9 +137,9 @@ public async Task Returns_no_group_for_an_unknown_id() { await Insert(InGroup(NewGroup("OrderPlaced"))); - var result = await GroupsStore.GetGroup(Guid.NewGuid().ToString(), null, null); + var result = await GroupsStore.GetUnresolvedGroup(Guid.NewGuid().ToString(), null, null); - Assert.That(result.Results, Is.Empty); + Assert.That(result.Results, Is.Null); } [Test] @@ -149,7 +149,7 @@ public async Task Returns_an_archived_group_view_by_id() await Insert(InGroup(group).ToFailedMessage(FailedMessageStatus.Archived)); - var result = await GroupsStore.GetFailureGroupView(group.Id, null, null); + var result = await GroupsStore.GetArchivedGroup(group.Id, null, null); using (Assert.EnterMultipleScope()) { @@ -163,7 +163,7 @@ public async Task Returns_no_archived_group_view_for_an_unknown_id() { await Insert(InGroup(NewGroup("OrderPlaced")).ToFailedMessage(FailedMessageStatus.Archived)); - var result = await GroupsStore.GetFailureGroupView(Guid.NewGuid().ToString(), null, null); + var result = await GroupsStore.GetArchivedGroup(Guid.NewGuid().ToString(), null, null); Assert.That(result.Results, Is.Null); } diff --git a/src/ServiceControl.Persistence/IGroupsDataStore.cs b/src/ServiceControl.Persistence/IGroupsDataStore.cs index 3180a7a67a..422082f25c 100644 --- a/src/ServiceControl.Persistence/IGroupsDataStore.cs +++ b/src/ServiceControl.Persistence/IGroupsDataStore.cs @@ -8,12 +8,11 @@ namespace ServiceControl.Persistence public interface IGroupsDataStore { - Task> GetFailureGroupsByClassifier(string classifier, string classifierFilter); - Task> GetArchivedFailureGroupsByClassifier(string classifier); - Task GetCurrentForwardingBatch(); + Task> GetUnresolvedGroupsByClassifier(string classifier, string classifierFilter); + Task> GetArchivedGroupsByClassifier(string classifier); - Task>> GetGroup(string groupId, string status, string modified); - Task> GetFailureGroupView(string groupId, string status, string modified); + Task> GetUnresolvedGroup(string groupId, string status, string modified); + Task> GetArchivedGroup(string groupId, string status, string modified); Task>> GetGroupErrors(string groupId, string status, string modified, SortInfo sortInfo, PagingInfo pagingInfo); Task GetGroupErrorsCount(string groupId, string status, string modified); diff --git a/src/ServiceControl.Persistence/IRetryDocumentDataStore.cs b/src/ServiceControl.Persistence/IRetryDocumentDataStore.cs index f28b4f640f..65dd85dd33 100644 --- a/src/ServiceControl.Persistence/IRetryDocumentDataStore.cs +++ b/src/ServiceControl.Persistence/IRetryDocumentDataStore.cs @@ -21,13 +21,13 @@ Task CreateBatchDocument(string retrySessionId, string requestId, RetryT Task>> QueryOrphanedBatches(string retrySessionId); Task> QueryAvailableBatches(); + // GroupFetcher + Task GetCurrentForwardingBatch(); + // RetriesGateway Task GetBatchesForAll(DateTime cutoff, Func callback); Task GetBatchesForEndpoint(DateTime cutoff, string endpoint, Func callback); Task GetBatchesForFailedQueueAddress(DateTime cutoff, string failedQueueAddresspoint, FailedMessageStatus status, Func callback); Task GetBatchesForFailureGroup(string groupId, string groupTitle, string groupType, DateTime cutoff, Func callback); - - // RetryAllInGroupHandler - Task QueryFailureGroupViewOnGroupId(string groupId); } } \ No newline at end of file diff --git a/src/ServiceControl/MessageFailures/Api/ArchiveMessagesController.cs b/src/ServiceControl/MessageFailures/Api/ArchiveMessagesController.cs index 8a7ff276ed..88c91451ed 100644 --- a/src/ServiceControl/MessageFailures/Api/ArchiveMessagesController.cs +++ b/src/ServiceControl/MessageFailures/Api/ArchiveMessagesController.cs @@ -47,7 +47,7 @@ await auditLog.AuditedOperation(user, MessageActionKind.Archive, Permissions.Err [HttpGet] public async Task GetArchiveMessageGroups(string classifier = "Exception Type and Stack Trace") { - var results = await dataStore.GetArchivedFailureGroupsByClassifier(classifier); + var results = await dataStore.GetArchivedGroupsByClassifier(classifier); Response.WithDeterministicEtag(EtagHelper.CalculateEtag(results)); @@ -74,11 +74,11 @@ await auditLog.AuditedOperation(user, MessageActionKind.Archive, Permissions.Err [HttpGet] public async Task> GetGroup(string groupId, string status = default, string modified = default) { - var result = await dataStore.GetFailureGroupView(groupId, status, modified); + var result = await dataStore.GetArchivedGroup(groupId, status, modified); Response.WithEtag(result.QueryStats.ETag); - return result.Results; + return result.Results == null ? NotFound() : result.Results; } } } \ No newline at end of file diff --git a/src/ServiceControl/Recoverability/API/FailureGroupsController.cs b/src/ServiceControl/Recoverability/API/FailureGroupsController.cs index e86b82dfbe..6da6addce4 100644 --- a/src/ServiceControl/Recoverability/API/FailureGroupsController.cs +++ b/src/ServiceControl/Recoverability/API/FailureGroupsController.cs @@ -107,13 +107,13 @@ public async Task GetRetryHistory() [Authorize(Policy = Permissions.ErrorRecoverabilityGroupsView)] [Route("recoverability/groups/id/{groupId:required:minlength(1)}")] [HttpGet] - public async Task GetGroup(string groupId, string status = default, string modified = default) + public async Task> GetGroup(string groupId, string status = default, string modified = default) { - var result = await store.GetGroup(groupId, status, modified); + var result = await store.GetUnresolvedGroup(groupId, status, modified); Response.WithEtag(result.QueryStats.ETag); - return result.Results.FirstOrDefault(); + return result.Results == null ? NotFound() : result.Results; } } } \ No newline at end of file diff --git a/src/ServiceControl/Recoverability/API/GroupFetcher.cs b/src/ServiceControl/Recoverability/API/GroupFetcher.cs index a2e85058ab..c9fec0b158 100644 --- a/src/ServiceControl/Recoverability/API/GroupFetcher.cs +++ b/src/ServiceControl/Recoverability/API/GroupFetcher.cs @@ -8,17 +8,18 @@ public class GroupFetcher { - public GroupFetcher(IGroupsDataStore store, IRetryHistoryDataStore retryStore, RetryingManager retryingManager, IArchiveMessages archiver) + public GroupFetcher(IGroupsDataStore store, IRetryHistoryDataStore retryStore, IRetryDocumentDataStore retryDocumentStore, RetryingManager retryingManager, IArchiveMessages archiver) { this.store = store; this.retryStore = retryStore; + this.retryDocumentStore = retryDocumentStore; this.retryingManager = retryingManager; this.archiver = archiver; } public async Task GetGroups(string classifier, string classifierFilter) { - var dbGroups = await store.GetFailureGroupsByClassifier(classifier, classifierFilter); + var dbGroups = await store.GetUnresolvedGroupsByClassifier(classifier, classifierFilter); var retryHistory = await retryStore.GetRetryHistory(); var unacknowledgedRetries = retryHistory.GetUnacknowledgedByClassifier(classifier); @@ -32,7 +33,7 @@ public async Task GetGroups(string classifier, string classifi openGroups = MapOpenGroups(openGroups, archiver.GetArchivalOperations()).ToList(); openGroups = openGroups.Where(group => !closedGroups.Any(closedGroup => closedGroup.Id == group.Id)).ToList(); - var currentForwardingBatch = await store.GetCurrentForwardingBatch(); + var currentForwardingBatch = await retryDocumentStore.GetCurrentForwardingBatch(); MakeSureForwardingBatchIsIncludedAsOpen(classifier, currentForwardingBatch, openGroups); var groups = openGroups.Union(closedGroups); @@ -190,6 +191,7 @@ static HistoricRetryOperation GetLatestHistoricOperation(RetryHistory history, s readonly IGroupsDataStore store; readonly IRetryHistoryDataStore retryStore; + readonly IRetryDocumentDataStore retryDocumentStore; readonly RetryingManager retryingManager; readonly IArchiveMessages archiver; } diff --git a/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs b/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs index 38da31b01f..9a4e8e7761 100644 --- a/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs +++ b/src/ServiceControl/Recoverability/Retrying/Handlers/RetryAllInGroupHandler.cs @@ -9,7 +9,7 @@ namespace ServiceControl.Recoverability using ServiceControl.Persistence.Recoverability; [Handler] - class RetryAllInGroupHandler(RetriesGateway retries, RetryingManager retryingManager, IArchiveMessages archiver, IRetryDocumentDataStore dataStore, ILogger logger) + class RetryAllInGroupHandler(RetriesGateway retries, RetryingManager retryingManager, IArchiveMessages archiver, IGroupsDataStore dataStore, ILogger logger) : IHandleMessages { public async Task Handle(RetryAllInGroup message, IMessageHandlerContext context) @@ -27,7 +27,7 @@ public async Task Handle(RetryAllInGroup message, IMessageHandlerContext context } - var group = await dataStore.QueryFailureGroupViewOnGroupId(message.GroupId); + var group = (await dataStore.GetUnresolvedGroup(message.GroupId, null, null)).Results; string originator = null; if (group?.Title != null)