Skip to content

Commit d16c110

Browse files
committed
fix(runtime): harden dissemination convergence
1 parent bee0660 commit d16c110

9 files changed

Lines changed: 267 additions & 27 deletions

File tree

src/Orleans.Runtime/Dissemination/DisseminationBroadcastQueue.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -227,9 +227,9 @@ private TimeSpan GetRetryDelay(int attempt)
227227
}
228228

229229
var cap = _options.CurrentValue.Overlay.AntiEntropyInterval;
230-
if (cap < floor)
230+
if (floor > cap)
231231
{
232-
cap = floor;
232+
floor = cap;
233233
}
234234

235235
var multiplier = Math.Pow(2, Math.Min(Math.Max(0, attempt - 1), 20));

src/Orleans.Runtime/Dissemination/DisseminationProtocol.cs

Lines changed: 44 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ namespace Orleans.Runtime.Dissemination;
99
// The protocol coordinates routing and application while namespaces remain authoritative for values and repair history.
1010
internal sealed partial class DisseminationProtocol
1111
{
12+
private const int MaxRetainedNonMemberResponseCursors = 64;
1213
private static readonly TimeSpan MaxAntiEntropyRoundLifetime = TimeSpan.FromMilliseconds(uint.MaxValue - 1);
1314
private readonly SiloAddress _localSilo;
1415
private readonly IInternalGrainFactory _grainFactory;
@@ -18,7 +19,8 @@ internal sealed partial class DisseminationProtocol
1819
private readonly ILogger<DisseminationProtocol> _logger;
1920
private readonly DisseminationBroadcastQueue _broadcastQueue;
2021
private readonly object _antiEntropyResponseCursorLock = new();
21-
private readonly Dictionary<SiloAddress, int> _antiEntropyResponseCursors = [];
22+
private readonly Dictionary<SiloAddress, AntiEntropyResponseCursor> _antiEntropyResponseCursors = [];
23+
private long _antiEntropyResponseCursorAccess;
2224
private readonly object _valueUpdateLock = new();
2325
private readonly Dictionary<DigestKey, ValueUpdate> _lastValueUpdates = [];
2426
private readonly FrozenDictionary<DisseminationNamespace, IDisseminationNamespace> _namespaces;
@@ -422,7 +424,6 @@ public ValueTask<DisseminationAntiEntropyResponse> ReceiveAntiEntropy(
422424
CancellationToken cancellationToken)
423425
{
424426
cancellationToken.ThrowIfCancellationRequested();
425-
PruneAntiEntropyResponseCursors(_membership.CurrentSnapshot);
426427
// Incoming digests are passive evidence for existing peer pumps, not a reason to create new ones.
427428
foreach (var (namespaceName, entries) in request.Digests)
428429
{
@@ -437,7 +438,9 @@ public ValueTask<DisseminationAntiEntropyResponse> ReceiveAntiEntropy(
437438
}
438439

439440
var options = _options.CurrentValue;
440-
return new(CreateAntiEntropyResponse(request, options, cancellationToken));
441+
var response = CreateAntiEntropyResponse(request, options, cancellationToken);
442+
PruneAntiEntropyResponseCursors(_membership.CurrentSnapshot, request.Sender);
443+
return new(response);
441444
}
442445

443446
private DisseminationAntiEntropyResponse CreateAntiEntropyResponse(
@@ -596,9 +599,16 @@ private int GetAntiEntropyResponseStart(SiloAddress peer, int candidateCount)
596599

597600
lock (_antiEntropyResponseCursorLock)
598601
{
599-
return _antiEntropyResponseCursors.TryGetValue(peer, out var cursor)
600-
? cursor % candidateCount
601-
: 0;
602+
if (!_antiEntropyResponseCursors.TryGetValue(peer, out var cursor))
603+
{
604+
return 0;
605+
}
606+
607+
_antiEntropyResponseCursors[peer] = cursor with
608+
{
609+
LastAccess = ++_antiEntropyResponseCursorAccess,
610+
};
611+
return cursor.Position % candidateCount;
602612
}
603613
}
604614

@@ -616,7 +626,9 @@ private void AdvanceAntiEntropyResponseCursor(
616626
}
617627
else
618628
{
619-
_antiEntropyResponseCursors[peer] = (start + Math.Max(1, examined)) % candidateCount;
629+
_antiEntropyResponseCursors[peer] = new(
630+
(start + Math.Max(1, examined)) % candidateCount,
631+
++_antiEntropyResponseCursorAccess);
620632
}
621633
}
622634
}
@@ -629,15 +641,35 @@ private void ClearAntiEntropyResponseCursor(SiloAddress peer)
629641
}
630642
}
631643

632-
private void PruneAntiEntropyResponseCursors(DisseminationMembershipSnapshot membership)
644+
private void PruneAntiEntropyResponseCursors(
645+
DisseminationMembershipSnapshot membership,
646+
SiloAddress currentRequester)
633647
{
634648
lock (_antiEntropyResponseCursorLock)
635649
{
636-
foreach (var peer in _antiEntropyResponseCursors.Keys.ToArray())
650+
var nonMemberCount = 0;
651+
foreach (var peer in _antiEntropyResponseCursors.Keys)
637652
{
638653
if (!membership.ContainsMember(peer))
639654
{
640-
_antiEntropyResponseCursors.Remove(peer);
655+
nonMemberCount++;
656+
}
657+
}
658+
659+
if (nonMemberCount <= MaxRetainedNonMemberResponseCursors)
660+
{
661+
return;
662+
}
663+
664+
foreach (var cursor in _antiEntropyResponseCursors
665+
.Where(entry => !Equals(entry.Key, currentRequester) && !membership.ContainsMember(entry.Key))
666+
.OrderBy(static entry => entry.Value.LastAccess)
667+
.ToArray())
668+
{
669+
_antiEntropyResponseCursors.Remove(cursor.Key);
670+
if (--nonMemberCount <= MaxRetainedNonMemberResponseCursors)
671+
{
672+
break;
641673
}
642674
}
643675
}
@@ -960,6 +992,8 @@ private static int CompareAntiEntropyRepairs(AntiEntropyRepair left, AntiEntropy
960992

961993
private readonly record struct ValueUpdate(long Version, long Timestamp);
962994

995+
private readonly record struct AntiEntropyResponseCursor(int Position, long LastAccess);
996+
963997
private readonly record struct AntiEntropyRepair(
964998
IDisseminationNamespace Namespace,
965999
List<DisseminationBroadcastValue> Items,

src/Orleans.Runtime/Dissemination/WakeTimer.cs

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -148,11 +148,19 @@ public void Dispose()
148148
}
149149

150150
_disposed = true;
151-
waiter = CompleteWaitUnsafe();
151+
waiter = _waiter;
152+
_waiter = null;
153+
_armed = false;
152154
}
153155

154-
waiter?.TrySetResult(false);
155-
_timer.Dispose();
156+
try
157+
{
158+
waiter?.TrySetResult(false);
159+
}
160+
finally
161+
{
162+
_timer.Dispose();
163+
}
156164
}
157165

158166
private TaskCompletionSource<bool>? CompleteWaitUnsafe()

src/Orleans.Runtime/Manifest/ClusterManifestProvider.cs

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -210,16 +210,23 @@ private async Task<bool> UpdateManifest(ClusterMembershipSnapshot clusterMembers
210210
if (await peerRepairTask)
211211
{
212212
modified = true;
213+
var repairedSilos = builder.ToImmutable();
214+
PruneManifestCache(repairedSilos);
215+
var repairedManifest = CreateClusterManifest(
216+
new MajorMinorVersion(clusterMembership.Version.Value, existingManifest.Version.Minor + 1),
217+
repairedSilos);
218+
if (!TryPublishManifest(repairedManifest))
219+
{
220+
return false;
221+
}
222+
223+
existingManifest = repairedManifest;
224+
modified = false;
213225
if (missingSilos.All(builder.ContainsKey))
214226
{
215-
// Peer repair already supplied every missing manifest. Publish immediately while redundant
216-
// direct fetches continue independently, so a hung direct target cannot delay convergence.
217-
var repairedSilos = builder.ToImmutable();
218-
PruneManifestCache(repairedSilos);
219-
var repairedManifest = CreateClusterManifest(
220-
new MajorMinorVersion(clusterMembership.Version.Value, existingManifest.Version.Minor + 1),
221-
repairedSilos);
222-
return TryPublishManifest(repairedManifest);
227+
// Peer repair already supplied every missing manifest. Redundant direct fetches can continue
228+
// independently without delaying convergence.
229+
return true;
223230
}
224231
}
225232

src/Orleans.Runtime/MembershipService/MembershipGossiper.cs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,9 @@ private async Task TryGossipViaDissemination(MembershipTableSnapshot snapshot, C
5555
return;
5656
}
5757

58-
await disseminationNamespace.PublishAsync(dissemination, snapshot, cancellationToken);
58+
await disseminationNamespace.PublishAsync(dissemination, snapshot, cancellationToken)
59+
.AsTask()
60+
.WaitAsync(cancellationToken);
5961
}
6062
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
6163
{

test/Orleans.Core.Tests/Manifest/ClusterManifestProviderTests.cs

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -630,6 +630,71 @@ public async Task PeerRepair_HungPeersAndHealthyLaterPeer_UsesAtMostThreeConcurr
630630
}
631631
}
632632

633+
[Fact]
634+
public async Task PeerRepair_PartialResult_PublishesBeforeHungDirectFetchesComplete()
635+
{
636+
var localSilo = CreateSiloAddress(11111, 1);
637+
var peers = Enumerable.Range(11112, 3).Select(port => CreateSiloAddress(port, 1)).OrderBy(static address => address).ToArray();
638+
var repairedPeer = peers[0];
639+
var remoteManifest = CreateGrainManifest();
640+
var remoteHash = ManifestHashCalculator.ComputeHash(remoteManifest);
641+
var summary = new ClusterManifestHashSummary(
642+
new MajorMinorVersion(1, 1),
643+
new Dictionary<SiloAddress, ManifestHash> { [repairedPeer] = remoteHash });
644+
var update = new ClusterManifestUpdate(
645+
new MajorMinorVersion(1, 1),
646+
ImmutableDictionary<SiloAddress, GrainManifest>.Empty.Add(repairedPeer, remoteManifest),
647+
includesAllActiveServers: false);
648+
var pendingDirectFetch = new TaskCompletionSource<GrainManifest>(TaskCreationOptions.RunContinuationsAsynchronously);
649+
var requestLog = new ManifestRequestLog(expectedProbeCount: peers.Length, expectedLegacyFetchCount: peers.Length);
650+
var targets = peers.ToDictionary(
651+
peer => peer,
652+
peer => new TestClusterManifestSystemTarget(
653+
getHashSummary: () =>
654+
{
655+
requestLog.RecordProbe(peer);
656+
return Task.FromResult(summary);
657+
},
658+
getUpdate: _ => Task.FromResult<ClusterManifestUpdate?>(update),
659+
getLegacyManifest: () =>
660+
{
661+
requestLog.RecordLegacyFetch(peer);
662+
return pendingDirectFetch.Task;
663+
}));
664+
var grainFactory = CreateGrainFactory(targets);
665+
var membership = new TestClusterMembershipService(CreateActiveMembershipSnapshot(1, localSilo, peers));
666+
var provider = CreateClusterManifestProvider(
667+
localSilo,
668+
membership,
669+
grainFactory,
670+
new FakeTimeProvider(),
671+
NullLogger<ClusterManifestProvider>.Instance);
672+
var repairedManifest = ObserveManifestAsync(provider, new MajorMinorVersion(1, 1));
673+
var lifecycle = await StartAsync(provider);
674+
675+
try
676+
{
677+
await Task.WhenAll(
678+
requestLog.WaitForProbeCountAsync(peers.Length),
679+
requestLog.WaitForLegacyFetchCountAsync(peers.Length));
680+
681+
var repaired = await repairedManifest.WaitAsync(
682+
TimeSpan.FromSeconds(5),
683+
TestContext.Current.CancellationToken);
684+
685+
Assert.Equal(remoteManifest, repaired.Silos[repairedPeer]);
686+
Assert.Contains(localSilo, repaired.Silos.Keys);
687+
Assert.DoesNotContain(peers.Skip(1), repaired.Silos.Keys.Contains);
688+
Assert.False(pendingDirectFetch.Task.IsCompleted);
689+
}
690+
finally
691+
{
692+
await lifecycle.OnStop(TestContext.Current.CancellationToken);
693+
provider.Dispose();
694+
membership.Dispose();
695+
}
696+
}
697+
633698
[Fact]
634699
public async Task PeerRepair_StopCancellation_CompletesHungProbeProcessing()
635700
{

test/Orleans.Core.Tests/Membership/MembershipGossiperTests.cs

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -171,6 +171,34 @@ public async Task GossipToRemoteSilos_CallerCancellation_CancelsDirectWrapperAnd
171171
}
172172
}
173173

174+
[Fact]
175+
public async Task GossipToRemoteSilos_CallerCancellationBoundsUnresponsiveDisseminationPublish()
176+
{
177+
var disseminationStarted = NewBarrier();
178+
var disseminationCompletion = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously);
179+
using var cancellation = new CancellationTokenSource();
180+
using var rig = CreateTestRig(
181+
_ => Task.CompletedTask,
182+
(_, _, _) =>
183+
{
184+
disseminationStarted.TrySetResult();
185+
return new ValueTask<bool>(disseminationCompletion.Task);
186+
});
187+
var snapshot = CreateSnapshot(rig.LocalSilo, rig.RemoteSilo, SiloStatus.Active);
188+
var gossipTask = rig.Gossiper.GossipToRemoteSilos(
189+
[rig.RemoteSilo],
190+
snapshot,
191+
rig.LocalSilo,
192+
SiloStatus.Stopping,
193+
cancellation.Token);
194+
195+
await disseminationStarted.Task;
196+
cancellation.Cancel();
197+
198+
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => gossipTask);
199+
Assert.False(disseminationCompletion.Task.IsCompleted);
200+
}
201+
174202
private static MembershipGossiperTestRig CreateTestRig(
175203
Func<MembershipTableSnapshot, Task> directGossip,
176204
Func<DisseminationKey, long, CancellationToken, ValueTask<bool>> disseminationPublish)

test/Orleans.Runtime.Internal.Tests/Dissemination/DisseminationProtocolTests.cs

Lines changed: 62 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -434,6 +434,59 @@ public async Task HighPriorityNamespaceFlushesWithoutWaitingForCoalescingWindow(
434434
await protocol.StopAsync(TestContext.Current.CancellationToken);
435435
}
436436

437+
[Fact]
438+
public async Task HighPriorityRetryIsCappedByAntiEntropyInterval()
439+
{
440+
var local = CreateSilo(11111);
441+
var peer = CreateSilo(11112);
442+
var transport = new FakeTransport(local, peer);
443+
var timeProvider = new FakeTimeProvider();
444+
var firstAttempt = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
445+
var secondAttempt = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
446+
var sendCount = 0;
447+
transport.SendBroadcastHandler = (target, batch, cancellationToken) =>
448+
{
449+
if (Interlocked.Increment(ref sendCount) == 1)
450+
{
451+
firstAttempt.TrySetResult();
452+
throw new InvalidOperationException("transient send failure");
453+
}
454+
455+
transport.BroadcastBatches.Add((target, batch));
456+
secondAttempt.TrySetResult();
457+
return Task.CompletedTask;
458+
};
459+
460+
var normalNamespace = new FakeNamespace(local, new DisseminationNamespace("normal"));
461+
normalNamespace.Options.MaxCoalescingDelay = TimeSpan.FromMinutes(1);
462+
var highNamespace = new FakeNamespace(local, new DisseminationNamespace("high"));
463+
highNamespace.Options.Priority = DisseminationPriority.High;
464+
var protocol = CreateProtocol(
465+
transport,
466+
new IDisseminationNamespace[] { normalNamespace, highNamespace },
467+
options => options.Overlay.AntiEntropyInterval = TimeSpan.FromSeconds(1),
468+
timeProvider);
469+
using var schedule = new BroadcastScheduleObserver();
470+
471+
Assert.True(await PublishValue(
472+
protocol,
473+
highNamespace,
474+
highNamespace.CreateValue(FakeNamespace.DefaultKey, sequence: 1),
475+
TestContext.Current.CancellationToken));
476+
timeProvider.Advance(TimeSpan.FromMilliseconds(1));
477+
await firstAttempt.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken);
478+
var retry = await schedule.WaitAsync(
479+
e => e.Peer.Equals(peer) && e.Reason == DisseminationBroadcastScheduleReason.Retry,
480+
TimeSpan.FromSeconds(5));
481+
482+
Assert.Equal(TimeSpan.FromSeconds(1), retry.DueTime);
483+
timeProvider.Advance(TimeSpan.FromMilliseconds(999));
484+
Assert.False(secondAttempt.Task.IsCompleted);
485+
timeProvider.Advance(TimeSpan.FromMilliseconds(1));
486+
await secondAttempt.Task.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken);
487+
await protocol.StopAsync(TestContext.Current.CancellationToken);
488+
}
489+
437490
[Fact]
438491
public async Task HighPriorityNotificationPullsPendingCoalescedFlushForward()
439492
{
@@ -2123,8 +2176,9 @@ public async Task AntiEntropyResponseOmitsNamespacesWithoutRepairs()
21232176
public async Task AntiEntropyTruncationRotatesPastContinuouslyAdvancingKey()
21242177
{
21252178
var local = CreateSilo(11111);
2126-
var peer = CreateSilo(11112);
2127-
var transport = new FakeTransport(local, peer);
2179+
var member = CreateSilo(11112);
2180+
var requester = CreateSilo(11113);
2181+
var transport = new FakeTransport(local, member);
21282182
var ns = new FakeNamespace(local);
21292183
DisseminationKey hotKey = "hot";
21302184
DisseminationKey waitingKey = "waiting";
@@ -2136,14 +2190,19 @@ public async Task AntiEntropyTruncationRotatesPastContinuouslyAdvancingKey()
21362190
options => options.MaxBatchItems = 1);
21372191
var request = new DisseminationAntiEntropyRequest
21382192
{
2139-
Sender = peer,
2193+
Sender = requester,
21402194
Digests = CreateAntiEntropyRequestDigest(
21412195
ns.Name,
21422196
(hotKey, 0),
21432197
(waitingKey, 0)),
21442198
};
21452199

21462200
var first = await protocol.ReceiveAntiEntropy(request, TestContext.Current.CancellationToken);
2201+
await protocol.ReceiveAntiEntropy(new DisseminationAntiEntropyRequest
2202+
{
2203+
Sender = member,
2204+
Digests = request.Digests,
2205+
}, TestContext.Current.CancellationToken);
21472206
ns.SetValue(hotKey, version: 2);
21482207
var second = await protocol.ReceiveAntiEntropy(request, TestContext.Current.CancellationToken);
21492208

0 commit comments

Comments
 (0)