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

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


The following commit(s) were added to refs/heads/main by this push:
     new c78e2d7610 chore: Add tests for not invoking `onComplete` twice for 
statefulMap operator.
c78e2d7610 is described below

commit c78e2d76102ee0f9768ce7bca4bbbd0f71a5e9a1
Author: He-Pin <[email protected]>
AuthorDate: Mon Dec 25 11:33:47 2023 +0800

    chore: Add tests for not invoking `onComplete` twice for statefulMap 
operator.
---
 .../stream/scaladsl/FlowStatefulMapSpec.scala      | 78 +++++++++++++++++++---
 1 file changed, 70 insertions(+), 8 deletions(-)

diff --git 
a/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/FlowStatefulMapSpec.scala
 
b/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/FlowStatefulMapSpec.scala
index 28f22d642a..2cf80a96d3 100644
--- 
a/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/FlowStatefulMapSpec.scala
+++ 
b/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/FlowStatefulMapSpec.scala
@@ -13,6 +13,15 @@
 
 package org.apache.pekko.stream.scaladsl
 
+import java.util.concurrent.atomic.AtomicInteger
+
+import scala.annotation.nowarn
+import scala.concurrent.Await
+import scala.concurrent.Promise
+import scala.concurrent.duration.DurationInt
+import scala.util.Success
+import scala.util.control.NoStackTrace
+
 import org.apache.pekko
 import pekko.Done
 import pekko.stream.AbruptStageTerminationException
@@ -21,16 +30,10 @@ import pekko.stream.ActorMaterializer
 import pekko.stream.Supervision
 import pekko.stream.testkit.StreamSpec
 import pekko.stream.testkit.TestSubscriber
+import pekko.stream.testkit.Utils.TE
 import pekko.stream.testkit.scaladsl.TestSink
 import pekko.stream.testkit.scaladsl.TestSource
-
-import java.util.concurrent.atomic.AtomicInteger
-import scala.annotation.nowarn
-import scala.concurrent.Await
-import scala.concurrent.Promise
-import scala.concurrent.duration.DurationInt
-import scala.util.Success
-import scala.util.control.NoStackTrace
+import pekko.testkit.EventFilter
 
 class FlowStatefulMapSpec extends StreamSpec {
 
@@ -371,5 +374,64 @@ class FlowStatefulMapSpec extends StreamSpec {
         .expectComplete()
       gate.ensure()
     }
+
+    "will not call `onComplete` twice if `f` fail" in {
+      val closedCounter = new AtomicInteger(0)
+      val probe = Source
+        .repeat(1)
+        .statefulMap(() => "opening resource")(
+          (_, _) => throw TE("failing read"),
+          _ => {
+            closedCounter.incrementAndGet()
+            None
+          })
+        .runWith(TestSink.probe[String])
+
+      probe.request(1)
+      probe.expectError(TE("failing read"))
+      closedCounter.get() should ===(1)
+    }
+
+    "will not call `onComplete` twice if both `f` and `onComplete` fail" in {
+      val closedCounter = new AtomicInteger(0)
+      val probe = Source
+        .repeat(1)
+        .statefulMap(() => "opening resource")((_, _) => throw TE("failing 
read"),
+          _ => {
+            if (closedCounter.incrementAndGet() == 1) {
+              throw TE("boom")
+            }
+            None
+          })
+        .runWith(TestSink.probe[Int])
+
+      EventFilter[TE](occurrences = 1).intercept {
+        probe.request(1)
+        probe.expectError(TE("boom"))
+      }
+      closedCounter.get() should ===(1)
+    }
+
+    "will not call `onComplete` twice if `onComplete` fail on upstream 
complete" in {
+      val closedCounter = new AtomicInteger(0)
+      val (pub, sub) = TestSource[Int]()
+        .statefulMap(() => "opening resource")((state, value) => (state, 
value),
+          _ => {
+            closedCounter.incrementAndGet()
+            throw TE("boom")
+          })
+        .toMat(TestSink.probe[Int])(Keep.both)
+        .run()
+
+      EventFilter[TE](occurrences = 1).intercept {
+        sub.request(1)
+        pub.sendNext(1)
+        sub.expectNext(1)
+        sub.request(1)
+        pub.sendComplete()
+        sub.expectError(TE("boom"))
+      }
+      closedCounter.get() shouldBe 1
+    }
   }
 }


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

Reply via email to