Skip to content

Commit 86b74ea

Browse files
authored
Merge pull request #5635 from Particular/rhys/cancellation-endpointsettings
Add cancellation token to GetAllEndpointSettings
2 parents ff72c01 + 7beda68 commit 86b74ea

6 files changed

Lines changed: 11 additions & 12 deletions

File tree

src/ServiceControl.Persistence.EFCore/Implementation/EndpointSettingsStore.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
22

33
public class EndpointSettingsStore : IEndpointSettingsStore
44
{
5-
public IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings() =>
5+
public IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings(CancellationToken cancellationToken) =>
66
throw new NotImplementedException();
77

88
public Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token) =>

src/ServiceControl.Persistence.RavenDB/EndpointSettingsStore.cs

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
namespace ServiceControl.Persistence.RavenDB;
22

33
using System.Collections.Generic;
4+
using System.Runtime.CompilerServices;
45
using System.Threading;
56
using System.Threading.Tasks;
67
using Infrastructure;
@@ -12,14 +13,14 @@ class EndpointSettingsStore(IRavenSessionProvider sessionProvider) : IEndpointSe
1213
static string MakeDocumentId(string name) =>
1314
$"{EndpointSettings.CollectionName}/{DeterministicGuid.MakeId(name)}";
1415

15-
public async IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings()
16+
public async IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings([EnumeratorCancellation] CancellationToken cancellationToken)
1617
{
17-
using IAsyncDocumentSession session = await sessionProvider.OpenSession();
18+
using IAsyncDocumentSession session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
1819
await using IAsyncEnumerator<StreamResult<EndpointSettings>> enumerator = await session
1920
.Advanced
20-
.StreamAsync<EndpointSettings>($"{EndpointSettings.CollectionName}/");
21+
.StreamAsync<EndpointSettings>($"{EndpointSettings.CollectionName}/", token: cancellationToken);
2122

22-
while (await enumerator.MoveNextAsync())
23+
while (await enumerator.MoveNextAsync() && !cancellationToken.IsCancellationRequested)
2324
{
2425
yield return enumerator.Current.Document;
2526
}

src/ServiceControl.Persistence/IEndpointSettingsStore.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66

77
public interface IEndpointSettingsStore
88
{
9-
IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings();
9+
IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings(CancellationToken token);
1010

1111
Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token);
1212
Task Delete(string name, CancellationToken cancellationToken);

src/ServiceControl.UnitTests/Monitoring/HeartbeatEndpointSettingsSyncHostedServiceTests.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -242,7 +242,7 @@ public void RecordHeartbeat(EndpointInstanceId endpointInstanceId, DateTime time
242242

243243
class MockEndpointSettingsStore(EndpointSettings[] settings) : IEndpointSettingsStore
244244
{
245-
public IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings() => settings.ToAsyncEnumerable();
245+
public IAsyncEnumerable<EndpointSettings> GetAllEndpointSettings(CancellationToken token) => settings.ToAsyncEnumerable();
246246

247247
public Task UpdateEndpointSettings(EndpointSettings settings, CancellationToken token)
248248
{

src/ServiceControl/Monitoring/HeartbeatEndpointSettingsSyncHostedService.cs

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -67,8 +67,7 @@ async Task PurgeMonitoringDataThatDoesNotNeedToBeTracked(CancellationToken cance
6767
ILookup<string, Guid> monitorEndpointsLookup = endpointsViews
6868
.Where(view => !view.IsSendingHeartbeats)
6969
.ToLookup(view => view.Name, view => view.Id);
70-
await foreach (EndpointSettings endpointSetting in endpointSettingsStore.GetAllEndpointSettings()
71-
.WithCancellation(cancellationToken))
70+
await foreach (EndpointSettings endpointSetting in endpointSettingsStore.GetAllEndpointSettings(cancellationToken))
7271
{
7372
if (!endpointSetting.TrackInstances)
7473
{
@@ -94,8 +93,7 @@ async Task InitialiseSettings(HashSet<string> monitorEndpoints, CancellationToke
9493
HashSet<string> settingsNames = [];
9594

9695
// Delete any endpoints data that no longer exists
97-
await foreach (EndpointSettings endpointSetting in endpointSettingsStore.GetAllEndpointSettings()
98-
.WithCancellation(cancellationToken))
96+
await foreach (EndpointSettings endpointSetting in endpointSettingsStore.GetAllEndpointSettings(cancellationToken))
9997
{
10098
if (endpointSetting.Name == string.Empty)
10199
{

src/ServiceControl/Monitoring/Web/EndpointsSettingsController.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ public class EndpointsSettingsController(
3333
public async IAsyncEnumerable<SettingsData> Endpoints([EnumeratorCancellation] CancellationToken token)
3434
{
3535
await using IAsyncEnumerator<EndpointSettings> enumerator =
36-
dataStore.GetAllEndpointSettings().GetAsyncEnumerator(token);
36+
dataStore.GetAllEndpointSettings(token).GetAsyncEnumerator(token);
3737
bool noResults = true;
3838
while (await enumerator.MoveNextAsync())
3939
{

0 commit comments

Comments
 (0)