@@ -680,8 +680,7 @@ private async Task BroadcastEventsAsync()
680680 {
681681 // Read until the channel is completed (not until cancelled)
682682 // This ensures all pending events are broadcast during graceful shutdown
683- // Note: ReadAllAsync() completes gracefully when the channel writer is completed,
684- // so no try-catch is needed for normal operation
683+ // Note: ReadAllAsync() completes gracefully when the channel writer is completed
685684 await foreach ( LeadershipChangedEventArgs eventArgs in broadcastChannel . Reader . ReadAllAsync ( ) . ConfigureAwait ( false ) )
686685 {
687686 // Broadcast to all subscriber channels
@@ -705,42 +704,18 @@ private async Task BroadcastEventsAsync()
705704 /// An async enumerable that yields leadership change events as they occur.
706705 /// The enumeration will continue until the service is disposed or the cancellation token is triggered.
707706 /// </returns>
708- public async IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsync (
709- [ EnumeratorCancellation ] CancellationToken cancellationToken = default )
710- {
711- ThrowIfDisposed ( ) ;
707+ public IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsync (
708+ CancellationToken cancellationToken = default )
709+ => GetLeadershipChangesAsyncCore ( eventTypes : null , cancellationToken , subscriberRegistered : null ) ;
712710
713- // Create a dedicated channel for this subscriber
714- var subscriberChannel = Channel . CreateUnbounded < LeadershipChangedEventArgs > ( new UnboundedChannelOptions
715- {
716- SingleWriter = true ,
717- SingleReader = true
718- } ) ;
719-
720- // Register the subscriber channel
721- lock ( subscriberLock )
722- {
723- subscriberChannels . Add ( subscriberChannel ) ;
724- }
725-
726- try
727- {
728- // Read from the subscriber channel
729- await foreach ( LeadershipChangedEventArgs eventArgs in subscriberChannel . Reader . ReadAllAsync ( cancellationToken ) . ConfigureAwait ( false ) )
730- {
731- yield return eventArgs ;
732- }
733- }
734- finally
735- {
736- // Unregister and complete the subscriber channel
737- lock ( subscriberLock )
738- {
739- subscriberChannels . Remove ( subscriberChannel ) ;
740- }
741- subscriberChannel . Writer . TryComplete ( ) ;
742- }
743- }
711+ /// <summary>
712+ /// Internal overload that signals when the subscriber is registered.
713+ /// Used by tests to avoid race conditions.
714+ /// </summary>
715+ internal IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsyncCore (
716+ CancellationToken cancellationToken ,
717+ TaskCompletionSource ? subscriberRegistered )
718+ => GetLeadershipChangesAsyncCore ( eventTypes : null , cancellationToken , subscriberRegistered ) ;
744719
745720 /// <summary>
746721 /// Gets an async enumerable stream of filtered leadership change events.
@@ -751,9 +726,28 @@ public async IAsyncEnumerable<LeadershipChangedEventArgs> GetLeadershipChangesAs
751726 /// An async enumerable that yields filtered leadership change events as they occur.
752727 /// The enumeration will continue until the service is disposed or the cancellation token is triggered.
753728 /// </returns>
754- public async IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsync (
729+ public IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsync (
730+ LeadershipEventType eventTypes ,
731+ CancellationToken cancellationToken = default )
732+ => GetLeadershipChangesAsyncCore ( eventTypes , cancellationToken , subscriberRegistered : null ) ;
733+
734+ /// <summary>
735+ /// Internal overload for filtered events that signals when the subscriber is registered.
736+ /// Used by tests to avoid race conditions.
737+ /// </summary>
738+ internal IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsyncCore (
755739 LeadershipEventType eventTypes ,
756- [ EnumeratorCancellation ] CancellationToken cancellationToken = default )
740+ CancellationToken cancellationToken ,
741+ TaskCompletionSource ? subscriberRegistered )
742+ => GetLeadershipChangesAsyncCore ( ( LeadershipEventType ? ) eventTypes , cancellationToken , subscriberRegistered ) ;
743+
744+ /// <summary>
745+ /// Core implementation that handles both filtered and unfiltered event streams.
746+ /// </summary>
747+ private async IAsyncEnumerable < LeadershipChangedEventArgs > GetLeadershipChangesAsyncCore (
748+ LeadershipEventType ? eventTypes ,
749+ [ EnumeratorCancellation ] CancellationToken cancellationToken ,
750+ TaskCompletionSource ? subscriberRegistered )
757751 {
758752 ThrowIfDisposed ( ) ;
759753
@@ -770,20 +764,28 @@ public async IAsyncEnumerable<LeadershipChangedEventArgs> GetLeadershipChangesAs
770764 subscriberChannels . Add ( subscriberChannel ) ;
771765 }
772766
767+ // Signal that the subscriber is now registered
768+ subscriberRegistered ? . TrySetResult ( ) ;
769+
773770 try
774771 {
775- // Read from the subscriber channel and filter
772+ // Read from the subscriber channel
776773 await foreach ( LeadershipChangedEventArgs eventArgs in subscriberChannel . Reader . ReadAllAsync ( cancellationToken ) . ConfigureAwait ( false ) )
777774 {
775+ // If no filter specified, yield all events
776+ if ( eventTypes is null )
777+ {
778+ yield return eventArgs ;
779+ continue ;
780+ }
781+
778782 // Check if this event matches the requested event types
779- bool shouldYield = ( ( eventTypes & LeadershipEventType . Acquired ) != 0 && eventArgs . BecameLeader ) ||
780- ( ( eventTypes & LeadershipEventType . Lost ) != 0 && eventArgs . LostLeadership ) ||
781- ( ( eventTypes & LeadershipEventType . Changed ) != 0 ) ;
783+ bool shouldYield = ( ( eventTypes . Value & LeadershipEventType . Acquired ) != 0 && eventArgs . BecameLeader ) ||
784+ ( ( eventTypes . Value & LeadershipEventType . Lost ) != 0 && eventArgs . LostLeadership ) ||
785+ ( ( eventTypes . Value & LeadershipEventType . Changed ) != 0 ) ;
782786
783787 if ( shouldYield )
784- {
785788 yield return eventArgs ;
786- }
787789 }
788790 }
789791 finally
0 commit comments