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]