This is an automated email from the ASF dual-hosted git repository.
He-Pin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko.git
The following commit(s) were added to refs/heads/main by this push:
new 5293186b72 test: harden PersistentActorRecoveryTimeoutSpec receive
timeout test (#3436)
5293186b72 is described below
commit 5293186b728d66cdfafd61b474db74477557cfb4
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Sat Aug 15 20:04:25 2026 +0800
test: harden PersistentActorRecoveryTimeoutSpec receive timeout test (#3436)
Motivation:
The "should not interfere with receive timeouts" test previously only
verified that the receive timeout value is preserved after recovery.
It did not verify that a short receive timeout doesn't accidentally
trigger the recovery timeout, nor that the recovery timeout still
fires when the journal is unresponsive.
Modification:
- Use separate probes for persist and replay phases to avoid
cross-contamination.
- Use a short receive timeout (20ms) during replay to verify it does
not trigger recovery timeout (expectNoMessage).
- Verify that recovery timeout (500ms) still fires when the journal
does not respond (expectMsgType[Failure] with RecoveryTimedOut).
- Reduce recovery-event-timeout from 30s to 500ms for the
receive-timeout journal to keep the test fast.
- Remove unused TestDuration import.
Result:
The test now encodes the intended contract: receive timeout and
recovery timeout are independent mechanisms. The test is deterministic
and fast (~700ms total for the replay phase).
Tests:
- sbt "persistence / Test / testOnly
org.apache.pekko.persistence.PersistentActorRecoveryTimeoutSpec"
References:
Refs akka/akka-core#32030
---
.../PersistentActorRecoveryTimeoutSpec.scala | 37 +++++++++++++---------
1 file changed, 22 insertions(+), 15 deletions(-)
diff --git
a/persistence/src/test/scala/org/apache/pekko/persistence/PersistentActorRecoveryTimeoutSpec.scala
b/persistence/src/test/scala/org/apache/pekko/persistence/PersistentActorRecoveryTimeoutSpec.scala
index bc466701b7..7e6a1e17d3 100644
---
a/persistence/src/test/scala/org/apache/pekko/persistence/PersistentActorRecoveryTimeoutSpec.scala
+++
b/persistence/src/test/scala/org/apache/pekko/persistence/PersistentActorRecoveryTimeoutSpec.scala
@@ -19,7 +19,7 @@ import org.apache.pekko
import pekko.actor.{ Actor, ActorLogging, ActorRef, Props }
import pekko.actor.Status.Failure
import pekko.persistence.journal.SteppingInmemJournal
-import pekko.testkit.{ ImplicitSender, PekkoSpec, TestDuration, TestProbe }
+import pekko.testkit.{ ImplicitSender, PekkoSpec, TestProbe }
import com.typesafe.config.ConfigFactory
@@ -35,7 +35,7 @@ object PersistentActorRecoveryTimeoutSpec {
|pekko.persistence.journal.stepping-inmem.recovery-event-timeout=3s
|$receiveTimeoutJournalPluginId.class=${classOf[SteppingInmemJournal].getName}
|$receiveTimeoutJournalPluginId.instance-id="$receiveTimeoutJournalId"
- |$receiveTimeoutJournalPluginId.recovery-event-timeout=30s
+ |$receiveTimeoutJournalPluginId.recovery-event-timeout=500ms
""".stripMargin))
.withFallback(PersistenceSpec.config("stepping-inmem",
"PersistentActorRecoveryTimeoutSpec"))
@@ -128,18 +128,17 @@ class PersistentActorRecoveryTimeoutSpec
}
"should not interfere with receive timeouts" in {
- val timeout = 42.days
-
- val probe = TestProbe()
+ val probe1 = TestProbe()
val persisting =
-
system.actorOf(Props(classOf[PersistentActorRecoveryTimeoutSpec.TestReceiveTimeoutActor],
timeout, probe.ref))
+ system.actorOf(
+
Props(classOf[PersistentActorRecoveryTimeoutSpec.TestReceiveTimeoutActor],
42.days, probe1.ref))
awaitAssert(SteppingInmemJournal.getRef(receiveTimeoutJournalId),
3.seconds)
val journal = SteppingInmemJournal.getRef(receiveTimeoutJournalId)
// initial read highest
SteppingInmemJournal.step(journal)
- probe.expectMsg(timeout)
+ probe1.expectMsg(42.days)
persisting ! "A"
SteppingInmemJournal.step(journal)
@@ -149,17 +148,25 @@ class PersistentActorRecoveryTimeoutSpec
system.stop(persisting)
expectTerminated(persisting)
- // now replay and verify that recovery keeps the actor's receive timeout
-
system.actorOf(Props(classOf[PersistentActorRecoveryTimeoutSpec.TestReceiveTimeoutActor],
timeout, probe.ref))
+ // now replay with a short receive timeout to verify it doesn't trigger
recovery timeout
+ val probe2 = TestProbe()
+ val timeout = 20.millis
+ val replaying =
+
system.actorOf(Props(classOf[PersistentActorRecoveryTimeoutSpec.TestReceiveTimeoutActor],
timeout, probe2.ref))
- // Release both recovery journal operations up front. Waiting for the
second stepped
- // operation can race with the recovery timeout under heavy CI load.
- journal ! SteppingInmemJournal.Token
- journal ! SteppingInmemJournal.Token
+ probe2.expectNoMessage(50.millis) // longer than the receive timeout
+ // initial read highest
+ SteppingInmemJournal.step(journal)
+ probe2.expectNoMessage(100.millis) // receive timeout should not trigger
recovery timeout
- // we should get initial receive timeout back from actor when replay
completes
- probe.expectMsg(30.seconds.dilated, timeout)
+ // but waiting longer without SteppingInmemJournal.step will be recovery
timeout
+ probe2.expectMsgType[Failure].cause shouldBe a[RecoveryTimedOut]
+ watch(replaying)
+ expectTerminated(replaying)
+ // avoid having it stuck in the next test from the
+ // last read request above
+ SteppingInmemJournal.step(journal)
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]