Skip to content

Commit

Permalink
chore: Fix RecipeAdhocSource test.
Browse files Browse the repository at this point in the history
  • Loading branch information
He-Pin committed Dec 28, 2023
1 parent 09023ce commit ff09b28
Show file tree
Hide file tree
Showing 2 changed files with 10 additions and 7 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import org.apache.pekko.actor.ActorSystem;
import org.apache.pekko.dispatch.Futures;
import org.apache.pekko.japi.pf.PFBuilder;
import org.apache.pekko.stream.BackpressureTimeoutException;
import org.apache.pekko.stream.javadsl.Keep;
import org.apache.pekko.stream.javadsl.Source;
import org.apache.pekko.stream.testkit.TestSubscriber;
Expand Down Expand Up @@ -232,7 +233,7 @@ public void restartUptoMaxRetries() throws Exception {
assertEquals(4, startedCount.get()); // startCount == 4, which means "re"-tried 3 times

Thread.sleep(500);
assertEquals(TimeoutException.class, probe.expectError().getClass());
assertEquals(BackpressureTimeoutException.class, probe.expectError().getClass());
probe.request(1); // send demand
probe.expectNoMessage(FiniteDuration.create(200, "milliseconds")); // but no more restart
}
Expand Down
14 changes: 8 additions & 6 deletions docs/src/test/scala/docs/stream/cookbook/RecipeAdhocSource.scala
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,14 @@

package docs.stream.cookbook

import java.util.concurrent.atomic.{ AtomicBoolean, AtomicInteger }

import org.apache.pekko.stream.scaladsl.Source
import org.apache.pekko.stream.testkit.scaladsl.TestSink
import org.apache.pekko.testkit.TimingTest
import org.apache.pekko.{ Done, NotUsed }
import java.util.concurrent.atomic.{AtomicBoolean, AtomicInteger}
import org.apache.pekko
import pekko.stream.scaladsl.Source
import pekko.stream.testkit.scaladsl.TestSink
import pekko.stream.BackpressureTimeoutException
import pekko.testkit.TimingTest
import pekko.{Done, NotUsed}

import scala.concurrent._
import scala.concurrent.duration._
Expand Down Expand Up @@ -127,7 +129,7 @@ class RecipeAdhocSource extends RecipeSpec {
startedCount.get() should be(4) // startCount == 4, which means "re"-tried 3 times

Thread.sleep(500)
sink.expectError().getClass should be(classOf[TimeoutException])
sink.expectError().getClass should be(classOf[BackpressureTimeoutException])
sink.request(1) // send demand
sink.expectNoMessage(200.milliseconds) // but no more restart
}
Expand Down

0 comments on commit ff09b28

Please sign in to comment.