From 67b460fb09095cef49fc0a07def63752314198ec Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Fri, 18 Sep 2026 11:07:55 +0100 Subject: [PATCH] test: cover that harmless quarantine does not trigger DownSelfQuarantinedByRemote Motivation: akka/akka-core#32117 (backport of akka/akka#32050, issue akka/akka#31095) fixed InboundQuarantineCheck replying with a Quarantined control message for later inbound messages from a harmlessly quarantined (gracefully shut down) node, which made the remote node down itself via DownSelfQuarantinedByRemote. Pekko already stops propagating harmless quarantines by default (#1555, propagate-harmless-quarantine-events = off), but the multi-jvm coverage for this scenario was missing. Modification: - DowningWhenOtherHasQuarantinedThisActorSystemSpec: quarantine the second node directly instead of relying on split brain resolver timing, add a fourth node and a test that a node shutting down during a network partition does not trigger ThisActorSystemQuarantinedEvent on the surviving side. - Association / InboundQuarantineCheck: log at info when a Quarantined control message is sent, and mention "Harmless quarantine" in the harmless quarantine log message, matching Akka. - TestContext: pass harmless = false explicitly. Result: Multi-jvm coverage of the DownSelfQuarantinedByRemote scenarios, including harmless quarantine during a network partition. No behavioural change beyond log messages. Tests: - sbt "project cluster" "MultiJvm/testOnly org.apache.pekko.cluster.DowningWhenOtherHasQuarantinedThisActorSystem" (4 nodes, 4 tests each, all passed) - sbt "remote/testOnly org.apache.pekko.remote.artery.HarmlessQuarantineSpec org.apache.pekko.remote.artery.OutboundIdleShutdownSpec org.apache.pekko.remote.artery.RemoteMessageSerializationSpec" (14 passed) - native scalafmt run on the changed files References: Refs #1555 - port of akka/akka-core#32117 (the InboundQuarantineCheck fix was already present) --- ...herHasQuarantinedThisActorSystemSpec.scala | 116 ++++++++++++------ .../pekko/remote/artery/Association.scala | 5 +- .../artery/InboundQuarantineCheck.scala | 4 +- .../pekko/remote/artery/TestContext.scala | 2 +- 4 files changed, 85 insertions(+), 42 deletions(-) diff --git a/cluster/src/multi-jvm/scala/org/apache/pekko/cluster/DowningWhenOtherHasQuarantinedThisActorSystemSpec.scala b/cluster/src/multi-jvm/scala/org/apache/pekko/cluster/DowningWhenOtherHasQuarantinedThisActorSystemSpec.scala index 0e882f2bd9f..22eb1ac0f14 100644 --- a/cluster/src/multi-jvm/scala/org/apache/pekko/cluster/DowningWhenOtherHasQuarantinedThisActorSystemSpec.scala +++ b/cluster/src/multi-jvm/scala/org/apache/pekko/cluster/DowningWhenOtherHasQuarantinedThisActorSystemSpec.scala @@ -16,14 +16,18 @@ package org.apache.pekko.cluster import scala.concurrent.duration._ import org.apache.pekko +import pekko.actor.ActorIdentity import pekko.actor.ActorRef import pekko.actor.Identify import pekko.actor.RootActorPath +import pekko.remote.RARP import pekko.remote.artery.ArterySettings import pekko.remote.artery.ThisActorSystemQuarantinedEvent import pekko.remote.testkit.MultiNodeConfig import pekko.remote.transport.ThrottlerTransportAdapter +import pekko.testkit.EventFilter import pekko.testkit.LongRunningTest +import pekko.testkit.TestEvent.Mute import com.typesafe.config.ConfigFactory @@ -31,6 +35,7 @@ object DowningWhenOtherHasQuarantinedThisActorSystemSpec extends MultiNodeConfig val first = role("first") val second = role("second") val third = role("third") + val fourth = role("fourth") commonConfig( debugConfig(on = false) @@ -39,15 +44,9 @@ object DowningWhenOtherHasQuarantinedThisActorSystemSpec extends MultiNodeConfig ConfigFactory.parseString(""" pekko.remote.artery.enabled = on pekko.cluster.downing-provider-class = "org.apache.pekko.cluster.sbr.SplitBrainResolverProvider" - # speed up decision - pekko.cluster.split-brain-resolver.stable-after = 5s + pekko.cluster.split-brain-resolver.stable-after = 10s """))) - // exaggerate the timing issue by ,making the second node decide slower - // this is to more consistently repeat the scenario where the other side completes downing - // while the isolated part still has not made a decision and then see quarantined connections from the other nodes - nodeConfig(second)(ConfigFactory.parseString("pekko.cluster.split-brain-resolver.stable-after = 15s")) - testTransport(on = true) } @@ -57,11 +56,16 @@ class DowningWhenOtherHasQuarantinedThisActorSystemMultiJvmNode2 extends DowningWhenOtherHasQuarantinedThisActorSystemSpec class DowningWhenOtherHasQuarantinedThisActorSystemMultiJvmNode3 extends DowningWhenOtherHasQuarantinedThisActorSystemSpec +class DowningWhenOtherHasQuarantinedThisActorSystemMultiJvmNode4 + extends DowningWhenOtherHasQuarantinedThisActorSystemSpec abstract class DowningWhenOtherHasQuarantinedThisActorSystemSpec extends MultiNodeClusterSpec(DowningWhenOtherHasQuarantinedThisActorSystemSpec) { import DowningWhenOtherHasQuarantinedThisActorSystemSpec._ + muteDeadLetters(classOf[ActorIdentity])() + system.eventStream.publish(Mute(EventFilter.info(pattern = ".*Ignoring received gossip from unknown.*"))) + "Cluster node downed by other" must { if (!ArterySettings(system.settings.config.getConfig("pekko.remote.artery")).Enabled) { @@ -71,43 +75,23 @@ abstract class DowningWhenOtherHasQuarantinedThisActorSystemSpec } "join cluster" taggedAs LongRunningTest in { - awaitClusterUp(first, second, third) + awaitClusterUp(first, second, third, fourth) enterBarrier("after-1") } - "down itself" taggedAs LongRunningTest in { - runOn(first) { - testConductor.blackhole(first, second, ThrottlerTransportAdapter.Direction.Both).await - testConductor.blackhole(third, second, ThrottlerTransportAdapter.Direction.Both).await - } - enterBarrier("blackhole") - - within(15.seconds) { - runOn(first) { - awaitAssert { - cluster.state.unreachable.map(_.address) should ===(Set(address(second))) - } - awaitAssert { - // second downed and removed - cluster.state.members.map(_.address) should ===(Set(address(first), address(third))) - } - } - runOn(second) { - awaitAssert { - cluster.state.unreachable.map(_.address) should ===(Set(address(first), address(third))) - } - } + "down itself with DownSelfQuarantinedByRemote when other has quarantined" taggedAs LongRunningTest in { + runOn(first, second) { + system.eventStream.subscribe(testActor, classOf[ThisActorSystemQuarantinedEvent]) } - enterBarrier("down-second") - runOn(first) { - testConductor.passThrough(first, second, ThrottlerTransportAdapter.Direction.Both).await - testConductor.passThrough(third, second, ThrottlerTransportAdapter.Direction.Both).await + val secondUniqueAddress = cluster.state.members.find(_.address == address(second)).get.uniqueAddress + RARP(system).provider + .quarantine(secondUniqueAddress.address, Some(secondUniqueAddress.longUid), "Quarantine from test") } - enterBarrier("pass-through") + enterBarrier("quarantined") runOn(second) { - within(10.seconds) { + within(5.seconds) { // this is shorter than split-brain-resolver.stable-after, so it's not normal downing awaitAssert { // try to ping first (Cluster Heartbeat messages will not trigger the Quarantine message) system.actorSelection(RootActorPath(first) / "user").tell(Identify(None), ActorRef.noSender) @@ -115,6 +99,21 @@ abstract class DowningWhenOtherHasQuarantinedThisActorSystemSpec cluster.isTerminated should ===(true) } } + expectMsgType[ThisActorSystemQuarantinedEvent] + } + enterBarrier("second-shutdown") + + runOn(first) { + expectNoMessage(1.second) // no ThisActorSystemQuarantinedEvent + } + enterBarrier("wait") + + runOn(first) { + val sel = system.actorSelection(RootActorPath(second) / "user") + (1 to 15).foreach { _ => + sel.tell(Identify(None), ActorRef.noSender) // try to ping second + expectNoMessage(200.millis) // no ThisActorSystemQuarantinedEvent + } } enterBarrier("after-2") @@ -130,15 +129,56 @@ abstract class DowningWhenOtherHasQuarantinedThisActorSystemSpec cluster.shutdown() } + runOn(first) { + expectNoMessage(1.second) // no ThisActorSystemQuarantinedEvent + } + enterBarrier("wait") + runOn(first) { val sel = system.actorSelection(RootActorPath(third) / "user") - (1 to 25).foreach { _ => + (1 to 15).foreach { _ => sel.tell(Identify(None), ActorRef.noSender) // try to ping third expectNoMessage(200.millis) // no ThisActorSystemQuarantinedEvent } } - enterBarrier("after-2") + enterBarrier("after-3") + } + + "not be triggered by another node shutting down during network partition" taggedAs LongRunningTest in { + runOn(first) { + system.eventStream.subscribe(testActor, classOf[ThisActorSystemQuarantinedEvent]) + } + enterBarrier("subscribing") + + runOn(first) { + testConductor.blackhole(first, fourth, ThrottlerTransportAdapter.Direction.Both).await + } + enterBarrier("blackhole") + + runOn(third) { + cluster.shutdown() + } + + runOn(first) { + expectNoMessage(2.second) // no ThisActorSystemQuarantinedEvent + testConductor.passThrough(first, fourth, ThrottlerTransportAdapter.Direction.Both).await + } + + runOn(first) { + expectNoMessage(1.second) // no ThisActorSystemQuarantinedEvent + } + enterBarrier("wait") + + runOn(first) { + val sel = system.actorSelection(RootActorPath(fourth) / "user") + (1 to 15).foreach { _ => + sel.tell(Identify(None), ActorRef.noSender) // try to ping fourth + expectNoMessage(200.millis) // no ThisActorSystemQuarantinedEvent + } + } + + enterBarrier("after-4") } } diff --git a/remote/src/main/scala/org/apache/pekko/remote/artery/Association.scala b/remote/src/main/scala/org/apache/pekko/remote/artery/Association.scala index f0fc9350366..19f6e3a0259 100644 --- a/remote/src/main/scala/org/apache/pekko/remote/artery/Association.scala +++ b/remote/src/main/scala/org/apache/pekko/remote/artery/Association.scala @@ -560,8 +560,8 @@ private[remote] class Association( // quarantine state change was performed if (harmless) { log.info( - "Association to [{}] having UID [{}] has been stopped. All " + - "messages to this UID will be delivered to dead letters. Reason: {}", + "Association to [{}] having UID [{}] has been stopped. Harmless quarantine. " + + "All messages to this UID will be delivered to dead letters. Reason: {}", remoteAddress, u, reason) @@ -585,6 +585,7 @@ private[remote] class Association( send(ClearSystemMessageDelivery(current.incarnation), OptionVal.None, OptionVal.None) if (!harmless) { // try to tell the other system that we have quarantined it + log.info("Sending Quarantined to [{}]", peer) sendControl(Quarantined(localAddress, peer)) } setupStopQuarantinedTimer() diff --git a/remote/src/main/scala/org/apache/pekko/remote/artery/InboundQuarantineCheck.scala b/remote/src/main/scala/org/apache/pekko/remote/artery/InboundQuarantineCheck.scala index c6d20941497..b6e90028892 100644 --- a/remote/src/main/scala/org/apache/pekko/remote/artery/InboundQuarantineCheck.scala +++ b/remote/src/main/scala/org/apache/pekko/remote/artery/InboundQuarantineCheck.scala @@ -61,10 +61,12 @@ private[remote] class InboundQuarantineCheck(inboundContext: InboundContext) association.remoteAddress, env.originUid) // avoid starting outbound stream for heartbeats - if (!env.message.isInstanceOf[Quarantined] && !isHeartbeat(env.message)) + if (!env.message.isInstanceOf[Quarantined] && !isHeartbeat(env.message)) { + log.info("Sending Quarantined to [{}]", association.remoteAddress) inboundContext.sendControl( association.remoteAddress, Quarantined(inboundContext.localAddress, UniqueAddress(association.remoteAddress, env.originUid))) + } } pull(in) } else diff --git a/remote/src/test/scala/org/apache/pekko/remote/artery/TestContext.scala b/remote/src/test/scala/org/apache/pekko/remote/artery/TestContext.scala index 0212214d836..e44960b55a1 100644 --- a/remote/src/test/scala/org/apache/pekko/remote/artery/TestContext.scala +++ b/remote/src/test/scala/org/apache/pekko/remote/artery/TestContext.scala @@ -103,7 +103,7 @@ private[remote] class TestOutboundContext( } override def quarantine(reason: String): Unit = synchronized { - _associationState = _associationState.newQuarantined() + _associationState = _associationState.newQuarantined(harmless = false) } override def isOrdinaryMessageStreamActive(): Boolean = true