Skip to content
Open
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 @@ -12,7 +12,7 @@ namespace ServiceControl.Persistence.EFCore.Implementation;

public class GroupsDataStore(IServiceScopeFactory scopeFactory) : DataStoreBase(scopeFactory), IGroupsDataStore
{
public Task<IList<FailureGroupView>> GetFailureGroupsByClassifier(string classifier, string classifierFilter) =>
public Task<IList<FailureGroupView>> GetUnresolvedGroupsByClassifier(string classifier, string classifierFilter) =>
ExecuteWithDbContext(dbContext =>
{
var groups = ByClassifier(dbContext, classifier);
Expand All @@ -25,30 +25,15 @@ public Task<IList<FailureGroupView>> GetFailureGroupsByClassifier(string classif
return MostRecent(groups.AggregateGroups(WithStatus(dbContext, FailedMessageStatus.Unresolved)));
});

public Task<IList<FailureGroupView>> GetArchivedFailureGroupsByClassifier(string classifier) =>
public Task<IList<FailureGroupView>> 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<RetryBatch> GetCurrentForwardingBatch() =>
throw new NotImplementedException();

public Task<QueryResult<IList<FailureGroupView>>> GetGroup(string groupId, string status, string modified) =>
ExecuteWithDbContext(async dbContext =>
{
var groups = await ById(dbContext, groupId, FailedMessageStatus.Unresolved, status, modified).ToListAsync();
public Task<QueryResult<FailureGroupView>> GetUnresolvedGroup(string groupId, string status, string modified) =>
ExecuteWithDbContext(dbContext => SingleGroup(dbContext, groupId, FailedMessageStatus.Unresolved, status, modified));

return new QueryResult<IList<FailureGroupView>>(groups, groups.ToQueryStatsInfo());
});

public Task<QueryResult<FailureGroupView>> 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<FailureGroupView>(groups.FirstOrDefault()!, groups.ToQueryStatsInfo());
});
public Task<QueryResult<FailureGroupView>> GetArchivedGroup(string groupId, string status, string modified) =>
ExecuteWithDbContext(dbContext => SingleGroup(dbContext, groupId, FailedMessageStatus.Archived, status, modified));

public Task<QueryResult<IList<FailedMessageView>>> GetGroupErrors(string groupId, string status, string modified, SortInfo sortInfo, PagingInfo pagingInfo) =>
ExecuteWithDbContext(dbContext => InGroup(dbContext, groupId, status, modified).ToPagedResult(pagingInfo, sortInfo));
Expand All @@ -67,18 +52,18 @@ static IQueryable<FailedMessageGroupEntity> ByClassifier(ServiceControlDbContext
.AsNoTracking()
.Where(group => group.Type == classifier);

/// <summary>
/// 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 <paramref name="baseline" /> stands in for here.
/// </summary>
static IQueryable<FailureGroupView> ById(ServiceControlDbContext dbContext, string groupId, FailedMessageStatus baseline, string status, string modified) =>
dbContext.FailedMessageGroups
static async Task<QueryResult<FailureGroupView>> 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<FailureGroupView>(groups.FirstOrDefault()!, groups.ToQueryStatsInfo());
}

static IQueryable<FailedMessageEntity> WithStatus(ServiceControlDbContext dbContext, FailedMessageStatus status) =>
dbContext.FailedMessages
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,6 @@ public Task GetBatchesForFailedQueueAddress(DateTime cutoff, string failedQueueA
public Task GetBatchesForFailureGroup(string groupId, string groupTitle, string groupType, DateTime cutoff, Func<string, DateTime, Task> callback) =>
throw new NotImplementedException();

public Task<FailureGroupView> QueryFailureGroupViewOnGroupId(string groupId) =>
public Task<RetryBatch> GetCurrentForwardingBatch() =>
throw new NotImplementedException();
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ namespace ServiceControl.Persistence.RavenDB.Recoverability

class GroupsDataStore(IRavenSessionProvider sessionProvider) : IGroupsDataStore
{
public async Task<IList<FailureGroupView>> GetFailureGroupsByClassifier(string classifier, string classifierFilter)
public async Task<IList<FailureGroupView>> GetUnresolvedGroupsByClassifier(string classifier, string classifierFilter)
{
using var session = await sessionProvider.OpenSession();
var query = Queryable.Where(session.Query<FailureGroupView, FailureGroupsViewIndex>(), v => v.Type == classifier);
Expand All @@ -40,7 +40,7 @@ public async Task<IList<FailureGroupView>> GetFailureGroupsByClassifier(string c
return groups;
}

public async Task<IList<FailureGroupView>> GetArchivedFailureGroupsByClassifier(string classifier)
public async Task<IList<FailureGroupView>> GetArchivedGroupsByClassifier(string classifier)
{
using var session = await sessionProvider.OpenSession();
var groups = session
Expand All @@ -55,30 +55,21 @@ public async Task<IList<FailureGroupView>> GetArchivedFailureGroupsByClassifier(
return results;
}

public async Task<RetryBatch> GetCurrentForwardingBatch()
public async Task<QueryResult<FailureGroupView>> GetUnresolvedGroup(string groupId, string status, string modified)
{
using var session = await sessionProvider.OpenSession();
var nowForwarding = await session.Include<RetryBatchNowForwarding, RetryBatch>(r => r.RetryBatchId)
.LoadAsync<RetryBatchNowForwarding>(RetryDocumentDataStore.NowForwardingDocumentId);

return nowForwarding == null ? null : await session.LoadAsync<RetryBatch>(nowForwarding.RetryBatchId);
}

public async Task<QueryResult<IList<FailureGroupView>>> GetGroup(string groupId, string status, string modified)
{
using var session = await sessionProvider.OpenSession();
var queryResult = await session.Advanced
var document = await session.Advanced
.AsyncDocumentQuery<FailureGroupView, FailureGroupsViewIndex>()
.Statistics(out var stats)
.WhereEquals(group => group.Id, groupId)
.FilterByStatusWhere(status)
.FilterByLastModifiedRange(modified)
.ToListAsync();
.FirstOrDefaultAsync();

return queryResult.ToQueryResult(stats);
return new QueryResult<FailureGroupView>(document, stats.ToQueryStatsInfo());
}

public async Task<QueryResult<FailureGroupView>> GetFailureGroupView(string groupId, string status, string modified)
public async Task<QueryResult<FailureGroupView>> GetArchivedGroup(string groupId, string status, string modified)
{
using var session = await sessionProvider.OpenSession();
var document = await session.Advanced
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -200,12 +200,13 @@ public async Task GetBatchesForFailureGroup(string groupId, string groupTitle, s
}
}

public async Task<FailureGroupView> QueryFailureGroupViewOnGroupId(string groupId)
public async Task<RetryBatch> GetCurrentForwardingBatch()
{
using var session = await sessionProvider.OpenSession();
var group = await session.Query<FailureGroupView, FailureGroupsViewIndex>()
.FirstOrDefaultAsync(x => x.Id == groupId);
return group;
var nowForwarding = await session.Include<RetryBatchNowForwarding, RetryBatch>(r => r.RetryBatchId)
.LoadAsync<RetryBatchNowForwarding>(NowForwardingDocumentId);

return nowForwarding == null ? null : await session.LoadAsync<RetryBatch>(nowForwarding.RetryBatchId);
}

public static string MakeDocumentId(string messageUniqueId) => "RetryBatches/" + messageUniqueId;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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())
{
Expand All @@ -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 }));
}
Expand All @@ -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 }));
}
Expand All @@ -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));
}
Expand All @@ -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())
{
Expand All @@ -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" }));
}
Expand All @@ -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())
{
Expand All @@ -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]
Expand All @@ -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())
{
Expand All @@ -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);
}
Expand Down
9 changes: 4 additions & 5 deletions src/ServiceControl.Persistence/IGroupsDataStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,11 @@ namespace ServiceControl.Persistence

public interface IGroupsDataStore
{
Task<IList<FailureGroupView>> GetFailureGroupsByClassifier(string classifier, string classifierFilter);
Task<IList<FailureGroupView>> GetArchivedFailureGroupsByClassifier(string classifier);
Task<RetryBatch> GetCurrentForwardingBatch();
Task<IList<FailureGroupView>> GetUnresolvedGroupsByClassifier(string classifier, string classifierFilter);
Task<IList<FailureGroupView>> GetArchivedGroupsByClassifier(string classifier);

Task<QueryResult<IList<FailureGroupView>>> GetGroup(string groupId, string status, string modified);
Task<QueryResult<FailureGroupView>> GetFailureGroupView(string groupId, string status, string modified);
Task<QueryResult<FailureGroupView>> GetUnresolvedGroup(string groupId, string status, string modified);
Task<QueryResult<FailureGroupView>> GetArchivedGroup(string groupId, string status, string modified);
Task<QueryResult<IList<FailedMessageView>>> GetGroupErrors(string groupId, string status, string modified, SortInfo sortInfo, PagingInfo pagingInfo);
Task<QueryStatsInfo> GetGroupErrorsCount(string groupId, string status, string modified);

Expand Down
6 changes: 3 additions & 3 deletions src/ServiceControl.Persistence/IRetryDocumentDataStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@ Task<string> CreateBatchDocument(string retrySessionId, string requestId, RetryT
Task<QueryResult<IList<RetryBatch>>> QueryOrphanedBatches(string retrySessionId);
Task<IList<RetryBatchGroup>> QueryAvailableBatches();

// GroupFetcher
Task<RetryBatch> GetCurrentForwardingBatch();

// RetriesGateway
Task GetBatchesForAll(DateTime cutoff, Func<string, DateTime, Task> callback);
Task GetBatchesForEndpoint(DateTime cutoff, string endpoint, Func<string, DateTime, Task> callback);
Task GetBatchesForFailedQueueAddress(DateTime cutoff, string failedQueueAddresspoint, FailedMessageStatus status, Func<string, DateTime, Task> callback);
Task GetBatchesForFailureGroup(string groupId, string groupTitle, string groupType, DateTime cutoff, Func<string, DateTime, Task> callback);

// RetryAllInGroupHandler
Task<FailureGroupView> QueryFailureGroupViewOnGroupId(string groupId);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ await auditLog.AuditedOperation(user, MessageActionKind.Archive, Permissions.Err
[HttpGet]
public async Task<IActionResult> 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));

Expand All @@ -74,11 +74,11 @@ await auditLog.AuditedOperation(user, MessageActionKind.Archive, Permissions.Err
[HttpGet]
public async Task<ActionResult<FailureGroupView>> 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;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -107,13 +107,13 @@ public async Task<RetryHistory> GetRetryHistory()
[Authorize(Policy = Permissions.ErrorRecoverabilityGroupsView)]
[Route("recoverability/groups/id/{groupId:required:minlength(1)}")]
[HttpGet]
public async Task<FailureGroupView> GetGroup(string groupId, string status = default, string modified = default)
public async Task<ActionResult<FailureGroupView>> 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;
}
}
}
8 changes: 5 additions & 3 deletions src/ServiceControl/Recoverability/API/GroupFetcher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<GroupOperation[]> 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);

Expand All @@ -32,7 +33,7 @@ public async Task<GroupOperation[]> 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);
Expand Down Expand Up @@ -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;
}
Expand Down
Loading
Loading