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

raboof pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors.git


The following commit(s) were added to refs/heads/main by this push:
     new 817d8f433 MQTT Streaming: fix BroadcastHub race in server flow tests 
(#1977)
817d8f433 is described below

commit 817d8f4334eefb33fd832a8bc54a05bd17812414
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Sep 30 17:00:06 2026 +0100

    MQTT Streaming: fix BroadcastHub race in server flow tests (#1977)
    
    Motivation:
    MqttFlowTest.establishServerBidirectionalConnectionAndSubscribeToATopic
    intermittently timed out: the server received the client's Connect but
    never replied with ConnAck. The connection handler feeds events into a
    BroadcastHub with two consumers - runForeach (which sends ConnAck,
    SubAck and PubAck) and the source returned to flatMapMerge (drained by
    Sink.ignore). BroadcastHub only holds elements while it has no
    consumers, so if the returned source attached first, the Connect could
    be delivered before runForeach attached and was never answered. The
    leaked streams then also failed the next test's assertAllStagesStopped.
    
    Modification:
    Use BroadcastHub with startAfterNrOfConsumers = 2 in the server flow of
    MqttFlowTest and MqttFlowSpec so no event is emitted until both
    consumers are attached. This also fixes the documentation snippet.
    
    Result:
    The handler always sees the Connect, so the server test no longer
    depends on consumer attach order.
    
    Tests:
    - sbt mqtt-streaming/Test/javafmt mqtt-streaming/Test/scalafmt 
mqtt-streaming/Test/compile
    - sbt 'mqtt-streaming/testOnly docs.javadsl.MqttFlowTest docs.scaladsl.*' 
(x3):
      server flow tests pass in MqttFlowTest, TypedMqttFlowSpec and
      UntypedMqttFlowSpec; broker-dependent tests fail locally only because
      no MQTT broker was running (connection refused on localhost:1883)
    
    References:
    Refs #468
---
 mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java    | 6 +++++-
 mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala | 4 +++-
 2 files changed, 8 insertions(+), 2 deletions(-)

diff --git a/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java 
b/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java
index 780648378..2e9470aea 100644
--- a/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java
+++ b/mqtt-streaming/src/test/java/docs/javadsl/MqttFlowTest.java
@@ -194,7 +194,11 @@ public class MqttFlowTest {
                       run =
                           Source.<Command<Object>>queue(2, 
OverflowStrategy.dropHead())
                               .via(mqttFlow)
-                              
.toMat(BroadcastHub.of(DecodeErrorOrEvent.classOf()), Keep.both())
+                              // wait for both consumers (runForeach and the 
returned source)
+                              // so the Connect cannot bypass the handler that 
sends ConnAck
+                              .toMat(
+                                  
BroadcastHub.of(DecodeErrorOrEvent.classOf(), 2, 256),
+                                  Keep.both())
                               .run(system);
 
                   SourceQueueWithComplete<Command<Object>> queue = run.first();
diff --git a/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala 
b/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala
index fa4d0ed1c..d6a56aadd 100644
--- a/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala
+++ b/mqtt-streaming/src/test/scala/docs/scaladsl/MqttFlowSpec.scala
@@ -133,7 +133,9 @@ abstract class MqttFlowSpecBase(clientId: String, topic: 
String, system: ActorSy
               val (queue, source) = Source
                 .queue[Command[Nothing]](3, OverflowStrategy.dropHead)
                 .via(mqttFlow)
-                .toMat(BroadcastHub.sink)(Keep.both)
+                // wait for both consumers (runForeach and the returned source)
+                // so the Connect cannot bypass the handler that sends ConnAck
+                .toMat(BroadcastHub.sink(2, 256))(Keep.both)
                 .run()
 
               val subscribed = Promise[Done]()


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

Reply via email to