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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste
services.AddSingleton<IMonitoringDataStore, MonitoringDataStore>();
services.AddSingleton<IQueueAddressStore, QueueAddressStore>();
services.AddSingleton<IRetryBatchesDataStore, RetryBatchesDataStore>();
services.AddSingleton<IRetryDocumentDataStore, RetryDocumentDataStore>();
services.AddSingleton<IRetryBatchStore, RetryBatchStore>();
services.AddSingleton<IRetryHistoryDataStore, RetryHistoryDataStore>();
services.AddSingleton<IEndpointSettingsStore, EndpointSettingsStore>();
services.AddSingleton<ITrialLicenseDataProvider, TrialLicenseDataProvider>();
Expand Down
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 @@ -4,9 +4,15 @@ namespace ServiceControl.Persistence.EFCore.Implementation;

public class MessageRedirectsDataStore : IMessageRedirectsDataStore
{
public Task<MessageRedirectsCollection> GetOrCreate() =>
public Task<IReadOnlyList<MessageRedirect>> GetRedirects() =>
throw new NotImplementedException();

public Task Save(MessageRedirectsCollection redirects) =>
public Task AddRedirect(MessageRedirect redirect) =>
throw new NotImplementedException();

public Task UpdateRedirect(MessageRedirect redirect) =>
throw new NotImplementedException();

public Task RemoveRedirect(MessageRedirect redirect) =>
throw new NotImplementedException();
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
namespace ServiceControl.Persistence.EFCore.Implementation;

using ServiceControl.MessageFailures;
using ServiceControl.Persistence.Infrastructure;
using ServiceControl.Recoverability;

public class RetryBatchStore : IRetryBatchStore
{
public Task<string> CreateBatch(string retrySessionId, string requestId, RetryType retryType,
string[] failedMessageRetryIds, string originator, DateTime startTime, DateTime? last = null,
string? batchName = null, string? classifier = null,
string? initiatedById = null, string? initiatedByName = null, string? operationId = null) =>
throw new NotImplementedException();

public Task AssignMessagesToBatch(string batchId, string[] messageIds) =>
throw new NotImplementedException();

public Task MoveBatchToStaging(string batchId) =>
throw new NotImplementedException();

public Task<QueryResult<IList<RetryBatch>>> GetOrphanedBatches(string retrySessionId) =>
throw new NotImplementedException();

public Task<IList<RetryBatchGroup>> GetAvailableBatchGroups() =>
throw new NotImplementedException();

public Task<ForwardingRetryBatch> GetCurrentForwardingBatch() =>
throw new NotImplementedException();

public Task ForEachUnresolvedMessage(Func<string, DateTime, Task> callback) =>
throw new NotImplementedException();

public Task ForEachUnresolvedMessageForEndpoint(string endpoint, Func<string, DateTime, Task> callback) =>
throw new NotImplementedException();

public Task ForEachMessageForQueueAddress(string failedQueueAddress, FailedMessageStatus status, Func<string, DateTime, Task> callback) =>
throw new NotImplementedException();

public Task ForEachUnresolvedMessageInGroup(string groupId, Func<string, DateTime, Task> callback) =>
throw new NotImplementedException();
}
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,6 @@ public Task<RetryBatch> GetStagingBatch() =>
public Task Store(RetryBatchNowForwarding retryBatchNowForwarding) =>
throw new NotImplementedException();

public Task<MessageRedirectsCollection> GetOrCreateMessageRedirectsCollection() =>
throw new NotImplementedException();

public Task CancelExpiration(FailedMessage failedMessage) =>
throw new NotImplementedException();

Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
@@ -1,32 +1,39 @@
namespace ServiceControl.Persistence.RavenDB.MessageRedirects
{
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using ServiceControl.Persistence.MessageRedirects;

class MessageRedirectsDataStore(IRavenSessionProvider sessionProvider) : IMessageRedirectsDataStore
{
public const string CollectionId = "messageredirects";

public async Task<MessageRedirectsCollection> GetOrCreate()
public async Task<IReadOnlyList<MessageRedirect>> GetRedirects()
{
using var session = await sessionProvider.OpenSession();
var redirects = await session.LoadAsync<MessageRedirectsCollection>(CollectionId);
var document = await session.LoadAsync<MessageRedirectsCollection>(CollectionId);

if (redirects != null)
{
redirects.ETag = session.Advanced.GetChangeVectorFor(redirects);
redirects.LastModified = session.Advanced.GetLastModifiedFor(redirects).Value;
return document == null ? [] : document.ToRedirects();
}

return redirects;
}
public Task AddRedirect(MessageRedirect redirect) => Mutate(document => document.Add(redirect));

return new MessageRedirectsCollection();
}
public Task UpdateRedirect(MessageRedirect redirect) => Mutate(document => document.Update(redirect));

public Task RemoveRedirect(MessageRedirect redirect) => Mutate(document => document.Remove(redirect));

public async Task Save(MessageRedirectsCollection redirects)
async Task Mutate(Action<MessageRedirectsCollection> mutate)
{
using var session = await sessionProvider.OpenSession();
await session.StoreAsync(redirects, redirects.ETag, CollectionId);
var document = await session.LoadAsync<MessageRedirectsCollection>(CollectionId);
var changeVector = document == null ? null : session.Advanced.GetChangeVectorFor(document);

document ??= new MessageRedirectsCollection();

mutate(document);

await session.StoreAsync(document, changeVector, CollectionId);
await session.SaveChangesAsync();
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
namespace ServiceControl.Persistence.RavenDB.MessageRedirects
{
using System;
using System.Collections.Generic;
using System.Linq;
using ServiceControl.Persistence.MessageRedirects;

class MessageRedirectsCollection
{
public List<StoredRedirect> Redirects { get; set; } = [];

public IReadOnlyList<MessageRedirect> ToRedirects() =>
[.. Redirects.Select(redirect => new MessageRedirect
{
FromPhysicalAddress = redirect.FromPhysicalAddress,
ToPhysicalAddress = redirect.ToPhysicalAddress,
LastModified = new DateTime(redirect.LastModifiedTicks, DateTimeKind.Utc)
})];

public void Add(MessageRedirect redirect) => Redirects.Add(new StoredRedirect
{
FromPhysicalAddress = redirect.FromPhysicalAddress,
ToPhysicalAddress = redirect.ToPhysicalAddress,
LastModifiedTicks = redirect.LastModified.Ticks
});

public void Update(MessageRedirect redirect)
{
var existing = Redirects.SingleOrDefault(stored => stored.FromPhysicalAddress == redirect.FromPhysicalAddress);

if (existing == null)
{
return;
}

existing.ToPhysicalAddress = redirect.ToPhysicalAddress;
existing.LastModifiedTicks = redirect.LastModified.Ticks;
}

public void Remove(MessageRedirect redirect) =>
Redirects.RemoveAll(stored => stored.FromPhysicalAddress == redirect.FromPhysicalAddress);

public class StoredRedirect
{
public string FromPhysicalAddress { get; set; }
public string ToPhysicalAddress { get; set; }
public long LastModifiedTicks { get; set; }
}
}
}
2 changes: 1 addition & 1 deletion src/ServiceControl.Persistence.RavenDB/RavenPersistence.cs
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ public void AddPersistence(IServiceCollection services)
services.AddSingleton<IMonitoringDataStore, RavenMonitoringDataStore>();
services.AddSingleton<IQueueAddressStore, QueueAddressStore>();
services.AddSingleton<IRetryBatchesDataStore, RetryBatchesDataStore>();
services.AddSingleton<IRetryDocumentDataStore, RetryDocumentDataStore>();
services.AddSingleton<IRetryBatchStore, RetryDocumentDataStore>();
services.AddSingleton<IRetryHistoryDataStore, RetryHistoryDataStore>();
services.AddSingleton<IEndpointSettingsStore, EndpointSettingsStore>();
services.AddSingleton<ITrialLicenseDataProvider, TrialLicenseDataProvider>();
Expand Down
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
14 changes: 0 additions & 14 deletions src/ServiceControl.Persistence.RavenDB/RetryBatchesManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -55,20 +55,6 @@ public async Task<RetryBatch> GetStagingBatch()
public async Task Store(RetryBatchNowForwarding retryBatchNowForwarding) =>
await Session.StoreAsync(retryBatchNowForwarding, RetryDocumentDataStore.NowForwardingDocumentId);

public async Task<MessageRedirectsCollection> GetOrCreateMessageRedirectsCollection()
{
var redirects = await Session.LoadAsync<MessageRedirectsCollection>(MessageRedirectsDataStore.CollectionId);

if (redirects != null)
{
redirects.ETag = Session.Advanced.GetChangeVectorFor(redirects);
redirects.LastModified = Session.Advanced.GetLastModifiedFor(redirects)!.Value;
return redirects;
}

return new MessageRedirectsCollection();
}

public Task CancelExpiration(FailedMessage failedMessage)
{
expirationManager.CancelExpiration(Session, failedMessage);
Expand Down
Loading
Loading