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 ac34d11ad7 fix: enforce resizer bounds and minimum pool size for
AdjustPoolSize (#3393)
ac34d11ad7 is described below
commit ac34d11ad7b0b1a5f5a88c0fe28ea025821d4c44
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Wed Jul 29 01:17:24 2026 +0800
fix: enforce resizer bounds and minimum pool size for AdjustPoolSize (#3393)
* fix: enforce resizer bounds and minimum pool size for AdjustPoolSize
Motivation:
AdjustPoolSize bypasses pool configuration bounds. When a resizer is
configured, the pool can grow beyond upperBound or shrink below
lowerBound. Without a resizer, the pool can shrink to zero routees.
Modification:
Add a clamp method to the Resizer trait (default: min 1) and override
it in DefaultResizer and DefaultOptimalSizeExploringResizer to enforce
their lowerBound/upperBound. ResizablePoolActor now intercepts
AdjustPoolSize and clamps the resulting size. RouterPoolActor enforces
a minimum of 1 routee for non-resizer pools.
Result:
AdjustPoolSize respects resizer bounds when configured, and never
shrinks a pool below 1 routee otherwise.
Tests:
- sbt "actor-tests / Test / testOnly
org.apache.pekko.routing.RoundRobinSpec" passed
- sbt "actor-tests / Test / testOnly org.apache.pekko.routing.ResizerSpec"
passed
- sbt "actor / mimaReportBinaryIssues" passed
References:
Fixes #3253
* fix: add @since 2.0.0 to Resizer.clamp Scaladoc
Motivation:
Reviewer requested version annotation on the new public API method.
Modification:
Add @since 2.0.0 tag to the clamp method Scaladoc in Resizer trait.
Result:
New API method is properly annotated with its introduction version.
Tests:
Not run - docs only
References:
Refs #3253
---
.../org/apache/pekko/routing/ResizerSpec.scala | 35 ++++++++++++++++++++++
.../org/apache/pekko/routing/RoundRobinSpec.scala | 10 +++++++
.../routing/OptimalSizeExploringResizer.scala | 3 ++
.../scala/org/apache/pekko/routing/Resizer.scala | 26 ++++++++++++++++
.../org/apache/pekko/routing/RoutedActorCell.scala | 12 ++++----
.../org/apache/pekko/routing/RouterConfig.scala | 9 ++----
6 files changed, 84 insertions(+), 11 deletions(-)
diff --git
a/actor-tests/src/test/scala/org/apache/pekko/routing/ResizerSpec.scala
b/actor-tests/src/test/scala/org/apache/pekko/routing/ResizerSpec.scala
index 198b7144cf..52284eec42 100644
--- a/actor-tests/src/test/scala/org/apache/pekko/routing/ResizerSpec.scala
+++ b/actor-tests/src/test/scala/org/apache/pekko/routing/ResizerSpec.scala
@@ -249,6 +249,41 @@ class ResizerSpec extends PekkoSpec(ResizerSpec.config)
with DefaultTimeout with
}
+ "clamp AdjustPoolSize growth to upperBound" in {
+ val resizer = DefaultResizer(lowerBound = 2, upperBound = 5,
messagesPerResize = 100)
+ val router = system.actorOf(
+ RoundRobinPool(nrOfInstances = 0, resizer =
Some(resizer)).props(Props[TestActor]()))
+
+ val latch = new TestLatch(2)
+ router ! latch
+ router ! latch
+ Await.ready(latch, remainingOrDefault)
+
+ routeeSize(router) should ===(2)
+
+ router ! AdjustPoolSize(10)
+ routeeSize(router) should ===(5)
+ }
+
+ "clamp AdjustPoolSize shrink to lowerBound" in {
+ val resizer = DefaultResizer(lowerBound = 2, upperBound = 5,
messagesPerResize = 100)
+ val router = system.actorOf(
+ RoundRobinPool(nrOfInstances = 0, resizer =
Some(resizer)).props(Props[TestActor]()))
+
+ val latch = new TestLatch(2)
+ router ! latch
+ router ! latch
+ Await.ready(latch, remainingOrDefault)
+
+ routeeSize(router) should ===(2)
+
+ router ! AdjustPoolSize(2)
+ routeeSize(router) should ===(4)
+
+ router ! AdjustPoolSize(-10)
+ routeeSize(router) should ===(2)
+ }
+
}
}
diff --git
a/actor-tests/src/test/scala/org/apache/pekko/routing/RoundRobinSpec.scala
b/actor-tests/src/test/scala/org/apache/pekko/routing/RoundRobinSpec.scala
index 2e47538c5d..a7d80b29e5 100644
--- a/actor-tests/src/test/scala/org/apache/pekko/routing/RoundRobinSpec.scala
+++ b/actor-tests/src/test/scala/org/apache/pekko/routing/RoundRobinSpec.scala
@@ -125,6 +125,16 @@ class RoundRobinSpec extends PekkoSpec with DefaultTimeout
with ImplicitSender {
actor ! RemoveRoutee(other)
routeeSize(actor) should ===(5)
}
+
+ "not shrink below 1 routee with AdjustPoolSize" in {
+ val actor = system.actorOf(RoundRobinPool(3).props(routeeProps =
Props(new Actor {
+ def receive = Actor.emptyBehavior
+ })), "round-robin-min-one")
+
+ routeeSize(actor) should ===(3)
+ actor ! AdjustPoolSize(-10)
+ routeeSize(actor) should ===(1)
+ }
}
"round robin group" must {
diff --git
a/actor/src/main/scala/org/apache/pekko/routing/OptimalSizeExploringResizer.scala
b/actor/src/main/scala/org/apache/pekko/routing/OptimalSizeExploringResizer.scala
index de1e5bf192..fb33004c59 100644
---
a/actor/src/main/scala/org/apache/pekko/routing/OptimalSizeExploringResizer.scala
+++
b/actor/src/main/scala/org/apache/pekko/routing/OptimalSizeExploringResizer.scala
@@ -284,6 +284,9 @@ case class DefaultOptimalSizeExploringResizer(
Math.max(lowerBound, Math.min(proposedChange + currentSize, upperBound)) -
currentSize
}
+ override def clamp(proposedSize: Int): Int =
+ Math.max(lowerBound, Math.min(proposedSize, upperBound))
+
private def optimize(currentSize: PoolSize): Int = {
val adjacentDispatchWaits: Map[PoolSize, Duration] = {
diff --git a/actor/src/main/scala/org/apache/pekko/routing/Resizer.scala
b/actor/src/main/scala/org/apache/pekko/routing/Resizer.scala
index 3208cd099d..0e0a1b2d79 100644
--- a/actor/src/main/scala/org/apache/pekko/routing/Resizer.scala
+++ b/actor/src/main/scala/org/apache/pekko/routing/Resizer.scala
@@ -63,6 +63,17 @@ trait Resizer {
*/
def resize(currentRoutees: immutable.IndexedSeq[Routee]): Int
+ /**
+ * Constrain a proposed total pool size to the bounds enforced by this
resizer.
+ * Used by [[AdjustPoolSize]] management messages to respect resizer bounds.
+ * The default implementation only ensures the size is at least 1.
+ *
+ * @param proposedSize the proposed total number of routees
+ * @return the clamped pool size
+ * @since 2.0.0
+ */
+ def clamp(proposedSize: Int): Int = math.max(1, proposedSize)
+
}
object Resizer {
@@ -188,6 +199,9 @@ case class DefaultResizer(
else delta
}
+ override def clamp(proposedSize: Int): Int =
+ math.max(lowerBound, math.min(proposedSize, upperBound))
+
/**
* Number of routees considered busy, or above 'pressure level'.
*
@@ -347,6 +361,18 @@ private[pekko] class
ResizablePoolActor(supervisorStrategy: SupervisorStrategy)
({
case Resize =>
resizerCell.resize(initial = false)
+ case AdjustPoolSize(change: Int) =>
+ val currentRoutees = cell.router.routees
+ val currentSize = currentRoutees.size
+ val clampedSize = resizerCell.resizer.clamp(currentSize + change)
+ val effectiveChange = clampedSize - currentSize
+ if (effectiveChange > 0) {
+ val newRoutees =
Vector.fill(effectiveChange)(pool.newRoutee(cell.routeeProps, context))
+ cell.addRoutees(newRoutees)
+ } else if (effectiveChange < 0) {
+ val abandon = currentRoutees.drop(currentRoutees.length +
effectiveChange)
+ cell.removeRoutees(abandon, stopChild = true)
+ }
}: Actor.Receive).orElse(super.receive)
}
diff --git
a/actor/src/main/scala/org/apache/pekko/routing/RoutedActorCell.scala
b/actor/src/main/scala/org/apache/pekko/routing/RoutedActorCell.scala
index 5150037ef5..9344b22578 100644
--- a/actor/src/main/scala/org/apache/pekko/routing/RoutedActorCell.scala
+++ b/actor/src/main/scala/org/apache/pekko/routing/RoutedActorCell.scala
@@ -203,12 +203,14 @@ private[pekko] class RouterPoolActor(override val
supervisorStrategy: Supervisor
override def receive =
({
case AdjustPoolSize(change: Int) =>
- if (change > 0) {
- val newRoutees =
Vector.fill(change)(pool.newRoutee(cell.routeeProps, context))
+ val currentRoutees = cell.router.routees
+ val currentSize = currentRoutees.size
+ val effectiveChange = math.max(1, currentSize + change) - currentSize
+ if (effectiveChange > 0) {
+ val newRoutees =
Vector.fill(effectiveChange)(pool.newRoutee(cell.routeeProps, context))
cell.addRoutees(newRoutees)
- } else if (change < 0) {
- val currentRoutees = cell.router.routees
- val abandon = currentRoutees.drop(currentRoutees.length + change)
+ } else if (effectiveChange < 0) {
+ val abandon = currentRoutees.drop(currentRoutees.length +
effectiveChange)
cell.removeRoutees(abandon, stopChild = true)
}
}: Actor.Receive).orElse(super.receive)
diff --git a/actor/src/main/scala/org/apache/pekko/routing/RouterConfig.scala
b/actor/src/main/scala/org/apache/pekko/routing/RouterConfig.scala
index 2dbd45a14a..e579de965f 100644
--- a/actor/src/main/scala/org/apache/pekko/routing/RouterConfig.scala
+++ b/actor/src/main/scala/org/apache/pekko/routing/RouterConfig.scala
@@ -444,12 +444,9 @@ final case class RemoveRoutee(routee: Routee) extends
RouterManagementMesssage
*
* Positive `change` will add that number of routees to the [[Pool]].
* Negative `change` will remove that number of routees from the [[Pool]].
- * This is a direct management operation, so `change` is not constrained by the
- * pool's initial `nrOfInstances` or by bounds configured on a [[Resizer]].
- * A resizer may adjust the pool back within its bounds the next time it runs.
- *
- * If a resizable pool is reduced to zero routees, the ordinary message that
- * triggers the next resize may be lost before new routees are added.
+ * If a [[Resizer]] is configured, the resulting pool size is clamped to the
+ * resizer's bounds (see [[Resizer.clamp]]). Without a resizer, the pool will
+ * not shrink below 1 routee.
*
* Routees are stopped by sending a [[pekko.actor.PoisonPill]] to the routee.
* Precautions are taken reduce the risk of dropping messages that are
concurrently
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]