Skip to content
Draft
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 @@ -11,7 +11,8 @@ public class DatabaseConfiguration(
int dataSpaceRemainingThreshold,
int minimumStorageLeftRequiredForIngestion,
ServerConfiguration serverConfiguration,
TimeSpan bulkInsertCommitTimeout)
TimeSpan bulkInsertCommitTimeout,
TimeSpan queryTimeout)
{
public string Name { get; } = name;

Expand All @@ -30,5 +31,7 @@ public class DatabaseConfiguration(
public int MinimumStorageLeftRequiredForIngestion { get; internal set; } = minimumStorageLeftRequiredForIngestion; //Setting for ATT only

public TimeSpan BulkInsertCommitTimeout { get; } = bulkInsertCommitTimeout;

public TimeSpan QueryTimeout { get; } = queryTimeout;
}
}
161 changes: 85 additions & 76 deletions src/ServiceControl.Audit.Persistence.RavenDB/RavenAuditDataStore.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
namespace ServiceControl.Audit.Persistence.RavenDB
namespace ServiceControl.Audit.Persistence.RavenDB
{
using System;
using System.Collections.Generic;
Expand All @@ -11,6 +11,7 @@
using Raven.Client.Documents;
using ServiceControl.Audit.Auditing;
using ServiceControl.Audit.Infrastructure;
using ServiceControl.Infrastructure;
using ServiceControl.SagaAudit;
using Transformers;

Expand All @@ -28,81 +29,86 @@ public async Task<QueryResult<SagaHistory>> QuerySagaHistoryById(Guid input, Can
return sagaHistory == null ? QueryResult<SagaHistory>.Empty() : new QueryResult<SagaHistory>(sagaHistory, stats.ToQueryStatsInfo());
}

public async Task<QueryResult<IList<MessagesView>>> GetMessages(bool includeSystemMessages, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.FilterBySentTimeRange(timeSentRange)
.IncludeSystemMessagesWhere(includeSystemMessages)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: cancellationToken);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}

public async Task<QueryResult<IList<MessagesView>>> QueryMessages(string searchParam, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.Search(x => x.Query, searchParam)
.FilterBySentTimeRange(timeSentRange)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: cancellationToken);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}

public async Task<QueryResult<IList<MessagesView>>> QueryMessagesByReceivingEndpointAndKeyword(string endpoint, string keyword, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.Search(x => x.Query, keyword)
.Where(m => m.ReceivingEndpointName == endpoint)
.FilterBySentTimeRange(timeSentRange)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: cancellationToken);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}

public async Task<QueryResult<IList<MessagesView>>> QueryMessagesByReceivingEndpoint(bool includeSystemMessages, string endpointName, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.IncludeSystemMessagesWhere(includeSystemMessages)
.Where(m => m.ReceivingEndpointName == endpointName)
.FilterBySentTimeRange(timeSentRange)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: cancellationToken);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}
public Task<QueryResult<IList<MessagesView>>> GetMessages(bool includeSystemMessages, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) =>
WithQueryTimeout(async token =>
{
using var session = await sessionProvider.OpenSession(cancellationToken: token);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.FilterBySentTimeRange(timeSentRange)
.IncludeSystemMessagesWhere(includeSystemMessages)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: token);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}, cancellationToken);

public Task<QueryResult<IList<MessagesView>>> QueryMessages(string searchParam, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) =>
WithQueryTimeout(async token =>
{
using var session = await sessionProvider.OpenSession(cancellationToken: token);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.Search(x => x.Query, searchParam)
.FilterBySentTimeRange(timeSentRange)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: token);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}, cancellationToken);

public Task<QueryResult<IList<MessagesView>>> QueryMessagesByReceivingEndpointAndKeyword(string endpoint, string keyword, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) =>
WithQueryTimeout(async token =>
{
using var session = await sessionProvider.OpenSession(cancellationToken: token);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.Search(x => x.Query, keyword)
.Where(m => m.ReceivingEndpointName == endpoint)
.FilterBySentTimeRange(timeSentRange)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: token);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}, cancellationToken);

public Task<QueryResult<IList<MessagesView>>> QueryMessagesByReceivingEndpoint(bool includeSystemMessages, string endpointName, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) =>
WithQueryTimeout(async token =>
{
using var session = await sessionProvider.OpenSession(cancellationToken: token);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.IncludeSystemMessagesWhere(includeSystemMessages)
.Where(m => m.ReceivingEndpointName == endpointName)
.FilterBySentTimeRange(timeSentRange)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: token);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}, cancellationToken);

public Task<QueryResult<IList<MessagesView>>> QueryMessagesByConversationId(string conversationId, PagingInfo pagingInfo, SortInfo sortInfo, CancellationToken cancellationToken = default) =>
WithQueryTimeout(async token =>
{
using var session = await sessionProvider.OpenSession(cancellationToken: token);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.Where(m => m.ConversationId == conversationId)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: token);

public async Task<QueryResult<IList<MessagesView>>> QueryMessagesByConversationId(string conversationId, PagingInfo pagingInfo, SortInfo sortInfo, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
var results = await session.Query<MessagesViewIndex.SortAndFilterOptions>(GetIndexName(isFullTextSearchEnabled))
.Statistics(out var stats)
.Where(m => m.ConversationId == conversationId)
.Sort(sortInfo)
.Paging(pagingInfo)
.ToMessagesView()
.ToListAsync(token: cancellationToken);

return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}
return new QueryResult<IList<MessagesView>>(results, stats.ToQueryStatsInfo());
}, cancellationToken);

public async Task<MessageBodyView> GetMessageBody(string messageId, CancellationToken cancellationToken = default)
{
Expand Down Expand Up @@ -184,8 +190,11 @@ public async Task<QueryResult<IList<AuditCount>>> QueryAuditCounts(string endpoi
return new QueryResult<IList<AuditCount>>(results, QueryStatsInfo.Zero);
}

Task<T> WithQueryTimeout<T>(Func<CancellationToken, Task<T>> query, CancellationToken cancellationToken) =>
QueryTimeLimit.Run(query, databaseConfiguration.QueryTimeout, RavenPersistenceConfiguration.QueryTimeoutSettingName, cancellationToken);

static string GetIndexName(bool isFullTextSearchEnabled) => isFullTextSearchEnabled ? "MessagesViewIndexWithFullTextSearch" : "MessagesViewIndex";

bool isFullTextSearchEnabled = databaseConfiguration.EnableFullTextSearch;
}
}
}
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
namespace ServiceControl.Audit.Persistence.RavenDB
namespace ServiceControl.Audit.Persistence.RavenDB
{
using System;
using System.Collections.Generic;
Expand All @@ -23,6 +23,8 @@ public class RavenPersistenceConfiguration : IPersistenceConfiguration
public const string MinimumStorageLeftRequiredForIngestionKey = "MinimumStorageLeftRequiredForIngestion";
public const string BulkInsertCommitTimeoutInSecondsKey = "BulkInsertCommitTimeoutInSeconds";
public const string DataSpaceRemainingThresholdKey = "DataSpaceRemainingThreshold";
public const string QueryTimeoutInSecondsKey = QueryTimeLimit.SettingName;
public const string QueryTimeoutSettingName = "ServiceControl.Audit/" + QueryTimeLimit.SettingName;

public IEnumerable<string> ConfigurationKeys => new[]{
DatabaseNameKey,
Expand All @@ -37,7 +39,8 @@ public class RavenPersistenceConfiguration : IPersistenceConfiguration
RavenDbLogLevelKey,
DataSpaceRemainingThresholdKey,
MinimumStorageLeftRequiredForIngestionKey,
BulkInsertCommitTimeoutInSecondsKey
BulkInsertCommitTimeoutInSecondsKey,
QueryTimeoutInSecondsKey
};

public string Name => "RavenDB";
Expand Down Expand Up @@ -121,6 +124,8 @@ internal static DatabaseConfiguration GetDatabaseConfiguration(PersistenceSettin

var bulkInsertTimeout = TimeSpan.FromSeconds(GetBulkInsertCommitTimeout(settings));

var queryTimeout = GetQueryTimeout(settings);

return new DatabaseConfiguration(
databaseName,
expirationProcessTimerInSeconds,
Expand All @@ -130,7 +135,8 @@ internal static DatabaseConfiguration GetDatabaseConfiguration(PersistenceSettin
dataSpaceRemainingThreshold,
minimumStorageLeftRequiredForIngestion,
serverConfiguration,
bulkInsertTimeout);
bulkInsertTimeout,
queryTimeout);
}

static int GetExpirationProcessTimerInSeconds(PersistenceSettings settings)
Expand Down Expand Up @@ -185,6 +191,18 @@ static int GetBulkInsertCommitTimeout(PersistenceSettings settings)
return bulkInsertCommitTimeoutInSeconds;
}

static TimeSpan GetQueryTimeout(PersistenceSettings settings)
{
var queryTimeoutInSeconds = QueryTimeLimit.DefaultSeconds;

if (settings.PersisterSpecificSettings.TryGetValue(QueryTimeoutInSecondsKey, out var queryTimeoutString))
{
queryTimeoutInSeconds = int.Parse(queryTimeoutString);
}

return QueryTimeLimit.Validate(queryTimeoutInSeconds, QueryTimeoutSettingName, Logger);
}

static string GetLogPath(PersistenceSettings settings)
{
if (!settings.PersisterSpecificSettings.TryGetValue(LogPathKey, out var logPath))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,57 @@ public void Should_throw_if_both_path_or_connection_string_is_configured()
Assert.Throws<InvalidOperationException>(() => RavenPersistenceConfiguration.GetDatabaseConfiguration(settings));
}

[Test]
public void Should_default_query_timeout_to_one_minute()
{
var settings = BuildSettings();

settings.PersisterSpecificSettings[RavenPersistenceConfiguration.ConnectionStringKey] = "connection string";

var configuration = RavenPersistenceConfiguration.GetDatabaseConfiguration(settings);

Assert.That(configuration.QueryTimeout, Is.EqualTo(TimeSpan.FromMinutes(1)));
}

[Test]
public void Should_apply_query_timeout_setting()
{
var settings = BuildSettings();

settings.PersisterSpecificSettings[RavenPersistenceConfiguration.ConnectionStringKey] = "connection string";
settings.PersisterSpecificSettings[RavenPersistenceConfiguration.QueryTimeoutInSecondsKey] = "120";

var configuration = RavenPersistenceConfiguration.GetDatabaseConfiguration(settings);

Assert.That(configuration.QueryTimeout, Is.EqualTo(TimeSpan.FromSeconds(120)));
}

[Test]
public void Should_fall_back_to_default_query_timeout_when_value_is_not_positive()
{
var settings = BuildSettings();

settings.PersisterSpecificSettings[RavenPersistenceConfiguration.ConnectionStringKey] = "connection string";
settings.PersisterSpecificSettings[RavenPersistenceConfiguration.QueryTimeoutInSecondsKey] = "0";

var configuration = RavenPersistenceConfiguration.GetDatabaseConfiguration(settings);

Assert.That(configuration.QueryTimeout, Is.EqualTo(TimeSpan.FromMinutes(1)));
}

[Test]
public void Should_fall_back_to_default_query_timeout_when_value_exceeds_maximum()
{
var settings = BuildSettings();

settings.PersisterSpecificSettings[RavenPersistenceConfiguration.ConnectionStringKey] = "connection string";
settings.PersisterSpecificSettings[RavenPersistenceConfiguration.QueryTimeoutInSecondsKey] = "3700";

var configuration = RavenPersistenceConfiguration.GetDatabaseConfiguration(settings);

Assert.That(configuration.QueryTimeout, Is.EqualTo(TimeSpan.FromMinutes(1)));
}

PersistenceSettings BuildSettings()
{
return new PersistenceSettings(TimeSpan.FromMinutes(2), true, 100000);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
namespace ServiceControl.UnitTests;

using System;
using System.Threading;
using System.Threading.Tasks;
using NUnit.Framework;
using Raven.Client.Documents.Session;
using ServiceControl.Audit.Infrastructure;
using ServiceControl.Audit.Persistence.RavenDB;

class QueryTimeoutTests
{
[Test]
public void Should_cancel_query_that_exceeds_the_allowed_query_time()
{
var dataStore = new RavenAuditDataStore(new NeverCompletingSessionProvider(), BuildConfiguration(queryTimeout: TimeSpan.FromMilliseconds(50)));

var exception = Assert.ThrowsAsync<TimeoutException>(() => dataStore.QueryMessages("search", new PagingInfo(), new SortInfo("time_sent", "desc"), new DateTimeRange((DateTime?)null, null)));

Assert.That(exception.InnerException, Is.InstanceOf<OperationCanceledException>());
}

[Test]
public void Should_propagate_caller_cancellation_instead_of_timing_out()
{
var dataStore = new RavenAuditDataStore(new NeverCompletingSessionProvider(), BuildConfiguration(queryTimeout: TimeSpan.FromSeconds(30)));

using var callerTokenSource = new CancellationTokenSource(TimeSpan.FromMilliseconds(50));

var exception = Assert.CatchAsync(() => dataStore.QueryMessages("search", new PagingInfo(), new SortInfo("time_sent", "desc"), new DateTimeRange((DateTime?)null, null), callerTokenSource.Token));

Assert.That(exception, Is.InstanceOf<OperationCanceledException>());
}

static DatabaseConfiguration BuildConfiguration(TimeSpan queryTimeout) =>
new("audit", 60, true, TimeSpan.FromMinutes(5), 120000, 5, 5, new ServerConfiguration("http://localhost:33334"), TimeSpan.FromSeconds(60), queryTimeout);

class NeverCompletingSessionProvider : IRavenSessionProvider
{
public async ValueTask<IAsyncDocumentSession> OpenSession(SessionOptions options = default, CancellationToken cancellationToken = default)
{
await Task.Delay(Timeout.Infinite, cancellationToken);
return null;
}
}
}
Loading
Loading