Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -16,21 +16,26 @@ 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

object DowningWhenOtherHasQuarantinedThisActorSystemSpec extends MultiNodeConfig {
val first = role("first")
val second = role("second")
val third = role("third")
val fourth = role("fourth")

commonConfig(
debugConfig(on = false)
Expand All @@ -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)
}

Expand All @@ -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) {
Expand All @@ -71,50 +75,45 @@ 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)
// shutting down itself triggered by ThisActorSystemQuarantinedEvent
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")
Expand All @@ -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")
}

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading