You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Deployment-load publication sends each silo's sample to every other silo: N*(N-1) requests per period. Membership dissemination also needs a small ordinary-update path while retaining reliable convergence after missed updates, topology changes, and partitions.
Solution
Add opt-in deterministic dissemination for membership and deployment load, with bounded broadcast queues and anti-entropy repair. The subsystem and both namespace enablement flags default to false; existing direct paths cover bootstrap, mixed-version operation, and failed admission.
Membership broadcasts are sparse deltas; anti-entropy repairs are full snapshots. The broadcast forest includes Joining, Active, ShuttingDown, and Stopping silos, with membership updates ahead of load updates. Each outbound peer/key retains one immutable comparison snapshot captured by an accepted, exactly acknowledged broadcast. After send admission, its next delta compares that snapshot with current owner state, carrying changed entries and removed identities. This preserves coalesced updates and reversions without a generic history or repair-chain protocol.
Receivers apply a delta at its base view, or merge a replay at the target view while preserving canonical fields and maximum heartbeat timestamps. Same-version pruning removes Dead rows only; all other statuses retain their canonical membership. Coalesced deltas across versions can span an intervening Dead transition and cleanup. Missing baselines become eligible for the next full anti-entropy repair. A peer without a comparison snapshot receives a small empty current-view probe; bootstrap and repair establish missing state. Relays compare their accepted current view with each child's baseline, including when the relay is ahead of an incoming update. Initial independent heartbeat and retained-Dead differences remain repair work.
Full repair uses the view version plus a liveness/inventory fingerprint. Same-version repair merges maximum actual IAmAliveTime independently of StartTime, and reconciles pruned Dead rows while preserving non-Dead entries. Application confirms requested effects against the resulting owner snapshot before producing an exact compact acknowledgment. The cached sender comparison snapshot comes from construction time, preserving updates which arrive while an acknowledgment is in flight.
Load publication uses receipt-synchronized cohorts. Non-root silos send one fresh sample to a deterministic root, which contributes locally. A cohort seals when every expected Active incarnation contributes or the publication period expires (one second by default). Held ingress RPCs return the remaining delay to the next cohort boundary, aligning actual sampling timers. Distribution uses fanout eight and forwards admitted sets immediately. Independent ingress admission keeps membership's send slots available.
Item, payload, pending-key, and concurrency limits bound work. Queue notifications retain identities; membership comparison snapshots are shared references and are pruned with peer/key knowledge. A slow set of peers can retain different snapshots, so worst-case comparison memory scales with outbound peers times membership size. Namespace payload caches are constant-size. Local cancellation releases waits while admitted state retains its owner. Shutdown seals cohorts before draining held admissions and then peer queues within the caller's budget.
Scope and rationale
Broadcast creation is separate from full repair construction. Public options focus on enablement, topology, cadence, and resource limits. The independent prerequisites #11294 and #11296 have been merged by humans; this branch is rebased onto their finalized implementations. Provider-owned fatal rollback handling and canonical snapshot merging come from main. Superseded manager reset special-casing and restart-as-recovery tests are removed. Dissemination follows ordinary monotonic membership views.
Broadcast, repair, and cohort integration remain one working convergence feature. Manual performance tooling is separate in #11272. At one sample and one fitting cohort per second, healthy load traffic is approximately 2*(N-1) RPCs/s: 18 at 10 silos, 198 at 100, and 3,998 at 2,000, plus repair and replies. These are projections; root ingress remains concentrated and value delivery remains quadratic. Earlier fixed-window measurements cover a different implementation. Receipt-driven process measurements remain pending separately authorized tooling execution.
See runtime dissemination architecture for invariants, timing, limits, and failure semantics. Old/new-binary compatibility covers rolling upgrade, rollback, partition recovery, isolated full-snapshot/heartbeat repair, cancellation, and bounded shutdown with exact state comparisons.
ReubenBond
changed the title
Add deterministic dissemination broadcast and repair
feat(runtime): add deterministic dissemination broadcast and repair
Jun 19, 2026
The reason will be displayed to describe this comment to others. Learn more.
Pull request overview
This PR introduces a new internal dissemination substrate for monotonically versioned runtime state, using deterministic fixed-tree broadcast for the fast path and periodic anti-entropy repair for convergence. It integrates the substrate into deployment load statistics and membership gossip, and improves manifest convergence by reusing content-addressed manifest hashes/caching and peer-assisted fills.
Wired dissemination into deployment load publishing and membership gossip with opt-in options and legacy fallbacks.
Added manifest hash-based fetch/caching and peer-based fill to reduce manifest convergence request fanout; expanded SiloAddress parsing APIs and added new dissemination-focused tests/docs.
_manifestCache is a plain Dictionary, but UpdateManifest fetches missing manifests via multiple concurrent tasks (Task.WhenAll) which all call GetSiloManifest(). That results in concurrent reads/writes to _manifestCache (TryGetValue, indexer assignment), which is not thread-safe and can corrupt the dictionary or throw at runtime. Consider switching _manifestCache to ConcurrentDictionary or guarding all accesses (including PruneManifestCache/FillFromCachedHashes) with a dedicated lock.
var remoteManifestProvider = _grainFactory!.GetSystemTarget<IClusterManifestSystemTarget>(Constants.ManifestProviderType, siloAddress);
var hash = await remoteManifestProvider.GetSiloManifestHash().AsTask().WaitAsync(_shutdownCts.Token);
if (_manifestCache.TryGetValue(hash, out var cached))
{
return cached;
AcceptAsync now passes the provided cancellationToken to Socket.AcceptAsync. When that token is canceled, AcceptAsync will typically throw OperationCanceledException, but this method does not handle it and will propagate the exception instead of returning null (which is how the other IConnectionListener implementations here behave). This can break graceful shutdown/cancellation paths.
try
{
var acceptSocket = await _listenSocket!.AcceptAsync(cancellationToken);
acceptSocket.NoDelay = _options.NoDelay;
if (_options.KeepAlive)
GetSiloManifest() is called concurrently from UpdateManifest() via Task.WhenAll, but it reads/writes _manifestCache (a Dictionary) without synchronization. Concurrent Dictionary access can throw (e.g., during resize) or corrupt internal state. Protect cache access (or switch to ConcurrentDictionary) so multiple GetSiloManifest calls can run safely in parallel.
var remoteManifestProvider = _grainFactory!.GetSystemTarget<IClusterManifestSystemTarget>(Constants.ManifestProviderType, siloAddress);
var hash = await remoteManifestProvider.GetSiloManifestHash().AsTask().WaitAsync(_shutdownCts.Token);
if (_manifestCache.TryGetValue(hash, out var cached))
{
return cached;
This new contract does not match existing providers. For example, src/Orleans.Runtime/MembershipService/InMemoryMembershipTable.cs:91 and src/Google/Orleans.Clustering.Firestore/FirestoreMembershipTable.cs:79 delete every old non-Active row, including Joining, ShuttingDown, and Stopping entries. The same-version merge logic added by this change relies on only Dead rows being pruned, so those providers can still produce snapshots which violate that invariant. Align the provider cleanup implementations and conformance tests before documenting this guarantee, or weaken the guarantee and update the merge logic accordingly.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
Deployment-load publication sends each silo's sample to every other silo:
N*(N-1)requests per period. Membership dissemination also needs a small ordinary-update path while retaining reliable convergence after missed updates, topology changes, and partitions.Solution
Add opt-in deterministic dissemination for membership and deployment load, with bounded broadcast queues and anti-entropy repair. The subsystem and both namespace enablement flags default to
false; existing direct paths cover bootstrap, mixed-version operation, and failed admission.Membership broadcasts are sparse deltas; anti-entropy repairs are full snapshots. The broadcast forest includes Joining, Active, ShuttingDown, and Stopping silos, with membership updates ahead of load updates. Each outbound peer/key retains one immutable comparison snapshot captured by an accepted, exactly acknowledged broadcast. After send admission, its next delta compares that snapshot with current owner state, carrying changed entries and removed identities. This preserves coalesced updates and reversions without a generic history or repair-chain protocol.
Receivers apply a delta at its base view, or merge a replay at the target view while preserving canonical fields and maximum heartbeat timestamps. Same-version pruning removes Dead rows only; all other statuses retain their canonical membership. Coalesced deltas across versions can span an intervening Dead transition and cleanup. Missing baselines become eligible for the next full anti-entropy repair. A peer without a comparison snapshot receives a small empty current-view probe; bootstrap and repair establish missing state. Relays compare their accepted current view with each child's baseline, including when the relay is ahead of an incoming update. Initial independent heartbeat and retained-Dead differences remain repair work.
Full repair uses the view version plus a liveness/inventory fingerprint. Same-version repair merges maximum actual
IAmAliveTimeindependently ofStartTime, and reconciles pruned Dead rows while preserving non-Dead entries. Application confirms requested effects against the resulting owner snapshot before producing an exact compact acknowledgment. The cached sender comparison snapshot comes from construction time, preserving updates which arrive while an acknowledgment is in flight.Load publication uses receipt-synchronized cohorts. Non-root silos send one fresh sample to a deterministic root, which contributes locally. A cohort seals when every expected Active incarnation contributes or the publication period expires (one second by default). Held ingress RPCs return the remaining delay to the next cohort boundary, aligning actual sampling timers. Distribution uses fanout eight and forwards admitted sets immediately. Independent ingress admission keeps membership's send slots available.
Item, payload, pending-key, and concurrency limits bound work. Queue notifications retain identities; membership comparison snapshots are shared references and are pruned with peer/key knowledge. A slow set of peers can retain different snapshots, so worst-case comparison memory scales with outbound peers times membership size. Namespace payload caches are constant-size. Local cancellation releases waits while admitted state retains its owner. Shutdown seals cohorts before draining held admissions and then peer queues within the caller's budget.
Scope and rationale
Broadcast creation is separate from full repair construction. Public options focus on enablement, topology, cadence, and resource limits. The independent prerequisites #11294 and #11296 have been merged by humans; this branch is rebased onto their finalized implementations. Provider-owned fatal rollback handling and canonical snapshot merging come from main. Superseded manager reset special-casing and restart-as-recovery tests are removed. Dissemination follows ordinary monotonic membership views.
Broadcast, repair, and cohort integration remain one working convergence feature. Manual performance tooling is separate in #11272. At one sample and one fitting cohort per second, healthy load traffic is approximately
2*(N-1)RPCs/s: 18 at 10 silos, 198 at 100, and 3,998 at 2,000, plus repair and replies. These are projections; root ingress remains concentrated and value delivery remains quadratic. Earlier fixed-window measurements cover a different implementation. Receipt-driven process measurements remain pending separately authorized tooling execution.See runtime dissemination architecture for invariants, timing, limits, and failure semantics. Old/new-binary compatibility covers rolling upgrade, rollback, partition recovery, isolated full-snapshot/heartbeat repair, cancellation, and bounded shutdown with exact state comparisons.