This is an automated email from the ASF dual-hosted git repository.

He-Pin pushed a commit to branch fix/broadcast-hub-poststop
in repository https://gitbox.apache.org/repos/asf/pekko.git

commit c751b6f102d1ff00e2a2010133f29d0090eab58a
Author: θ™ŽιΈ£ <[email protected]>
AuthorDate: Tue Jul 14 18:45:44 2026 +0800

    fix(stream): notify BroadcastHub consumers on materializer shutdown
    
    Motivation:
    When the materializer running a BroadcastHub.sink is shut down (e.g. the
    hosting actor stops), consumers materialized on a different materializer
    may not be notified of the hub's termination (see #3345).
    
    Modification:
    Move registered-consumer (consumerWheel) notification from
    onUpstreamFailure to postStop so postStop is the single notification
    point. postStop now handles both Open (normal termination) and Closed
    (after upstream failure) states, ensuring registered consumers always
    receive the appropriate signal.
    
    Result:
    BroadcastHub consumers are reliably notified when the hub's materializer
    shuts down, preventing consumers from hanging indefinitely.
    
    Tests:
    - sbt "stream-tests / Test / testOnly 
org.apache.pekko.stream.scaladsl.HubSpec"
      51 passed, 0 failed
    
    References:
    Fixes #3345
---
 .../org/apache/pekko/stream/scaladsl/HubSpec.scala | 24 +++++++++++++-
 .../org/apache/pekko/stream/scaladsl/Hub.scala     | 37 ++++++++++++++--------
 2 files changed, 47 insertions(+), 14 deletions(-)

diff --git 
a/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala 
b/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
index b32c8746b1..628585cca0 100644
--- a/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
+++ b/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
@@ -19,7 +19,7 @@ import scala.concurrent.duration._
 
 import org.apache.pekko
 import pekko.Done
-import pekko.stream.KillSwitches
+import pekko.stream.{ KillSwitches, Materializer }
 import pekko.stream.ThrottleMode
 import pekko.stream.impl.ActorPublisher
 import pekko.stream.testkit.StreamSpec
@@ -559,6 +559,28 @@ class HubSpec extends StreamSpec {
       }
     }
 
+    "notify consumers when hub materializer is shut down" in {
+      val hubMat = Materializer(system)
+      val upstream = TestPublisher.probe[Int]()
+      val source = 
Source.fromPublisher(upstream).runWith(BroadcastHub.sink(8))(hubMat)
+
+      val downstream = TestSubscriber.probe[Int]()
+      source.runWith(Sink.fromSubscriber(downstream))
+
+      downstream.request(1)
+      Thread.sleep(100)
+
+      upstream.sendNext(1)
+      downstream.expectNext(1)
+
+      hubMat.shutdown()
+
+      // The consumer must be notified (not left hanging), which is the fix 
for #3345.
+      // The signal is an error because the SubSink output boundary failure
+      // arrives before the postStop completion callback.
+      downstream.expectError()
+    }
+
     "handle cancelled Sink" in {
       val in = TestPublisher.probe[Int]()
       val hubSource = Source.fromPublisher(in).runWith(BroadcastHub.sink(4))
diff --git a/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala 
b/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
index 32e647841f..f4ce5cdebd 100644
--- a/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
+++ b/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
@@ -652,21 +652,12 @@ private[pekko] class 
BroadcastHub[T](startAfterNrOfConsumers: Int, bufferSize: I
     override def onUpstreamFailure(ex: Throwable): Unit = {
       val failMessage = HubCompleted(Some(ex))
 
-      // Notify pending consumers and set tombstone
+      // Notify pending consumers and set tombstone.
+      // Registered consumers in the consumerWheel are notified by postStop.
       
state.getAndSet(Closed(Some(ex))).asInstanceOf[Open].registrations.foreach { 
consumer =>
         consumer.callback.invoke(failMessage)
       }
 
-      // Notify registered consumers β€” skip null (empty) slots
-      var idx = 0
-      while (idx < consumerWheel.length) {
-        val bucket = consumerWheel(idx)
-        if (bucket ne null) {
-          val itr = bucket.valuesIterator
-          while (itr.hasNext) itr.next().callback.invoke(failMessage)
-        }
-        idx += 1
-      }
       failStage(ex)
     }
 
@@ -759,19 +750,39 @@ private[pekko] class 
BroadcastHub[T](startAfterNrOfConsumers: Int, bufferSize: I
     }
 
     override def postStop(): Unit = {
-      // Notify pending consumers and set tombstone
+      // Notify all consumers (pending and registered) when the stage stops.
+      // Registered consumers in the consumerWheel are notified here (not in 
onUpstreamFailure)
+      // so that materializer shutdown produces the correct signal.
 
       @tailrec def tryClose(): Unit = state.get() match {
-        case Closed(_)  => // Already closed, ignore
+        case Closed(_) => // Already closed by onUpstreamFailure β€” fall 
through to notify wheel below
+          notifyRegisteredConsumers()
         case open: Open =>
           if (state.compareAndSet(open, Closed(None))) {
             val completedMessage = HubCompleted(None)
             open.registrations.foreach { consumer =>
               consumer.callback.invoke(completedMessage)
             }
+            notifyRegisteredConsumers()
           } else tryClose()
       }
 
+      def notifyRegisteredConsumers(): Unit = {
+        val message = state.get() match {
+          case Closed(Some(ex)) => HubCompleted(Some(ex))
+          case _                => HubCompleted(None)
+        }
+        var idx = 0
+        while (idx < consumerWheel.length) {
+          val bucket = consumerWheel(idx)
+          if (bucket ne null) {
+            val itr = bucket.valuesIterator
+            while (itr.hasNext) itr.next().callback.invoke(message)
+          }
+          idx += 1
+        }
+      }
+
       tryClose()
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to