Skip to content

Commit c0b41e8

Browse files
committed
feat: add cancellation-aware RPC contracts
1 parent c9649c3 commit c0b41e8

185 files changed

Lines changed: 8007 additions & 1262 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
Lines changed: 60 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
using System.Collections.Generic;
22
using System.Linq;
3+
using System.Threading;
34
using System.Threading.Tasks;
45
using Orleans.Concurrency;
56
using Orleans.Runtime;
@@ -14,40 +15,77 @@ internal sealed class DashboardClient(IGrainFactory grainFactory) : IDashboardCl
1415
private readonly IDashboardRemindersGrain _remindersGrain = grainFactory.GetGrain<IDashboardRemindersGrain>(0);
1516
private readonly IGrainFactory _grainFactory = grainFactory;
1617

17-
public async Task<Immutable<DashboardCounters>> DashboardCounters(string[]? exclusions) => await _dashboardGrain.GetCounters(exclusions);
18-
19-
public async Task<Immutable<Dictionary<string, GrainTraceEntry>>> ClusterStats() => await _dashboardGrain.GetClusterTracing();
20-
21-
public async Task<Immutable<ReminderResponse>> GetReminders(int pageNumber, int pageSize) => await _remindersGrain.GetReminders(pageNumber, pageSize);
22-
23-
public async Task<Immutable<SiloRuntimeStatistics?[]>> HistoricalStats(string siloAddress) => await Silo(siloAddress).GetRuntimeStatistics();
24-
25-
public async Task<Immutable<Dictionary<string, string?>>> SiloProperties(string siloAddress) => await Silo(siloAddress).GetExtendedProperties();
26-
27-
public async Task<Immutable<Dictionary<string, string>>> SiloMetadata(string siloAddress) => await Silo(siloAddress).GetMetadata();
28-
29-
public async Task<Immutable<Dictionary<string, GrainTraceEntry>>> SiloStats(string siloAddress) => await _dashboardGrain.GetSiloTracing(siloAddress);
30-
31-
public async Task<Immutable<StatCounter[]>> GetCounters(string siloAddress) => await Silo(siloAddress).GetCounters();
18+
public async Task<Immutable<DashboardCounters>> DashboardCounters(
19+
string[]? exclusions,
20+
CancellationToken cancellationToken = default) =>
21+
await _dashboardGrain.GetCounters(exclusions, cancellationToken);
22+
23+
public async Task<Immutable<Dictionary<string, GrainTraceEntry>>> ClusterStats(
24+
CancellationToken cancellationToken = default) =>
25+
await _dashboardGrain.GetClusterTracing(cancellationToken);
26+
27+
public async Task<Immutable<ReminderResponse>> GetReminders(
28+
int pageNumber,
29+
int pageSize,
30+
CancellationToken cancellationToken = default) =>
31+
await _remindersGrain.GetReminders(pageNumber, pageSize, cancellationToken);
32+
33+
public async Task<Immutable<SiloRuntimeStatistics?[]>> HistoricalStats(
34+
string siloAddress,
35+
CancellationToken cancellationToken = default) =>
36+
await Silo(siloAddress).GetRuntimeStatistics(cancellationToken);
37+
38+
public async Task<Immutable<Dictionary<string, string?>>> SiloProperties(
39+
string siloAddress,
40+
CancellationToken cancellationToken = default) =>
41+
await Silo(siloAddress).GetExtendedProperties(cancellationToken);
42+
43+
public async Task<Immutable<Dictionary<string, string>>> SiloMetadata(
44+
string siloAddress,
45+
CancellationToken cancellationToken = default) =>
46+
await Silo(siloAddress).GetMetadata(cancellationToken);
47+
48+
public async Task<Immutable<Dictionary<string, GrainTraceEntry>>> SiloStats(
49+
string siloAddress,
50+
CancellationToken cancellationToken = default) =>
51+
await _dashboardGrain.GetSiloTracing(siloAddress, cancellationToken);
52+
53+
public async Task<Immutable<StatCounter[]>> GetCounters(
54+
string siloAddress,
55+
CancellationToken cancellationToken = default) =>
56+
await Silo(siloAddress).GetCounters(cancellationToken);
3257

3358
public async Task<Immutable<Dictionary<string, Dictionary<string, GrainTraceEntry>>>> GrainStats(
34-
string grainName) => await _dashboardGrain.GetGrainTracing(grainName);
59+
string grainName,
60+
CancellationToken cancellationToken = default) =>
61+
await _dashboardGrain.GetGrainTracing(grainName, cancellationToken);
3562

36-
public async Task<Immutable<Dictionary<string, GrainMethodAggregate[]>>> TopGrainMethods(int take, string[]? exclusions) => await _dashboardGrain.TopGrainMethods(take, exclusions);
63+
public async Task<Immutable<Dictionary<string, GrainMethodAggregate[]>>> TopGrainMethods(
64+
int take,
65+
string[]? exclusions,
66+
CancellationToken cancellationToken = default) =>
67+
await _dashboardGrain.TopGrainMethods(take, exclusions, cancellationToken);
3768

3869
private ISiloGrainProxy Silo(string siloAddress) => _grainFactory.GetGrain<ISiloGrainProxy>(siloAddress);
3970

40-
public async Task<Immutable<string>> GetGrainState(string? id, string? grainType) => await _dashboardGrain.GetGrainState(id, grainType);
71+
public async Task<Immutable<string>> GetGrainState(
72+
string? id,
73+
string? grainType,
74+
CancellationToken cancellationToken = default) =>
75+
await _dashboardGrain.GetGrainState(id, grainType, cancellationToken);
4176

42-
public async Task<Immutable<string[]>> GetGrainTypes(string[]? exclusions = null) => await _dashboardGrain.GetGrainTypes(exclusions);
77+
public async Task<Immutable<string[]>> GetGrainTypes(
78+
string[]? exclusions = null,
79+
CancellationToken cancellationToken = default) =>
80+
await _dashboardGrain.GetGrainTypes(exclusions, cancellationToken);
4381

44-
public async Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages()
82+
public async Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages(CancellationToken cancellationToken = default)
4583
{
4684
// All silos run an identical lifecycle, so we only need to ask one of them.
4785
// Use the management grain to find an active host, then call its dashboard
4886
// silo grain proxy.
4987
var management = _grainFactory.GetGrain<IManagementGrain>(0);
50-
var hosts = await management.GetHosts(onlyActive: true);
88+
var hosts = await management.GetHosts(onlyActive: true, cancellationToken);
5189
var siloAddress = hosts
5290
.Where(x => x.Value == SiloStatus.Active)
5391
.Select(x => x.Key)
@@ -57,6 +95,6 @@ public async Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages()
5795
return new LifecycleStageInfo[0].AsImmutable();
5896
}
5997

60-
return await Silo(siloAddress.ToParsableString()).GetLifecycleStages();
98+
return await Silo(siloAddress.ToParsableString()).GetLifecycleStages(cancellationToken);
6199
}
62100
}
Lines changed: 39 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
using System.Collections.Generic;
2+
using System.Threading;
23
using System.Threading.Tasks;
34
using Orleans.Concurrency;
45
using Orleans.Runtime;
@@ -9,29 +10,54 @@ namespace Orleans.Dashboard.Core;
910

1011
internal interface IDashboardClient
1112
{
12-
Task<Immutable<DashboardCounters>> DashboardCounters(string[]? exclusions = null);
13+
Task<Immutable<DashboardCounters>> DashboardCounters(
14+
string[]? exclusions = null,
15+
CancellationToken cancellationToken = default);
1316

14-
Task<Immutable<Dictionary<string, GrainTraceEntry>>> ClusterStats();
17+
Task<Immutable<Dictionary<string, GrainTraceEntry>>> ClusterStats(CancellationToken cancellationToken = default);
1518

16-
Task<Immutable<ReminderResponse>> GetReminders(int pageNumber, int pageSize);
19+
Task<Immutable<ReminderResponse>> GetReminders(
20+
int pageNumber,
21+
int pageSize,
22+
CancellationToken cancellationToken = default);
1723

18-
Task<Immutable<SiloRuntimeStatistics?[]>> HistoricalStats(string siloAddress);
24+
Task<Immutable<SiloRuntimeStatistics?[]>> HistoricalStats(
25+
string siloAddress,
26+
CancellationToken cancellationToken = default);
1927

20-
Task<Immutable<Dictionary<string, string?>>> SiloProperties(string siloAddress);
28+
Task<Immutable<Dictionary<string, string?>>> SiloProperties(
29+
string siloAddress,
30+
CancellationToken cancellationToken = default);
2131

22-
Task<Immutable<Dictionary<string, string>>> SiloMetadata(string siloAddress);
32+
Task<Immutable<Dictionary<string, string>>> SiloMetadata(
33+
string siloAddress,
34+
CancellationToken cancellationToken = default);
2335

24-
Task<Immutable<Dictionary<string, GrainTraceEntry>>> SiloStats(string siloAddress);
36+
Task<Immutable<Dictionary<string, GrainTraceEntry>>> SiloStats(
37+
string siloAddress,
38+
CancellationToken cancellationToken = default);
2539

26-
Task<Immutable<StatCounter[]>> GetCounters(string siloAddress);
40+
Task<Immutable<StatCounter[]>> GetCounters(
41+
string siloAddress,
42+
CancellationToken cancellationToken = default);
2743

28-
Task<Immutable<Dictionary<string, Dictionary<string, GrainTraceEntry>>>> GrainStats(string grainName);
44+
Task<Immutable<Dictionary<string, Dictionary<string, GrainTraceEntry>>>> GrainStats(
45+
string grainName,
46+
CancellationToken cancellationToken = default);
2947

30-
Task<Immutable<Dictionary<string, GrainMethodAggregate[]>>> TopGrainMethods(int take, string[]? exclusions = null);
48+
Task<Immutable<Dictionary<string, GrainMethodAggregate[]>>> TopGrainMethods(
49+
int take,
50+
string[]? exclusions = null,
51+
CancellationToken cancellationToken = default);
3152

32-
Task<Immutable<string>> GetGrainState(string? id, string? grainType);
53+
Task<Immutable<string>> GetGrainState(
54+
string? id,
55+
string? grainType,
56+
CancellationToken cancellationToken = default);
3357

34-
Task<Immutable<string[]>> GetGrainTypes(string[]? exclusions = null);
58+
Task<Immutable<string[]>> GetGrainTypes(
59+
string[]? exclusions = null,
60+
CancellationToken cancellationToken = default);
3561

36-
Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages();
62+
Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages(CancellationToken cancellationToken = default);
3763
}
Lines changed: 27 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
using System.Collections.Generic;
2+
using System.Threading;
23
using System.Threading.Tasks;
34
using Orleans.Concurrency;
45
using Orleans.Dashboard.Model;
@@ -11,30 +12,47 @@ internal interface IDashboardGrain : IGrainWithIntegerKey
1112
{
1213
[OneWay]
1314
[Alias("InitializeAsync")]
14-
Task InitializeAsync();
15+
Task InitializeAsync(CancellationToken cancellationToken = default);
1516

1617
[OneWay]
1718
[Alias("SubmitTracing")]
18-
Task SubmitTracing(string siloAddress, Immutable<SiloGrainTraceEntry[]> grainCallTime);
19+
Task SubmitTracing(
20+
string siloAddress,
21+
Immutable<SiloGrainTraceEntry[]> grainCallTime,
22+
CancellationToken cancellationToken = default);
1923

2024
[Alias("GetCounters")]
21-
Task<Immutable<DashboardCounters>> GetCounters(string[]? exclusions = null);
25+
Task<Immutable<DashboardCounters>> GetCounters(
26+
string[]? exclusions = null,
27+
CancellationToken cancellationToken = default);
2228

2329
[Alias("GetGrainTracing")]
24-
Task<Immutable<Dictionary<string, Dictionary<string, GrainTraceEntry>>>> GetGrainTracing(string grain);
30+
Task<Immutable<Dictionary<string, Dictionary<string, GrainTraceEntry>>>> GetGrainTracing(
31+
string grain,
32+
CancellationToken cancellationToken = default);
2533

2634
[Alias("GetClusterTracing")]
27-
Task<Immutable<Dictionary<string, GrainTraceEntry>>> GetClusterTracing();
35+
Task<Immutable<Dictionary<string, GrainTraceEntry>>> GetClusterTracing(CancellationToken cancellationToken = default);
2836

2937
[Alias("GetSiloTracing")]
30-
Task<Immutable<Dictionary<string, GrainTraceEntry>>> GetSiloTracing(string address);
38+
Task<Immutable<Dictionary<string, GrainTraceEntry>>> GetSiloTracing(
39+
string address,
40+
CancellationToken cancellationToken = default);
3141

3242
[Alias("TopGrainMethods")]
33-
Task<Immutable<Dictionary<string, GrainMethodAggregate[]>>> TopGrainMethods(int take, string[]? exclusions = null);
43+
Task<Immutable<Dictionary<string, GrainMethodAggregate[]>>> TopGrainMethods(
44+
int take,
45+
string[]? exclusions = null,
46+
CancellationToken cancellationToken = default);
3447

3548
[Alias("GetGrainState")]
36-
Task<Immutable<string>> GetGrainState(string? id, string? grainType);
49+
Task<Immutable<string>> GetGrainState(
50+
string? id,
51+
string? grainType,
52+
CancellationToken cancellationToken = default);
3753

3854
[Alias("GetGrainTypes")]
39-
Task<Immutable<string[]>> GetGrainTypes(string[]? exclusions = null);
55+
Task<Immutable<string[]>> GetGrainTypes(
56+
string[]? exclusions = null,
57+
CancellationToken cancellationToken = default);
4058
}

src/Dashboard/Orleans.Dashboard/Core/IDashboardRemindersGrain.cs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
using System.Threading;
12
using System.Threading.Tasks;
23
using Orleans.Concurrency;
34
using Orleans.Dashboard.Model;
@@ -8,5 +9,8 @@ namespace Orleans.Dashboard.Core;
89
internal interface IDashboardRemindersGrain : IGrainWithIntegerKey
910
{
1011
[Alias("GetReminders")]
11-
Task<Immutable<ReminderResponse>> GetReminders(int pageNumber, int pageSize);
12+
Task<Immutable<ReminderResponse>> GetReminders(
13+
int pageNumber,
14+
int pageSize,
15+
CancellationToken cancellationToken = default);
1216
}
Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
using Orleans.Concurrency;
2+
using System.Threading;
23

34
namespace Orleans.Dashboard.Core;
45

@@ -7,5 +8,5 @@ internal interface ISiloGrainProxy : IGrainWithStringKey, ISiloGrainService
78
{
89

910
[Alias("GetMetadata")]
10-
Task<Immutable<Dictionary<string, string>>> GetMetadata();
11+
Task<Immutable<Dictionary<string, string>>> GetMetadata(CancellationToken cancellationToken = default);
1112
}
Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
using System.Collections.Generic;
2+
using System.Threading;
23
using System.Threading.Tasks;
34
using Orleans.Concurrency;
45
using Orleans.Runtime;
@@ -11,24 +12,24 @@ namespace Orleans.Dashboard.Core;
1112
internal interface ISiloGrainService : IGrainService
1213
{
1314
[Alias("SetVersion")]
14-
Task SetVersion(string orleans, string host);
15+
Task SetVersion(string orleans, string host, CancellationToken cancellationToken = default);
1516

1617
[OneWay]
1718
[Alias("ReportCounters")]
18-
Task ReportCounters(Immutable<StatCounter[]> stats);
19+
Task ReportCounters(Immutable<StatCounter[]> stats, CancellationToken cancellationToken = default);
1920

2021
[Alias("Enable")]
21-
Task Enable(bool enabled);
22+
Task Enable(bool enabled, CancellationToken cancellationToken = default);
2223

2324
[Alias("GetExtendedProperties")]
24-
Task<Immutable<Dictionary<string, string?>>> GetExtendedProperties();
25+
Task<Immutable<Dictionary<string, string?>>> GetExtendedProperties(CancellationToken cancellationToken = default);
2526

2627
[Alias("GetRuntimeStatistics")]
27-
Task<Immutable<SiloRuntimeStatistics?[]>> GetRuntimeStatistics();
28+
Task<Immutable<SiloRuntimeStatistics?[]>> GetRuntimeStatistics(CancellationToken cancellationToken = default);
2829

2930
[Alias("GetCounters")]
30-
Task<Immutable<StatCounter[]>> GetCounters();
31+
Task<Immutable<StatCounter[]>> GetCounters(CancellationToken cancellationToken = default);
3132

3233
[Alias("GetLifecycleStages")]
33-
Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages();
34+
Task<Immutable<LifecycleStageInfo[]>> GetLifecycleStages(CancellationToken cancellationToken = default);
3435
}

src/Dashboard/Orleans.Dashboard/DashboardHost.cs

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -24,34 +24,34 @@ internal sealed partial class DashboardHost(
2424
public async Task StartAsync(CancellationToken cancellationToken)
2525
{
2626
await Task.WhenAll(
27-
ActivateDashboardGrainAsync(),
28-
ActivateSiloGrainAsync(),
27+
ActivateDashboardGrainAsync(cancellationToken),
28+
ActivateSiloGrainAsync(cancellationToken),
2929
StartOpenTelemetryConsumerAsync()).ConfigureAwait(false);
3030
}
3131

3232
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
3333

34-
private async Task ActivateSiloGrainAsync()
34+
private async Task ActivateSiloGrainAsync(CancellationToken cancellationToken)
3535
{
3636
try
3737
{
3838
var siloGrain = siloGrainClient.GrainService(localSiloDetails.SiloAddress);
39-
await siloGrain.SetVersion(GetOrleansVersion(), GetHostVersion()).ConfigureAwait(false);
39+
await siloGrain.SetVersion(GetOrleansVersion(), GetHostVersion(), cancellationToken).ConfigureAwait(false);
4040
}
41-
catch (Exception ex)
41+
catch (Exception ex) when (!cancellationToken.IsCancellationRequested)
4242
{
4343
LogWarningActivateSiloGrainServiceStartupFailed(logger, ex);
4444
}
4545
}
4646

47-
private async Task ActivateDashboardGrainAsync()
47+
private async Task ActivateDashboardGrainAsync(CancellationToken cancellationToken)
4848
{
4949
try
5050
{
5151
var dashboardGrain = grainFactory.GetGrain<IDashboardGrain>(0);
52-
await dashboardGrain.InitializeAsync().ConfigureAwait(false);
52+
await dashboardGrain.InitializeAsync(cancellationToken).ConfigureAwait(false);
5353
}
54-
catch (Exception ex)
54+
catch (Exception ex) when (!cancellationToken.IsCancellationRequested)
5555
{
5656
LogWarningActivateDashboardGrainStartupFailed(logger, ex);
5757
}

0 commit comments

Comments
 (0)