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 c53fa8b622 fix: suppress dead-letter logging during a terminating
shutdown when disabled (#3262)
c53fa8b622 is described below
commit c53fa8b622f23324fb6d46d834447fc72377f21c
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Thu Jul 9 19:26:44 2026 +0800
fix: suppress dead-letter logging during a terminating shutdown when
disabled (#3262)
Motivation:
With pekko.log-dead-letters-during-shutdown = off (the default), dead
letters
were still logged during the entire CoordinatedShutdown sequence. The
DeadLetterListener had no shutdown awareness and was only stopped at the
very
end of shutdown (finalTerminate), so dead letters produced while the system
was
terminating - potentially for the whole shutdown duration - were logged
unconditionally.
Modification:
- DeadLetterListener now suppresses logging while a coordinated shutdown
that
actually terminates this ActorSystem is in progress, when
log-dead-letters-during-shutdown is disabled. Whether the shutdown
terminates
the system is taken from the EFFECTIVE per-reason configuration
(CoordinatedShutdown.confWithOverrides), the same resolution
CoordinatedShutdown
uses, so a non-terminating run - including one where only a specific
shutdown
reason overrides terminate-actor-system to off - keeps logging enabled
instead
of silently disabling it for the rest of a still-running system's life.
The
CoordinatedShutdown extension reference is cached and the effective
decision is
computed at most once to keep the dead-letter path cheap.
Result:
Dead letters are no longer logged during a terminating shutdown when
log-dead-letters-during-shutdown = off, while logging continues normally
for a
non-terminating coordinated shutdown (globally or per reason) and during
normal
operation.
Tests:
- actor-tests/testOnly org.apache.pekko.event.DeadLetterListenerShutdownSpec
(plus DeadLetterSupressionSpec, DeadLetterSuspensionSpec) - all passed.
Added directional tests that fail before the fix: a terminating shutdown
that
must suppress, and non-terminating shutdowns (global and per-reason
override)
that must keep logging.
- actor/mimaReportBinaryIssues - no binary issues.
References:
Fixes #3256
---
.../event/DeadLetterListenerShutdownSpec.scala | 185 +++++++++++++++++++++
.../apache/pekko/event/DeadLetterListener.scala | 53 +++++-
2 files changed, 234 insertions(+), 4 deletions(-)
diff --git
a/actor-tests/src/test/scala/org/apache/pekko/event/DeadLetterListenerShutdownSpec.scala
b/actor-tests/src/test/scala/org/apache/pekko/event/DeadLetterListenerShutdownSpec.scala
new file mode 100644
index 0000000000..dd4b19eb58
--- /dev/null
+++
b/actor-tests/src/test/scala/org/apache/pekko/event/DeadLetterListenerShutdownSpec.scala
@@ -0,0 +1,185 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.pekko.event
+
+import scala.concurrent.Await
+import scala.concurrent.Future
+import scala.concurrent.Promise
+import scala.concurrent.duration._
+
+import com.typesafe.config.ConfigFactory
+
+import org.apache.pekko
+import pekko.Done
+import pekko.actor.Actor
+import pekko.actor.ActorRef
+import pekko.actor.ActorSystem
+import pekko.actor.CoordinatedShutdown
+import pekko.actor.DeadLetter
+import pekko.actor.Props
+import pekko.testkit.PekkoSpec
+import pekko.testkit.TestDuration
+import pekko.testkit.TestKit
+import pekko.testkit.TestProbe
+
+object DeadLetterListenerShutdownSpec {
+ case object RunningMsg
+ case object ShutdownMsg
+
+ // A shutdown reason whose `terminate-actor-system` is overridden to `off`
in the test config, so a
+ // CoordinatedShutdown run for this reason executes its phases without
terminating the system.
+ case object NonTerminatingReason extends CoordinatedShutdown.Reason
+
+ // Forwards only the DeadLetterListener's log events (identified by the
shared "dead letters
+ // encountered" phrase) so the assertions are not disturbed by other INFO
log output produced
+ // during shutdown (e.g. the CoordinatedShutdown "Running ..." message).
+ class DeadLetterLogForwarder(target: ActorRef) extends Actor {
+ def receive: Receive = {
+ case info: Logging.Info if info.message.toString.contains("dead letters
encountered") =>
+ target ! info
+ case _ => // ignore any other log event so it is not reported as an
unhandled message
+ }
+ }
+}
+
+class DeadLetterListenerShutdownSpec extends PekkoSpec {
+ import DeadLetterListenerShutdownSpec._
+
+ private def newSystem(name: String, extraConfig: String): ActorSystem =
+ ActorSystem(
+ name,
+ ConfigFactory
+ .parseString(s"""
+ pekko.log-dead-letters = on
+ pekko.loglevel = INFO
+ $extraConfig
+ """)
+ .withFallback(system.settings.config))
+
+ private def collectorFor(sys: ActorSystem): TestProbe = {
+ val collector = TestProbe()(sys)
+ val forwarder = sys.actorOf(Props(new
DeadLetterLogForwarder(collector.ref)))
+ sys.eventStream.subscribe(forwarder, classOf[Logging.Info])
+ collector
+ }
+
+ // Registers a task that pauses the first CoordinatedShutdown phase, then
runs `during` while the
+ // shutdown is paused (the system is mid-shutdown and the DeadLetterListener
is still alive, since
+ // the listener is only stopped in the last phase). `triggerShutdown` starts
the shutdown
+ // (`terminate()` or `run(...)`) and yields the future to await once the
pause is released, so this
+ // works for both terminating and non-terminating runs. The release always
happens, even if an
+ // assertion throws, so the system is never left blocked.
+ private def whilePausedDuringShutdown(sys: ActorSystem, triggerShutdown: =>
Future[?])(during: => Unit): Unit = {
+ val reached = Promise[Done]()
+ val release = Promise[Done]()
+
CoordinatedShutdown(sys).addTask(CoordinatedShutdown.PhaseBeforeServiceUnbind,
"pause-for-test") { () =>
+ reached.success(Done)
+ release.future
+ }
+ val shutdown = triggerShutdown
+ try {
+ Await.result(reached.future, 5.seconds)
+ during
+ } finally {
+ release.trySuccess(Done)
+ Await.result(shutdown, 10.seconds)
+ }
+ }
+
+ "DeadLetterListener" must {
+
+ "not log dead letters during shutdown when
log-dead-letters-during-shutdown is off (#3256)" in {
+ val sys = newSystem("DeadLetterListenerShutdownSpec-off",
"pekko.log-dead-letters-during-shutdown = off")
+ try {
+ val collector = collectorFor(sys)
+
+ // While the system is running the dead letter is logged.
+ sys.eventStream.publish(DeadLetter(RunningMsg, sys.deadLetters,
sys.deadLetters))
+ collector.expectMsgType[Logging.Info](3.seconds).message.toString
should include("RunningMsg")
+
+ whilePausedDuringShutdown(sys, sys.terminate()) {
+ // A dead letter produced while terminating must not be logged. The
window is dilated so the
+ // regression is still caught on a slow CI box (a reintroduced bug
would log it here).
+ sys.eventStream.publish(DeadLetter(ShutdownMsg, sys.deadLetters,
sys.deadLetters))
+ collector.expectNoMessage(1.second.dilated)
+ }
+ } finally if (!sys.whenTerminated.isCompleted)
TestKit.shutdownActorSystem(sys)
+ }
+
+ "still log dead letters during shutdown when
log-dead-letters-during-shutdown is on" in {
+ val sys = newSystem("DeadLetterListenerShutdownSpec-on",
"pekko.log-dead-letters-during-shutdown = on")
+ try {
+ val collector = collectorFor(sys)
+
+ whilePausedDuringShutdown(sys, sys.terminate()) {
+ // With shutdown logging enabled the dead letter is still logged
while terminating.
+ sys.eventStream.publish(DeadLetter(ShutdownMsg, sys.deadLetters,
sys.deadLetters))
+ collector.expectMsgType[Logging.Info](3.seconds).message.toString
should include("ShutdownMsg")
+ }
+ } finally if (!sys.whenTerminated.isCompleted)
TestKit.shutdownActorSystem(sys)
+ }
+
+ "still log dead letters during a coordinated shutdown that does not
terminate the system (#3256)" in {
+ // With terminate-actor-system = off a CoordinatedShutdown.run()
executes its phases but keeps
+ // the system alive, leaving a shutdown reason set even though the
system is not terminating.
+ // Dead-letter logging must keep working in that case rather than being
silently disabled for
+ // the remaining lifetime of the still-running system.
+ val sys = newSystem(
+ "DeadLetterListenerShutdownSpec-nonterminating",
+ """
+ pekko.log-dead-letters-during-shutdown = off
+ pekko.coordinated-shutdown.terminate-actor-system = off
+ pekko.coordinated-shutdown.run-by-actor-system-terminate = off
+ """)
+ try {
+ val collector = collectorFor(sys)
+
+ whilePausedDuringShutdown(sys,
CoordinatedShutdown(sys).run(CoordinatedShutdown.UnknownReason)) {
+ sys.eventStream.publish(DeadLetter(ShutdownMsg, sys.deadLetters,
sys.deadLetters))
+ collector.expectMsgType[Logging.Info](3.seconds).message.toString
should include("ShutdownMsg")
+ }
+ } finally TestKit.shutdownActorSystem(sys)
+ }
+
+ "still log dead letters when a shutdown reason overrides
terminate-actor-system to off (#3256)" in {
+ // The base config terminates the system, but the specific shutdown
reason overrides
+ // terminate-actor-system to off, so this run does NOT terminate the
system. The effective
+ // (per-reason) value must be honored, not the base setting, otherwise
logging would be
+ // silently suppressed forever on a still-running system.
+ val reasonClass = NonTerminatingReason.getClass.getName
+ val sys = newSystem(
+ "DeadLetterListenerShutdownSpec-reason-override",
+ "pekko.log-dead-letters-during-shutdown = off\n" +
+ "pekko.coordinated-shutdown.run-by-actor-system-terminate = off\n" +
+ "pekko.coordinated-shutdown.reason-overrides.\"" + reasonClass +
"\".terminate-actor-system = off")
+ try {
+ val collector = collectorFor(sys)
+
+ whilePausedDuringShutdown(sys,
CoordinatedShutdown(sys).run(NonTerminatingReason)) {
+ sys.eventStream.publish(DeadLetter(ShutdownMsg, sys.deadLetters,
sys.deadLetters))
+ collector.expectMsgType[Logging.Info](3.seconds).message.toString
should include("ShutdownMsg")
+ }
+
+ // The run has completed but the system is still alive and its
shutdown reason stays set;
+ // logging must keep working for the rest of the system's life, not be
suppressed forever.
+ sys.eventStream.publish(DeadLetter(ShutdownMsg, sys.deadLetters,
sys.deadLetters))
+ collector.expectMsgType[Logging.Info](3.seconds).message.toString
should include("ShutdownMsg")
+ } finally TestKit.shutdownActorSystem(sys)
+ }
+ }
+}
diff --git
a/actor/src/main/scala/org/apache/pekko/event/DeadLetterListener.scala
b/actor/src/main/scala/org/apache/pekko/event/DeadLetterListener.scala
index 398dba099b..c3ee6dd2d8 100644
--- a/actor/src/main/scala/org/apache/pekko/event/DeadLetterListener.scala
+++ b/actor/src/main/scala/org/apache/pekko/event/DeadLetterListener.scala
@@ -21,6 +21,7 @@ import pekko.actor.Actor
import pekko.actor.ActorLogMarker
import pekko.actor.ActorRef
import pekko.actor.AllDeadLetters
+import pekko.actor.CoordinatedShutdown
import pekko.actor.DeadLetter
import pekko.actor.DeadLetterActorRef
import pekko.actor.DeadLetterSuppression
@@ -35,6 +36,22 @@ class DeadLetterListener extends Actor {
val eventStream: EventStream = context.system.eventStream
protected val maxCount: Int = context.system.settings.LogDeadLetters
private val isAlwaysLoggingDeadLetters = maxCount == Int.MaxValue
+ // When dead-letter logging during shutdown is disabled we must stop logging
as soon as the
+ // system starts terminating, not only once this listener is finally stopped
at the very end of
+ // CoordinatedShutdown. Otherwise dead letters produced throughout the
(potentially long) shutdown
+ // are still logged unconditionally. See issue #3256.
+ private val logDeadLettersDuringShutdown =
context.system.settings.LogDeadLettersDuringShutdown
+ // The `pekko.coordinated-shutdown` config, used to resolve the EFFECTIVE
(per-reason)
+ // `terminate-actor-system` value of an in-progress shutdown; see
`terminatingShutdownInProgress`.
+ private val coordinatedShutdownConfig =
context.system.settings.config.getConfig("pekko.coordinated-shutdown")
+ // Resolved lazily so that handling a dead letter never forces
CoordinatedShutdown initialization
+ // from the actor's constructor. At worst the first dead letter triggers the
one-time extension
+ // initialization; afterwards this just caches the reference and avoids a
per-dead-letter lookup.
+ private lazy val coordinatedShutdown = CoordinatedShutdown(context.system)
+ // Cached decision of whether the in-progress coordinated shutdown actually
terminates the system.
+ // The shutdown reason is set once and never changes, so this is computed at
most once. Accessed
+ // only from the actor's own thread, so a plain var is sufficient.
+ private var shutdownTerminatesSystem: Option[Boolean] = None
protected var count: Int = 0
override def preStart(): Unit = {
@@ -72,7 +89,7 @@ class DeadLetterListener extends Actor {
private def receiveWithAlwaysLogging: Receive = {
case d: AllDeadLetters =>
- if (!isWrappedSuppressed(d)) {
+ if (!isSuppressed(d)) {
incrementCount()
logDeadLetter(d, doneMsg = "")
}
@@ -80,7 +97,7 @@ class DeadLetterListener extends Actor {
private def receiveWithMaxCountLogging: Receive = {
case d: AllDeadLetters =>
- if (!isWrappedSuppressed(d)) {
+ if (!isSuppressed(d)) {
incrementCount()
if (count == maxCount) {
logDeadLetter(d, ", no more dead letters will be logged")
@@ -93,7 +110,7 @@ class DeadLetterListener extends Actor {
private def receiveWithSuspendLogging(suspendDuration: FiniteDuration):
Receive = {
case d: AllDeadLetters =>
- if (!isWrappedSuppressed(d)) {
+ if (!isSuppressed(d)) {
incrementCount()
if (count == maxCount) {
val doneMsg = s", no more dead letters will be logged in next
[${suspendDuration.pretty}]"
@@ -106,7 +123,7 @@ class DeadLetterListener extends Actor {
private def receiveWhenSuspended(suspendDuration: FiniteDuration,
suspendDeadline: Deadline): Receive = {
case d: AllDeadLetters =>
- if (!isWrappedSuppressed(d)) {
+ if (!isSuppressed(d)) {
incrementCount()
if (suspendDeadline.isOverdue()) {
val doneMsg = s", of which ${count - maxCount - 1} were not logged.
The counter will be reset now"
@@ -151,6 +168,34 @@ class DeadLetterListener extends Actor {
(snd ne ActorRef.noSender) && (snd ne context.system.deadLetters) &&
!snd.isInstanceOf[DeadLetterActorRef]
}
+ private def isSuppressed(d: AllDeadLetters): Boolean =
+ isWrappedSuppressed(d) || suppressedDuringShutdown
+
+ // Suppress logging while a terminating coordinated shutdown is in progress
and dead-letter logging
+ // during shutdown is disabled (`pekko.log-dead-letters-during-shutdown`,
`off` by default).
+ private def suppressedDuringShutdown: Boolean =
+ !logDeadLettersDuringShutdown && terminatingShutdownInProgress
+
+ // True while a coordinated shutdown that actually terminates this
ActorSystem is in progress.
+ // `terminate-actor-system` can be overridden per shutdown reason
(CoordinatedShutdown
+ // reason-overrides), so the decision must be taken from the EFFECTIVE
config for the active reason
+ // rather than from the base setting: a non-terminating run leaves a
shutdown reason set on a
+ // still-running system and must NOT suppress logging, otherwise dead-letter
logging would be
+ // silently disabled for the rest of that system's life (issue #3256).
`confWithOverrides` is the
+ // same resolution CoordinatedShutdown itself uses to decide whether to
terminate the system.
+ private def terminatingShutdownInProgress: Boolean =
+ coordinatedShutdown.shutdownReason() match {
+ case reason @ Some(_) =>
+ shutdownTerminatesSystem.getOrElse {
+ val terminates =
+ CoordinatedShutdown.confWithOverrides(coordinatedShutdownConfig,
reason).getBoolean(
+ "terminate-actor-system")
+ shutdownTerminatesSystem = Some(terminates)
+ terminates
+ }
+ case None => false
+ }
+
private def isWrappedSuppressed(d: AllDeadLetters): Boolean = {
d.message match {
case w: WrappedMessage if w.message.isInstanceOf[DeadLetterSuppression]
=> true
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]