This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-http.git
The following commit(s) were added to refs/heads/main by this push:
new 561929f96 feat: add client-side max connection age for HTTP/2 (#1320)
561929f96 is described below
commit 561929f9659961d89ed7e6e5b97b701362c5a832
Author: Rayan-and-beyond <[email protected]>
AuthorDate: Sat Oct 3 13:58:36 2026 +0300
feat: add client-side max connection age for HTTP/2 (#1320)
* http2: add persistent client connection max age #1319
Motivation:
Long-lived managed HTTP/2 clients can stay pinned to existing server
instances.
Modification:
Add a configurable maximum age that retires managed persistent HTTP/2
connections after in-flight requests drain.
Result:
Requests arriving after retirement establish a fresh connection while
in-flight requests complete normally.
Tests:
- sbt validatePullRequest
- sbt http2-tests/test
- sbt +http-core/mimaReportBinaryIssues
- sbt scalafmtCheckAll scalafmtSbtCheck
- sbt +headerCheckAll
- sbt docs/paradox
- sbt checkCodeStyle
- git diff --check
References:
Refs #1319
* http2: address max-age review feedback
Expose client max-age and jitter as public settings, add per-connection
jitter, and preserve a buffered request when the request source completes
during retirement. Document break-before-make behavior and extend
retirement/reconnect coverage.
* http2: preserve Java max-age precision
* Align client max connection age with the server-side setting from #1316
Motivation:
#1316 merged with `infinite` as the disabled value for
`max-connection-age`, a `Duration` setting type and
`JavaDurationConverter` for the Java API. The client setting used `0s`
and `FiniteDuration`.
Modification:
- `persistent-connection-max-age` defaults to `infinite`, is a
`Duration`, and must be > 0 or `infinite`, as on the server
- Java accessors use `JavaDurationConverter`, so
`ChronoUnit.FOREVER.getDuration` round-trips to `Duration.Inf`
- reference.conf, scaladoc and docs follow the server wording
- the jitter scheduling matches `Http2Demux`
- settings tests move to a new `Http2ClientSettingsSpec`, mirroring
`Http2ServerSettingsSpec`
- MiMa excludes are reduced to the abstract members that need them
- pekko-style imports in the changed files
Result:
The client and server max connection age settings share naming
conventions, defaults, validation and Java conversion behaviour.
Tests:
- sbt "http-core/testOnly ...Http2ClientSettingsSpec
...Http2CommonSettingsSpec ...Http2ServerSettingsSpec": 16 passed
- sbt "http2-tests/testOnly ...Http2PersistentClient*": 26 passed
- sbt "http-core/mimaReportBinaryIssues": clean
- sbt scalafmtAll and headerCreateAll run on changed modules
References:
Refs #1319, #1316
---------
Co-authored-by: PJ Fanning <[email protected]>
---
docs/src/main/paradox/client-side/http2.md | 29 +++++++
.../http2-client-max-connection-age.excludes | 21 +++++
http-core/src/main/resources/reference.conf | 23 ++++-
.../pekko/http/impl/engine/http2/Http2Demux.scala | 3 +-
.../engine/http2/client/PersistentConnection.scala | 65 +++++++++++++--
.../javadsl/settings/Http2ClientSettings.scala | 38 ++++++++-
.../scaladsl/settings/Http2ServerSettings.scala | 40 +++++++++
.../settings/Http2ClientSettingsSpec.scala | 97 ++++++++++++++++++++++
.../engine/http2/Http2PersistentClientSpec.scala | 68 +++++++++++++++
9 files changed, 372 insertions(+), 12 deletions(-)
diff --git a/docs/src/main/paradox/client-side/http2.md
b/docs/src/main/paradox/client-side/http2.md
index 69261215e..125844621 100644
--- a/docs/src/main/paradox/client-side/http2.md
+++ b/docs/src/main/paradox/client-side/http2.md
@@ -55,6 +55,35 @@ Java
The Apache Pekko HTTP client doesn't support HTTP/1 to HTTP/2 negotiation over
plaintext using the `Upgrade` mechanism.
+## Limiting managed persistent connection lifetime
+
+Managed persistent HTTP/2 clients can periodically retire long-lived
connections so that later requests establish
+fresh connections. This is useful when clients would otherwise remain pinned
to the same server instances after
+a scale-out or rolling deployment.
+
+Configure the maximum age and optional jitter under the HTTP/2 client settings:
+
+```
+pekko.http.client.http2.persistent-connection-max-age = 10m
+pekko.http.client.http2.persistent-connection-max-age-jitter = 0.1
+```
+
+The setting applies to `managedPersistentHttp2()` and
`managedPersistentHttp2WithPriorKnowledge()`. Each
+connection gets an independently jittered age; with the default jitter of
`0.1`, retirement happens between 90%
+and 110% of the configured maximum age. Set the jitter to `0` to disable it.
+
+Retirement is currently **break-before-make**. When a connection reaches its
age, the managed client stops
+assigning new requests to it and lets requests already in flight drain before
closing the connection. New requests
+are backpressured until the old connection disconnects, then a replacement
connection is established. As a result,
+a retiring connection can add its drain time plus TCP/TLS/HTTP2 setup time to
new-request latency.
+
+The existing `completion-timeout` bounds the drain period. If that timeout
expires, the old connection is closed
+even when requests are still in flight, terminating long-lived requests such
as streaming responses.
+
+The default maximum age is `infinite`, which disables age-based retirement.
The settings can also be changed
+programmatically with `Http2ClientSettings.withPersistentConnectionMaxAge` and
+`Http2ClientSettings.withPersistentConnectionMaxAgeJitter`.
+
## Request-response ordering
For HTTP/2 connections the responses are not guaranteed to arrive in the same
order that the requests were emitted to
diff --git
a/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-client-max-connection-age.excludes
b/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-client-max-connection-age.excludes
new file mode 100644
index 000000000..be4f838ef
--- /dev/null
+++
b/http-core/src/main/mima-filters/2.0.x.backwards.excludes/http2-client-max-connection-age.excludes
@@ -0,0 +1,21 @@
+# 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.
+
+# new persistent connection max-age settings for HTTP/2 clients
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ClientSettings.persistentConnectionMaxAge")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.scaladsl.settings.Http2ClientSettings.persistentConnectionMaxAgeJitter")
+ProblemFilters.exclude[ReversedMissingMethodProblem]("org.apache.pekko.http.javadsl.settings.Http2ClientSettings.withPersistentConnectionMaxAgeJitter")
diff --git a/http-core/src/main/resources/reference.conf
b/http-core/src/main/resources/reference.conf
index e591cead0..7fe20b171 100644
--- a/http-core/src/main/resources/reference.conf
+++ b/http-core/src/main/resources/reference.conf
@@ -633,6 +633,26 @@ pekko.http {
# Set to zero to retry indefinitely.
max-persistent-attempts = 0
+ # The maximum time a connection created by managedPersistentHttp2 or
managedPersistentHttp2WithPriorKnowledge
+ # is used before it is retired. When the age of a connection exceeds
this value, the connection stops
+ # accepting new requests, lets requests that are already in flight
complete (see `completion-timeout`), and
+ # is then closed. Later requests are sent on a new connection.
+ #
+ # Retiring connections regularly helps to rebalance long-lived HTTP/2
connections (as used by gRPC) across
+ # server instances, for example after a scale-out or a rolling deploy.
+ #
+ # Retirement is break-before-make: new requests wait for the retiring
connection to close before a
+ # new connection is established, so they can observe the drain plus the
connection setup latency.
+ #
+ # The value `infinite` disables this mechanism and is the default.
+ persistent-connection-max-age = infinite
+
+ # The jitter applied to `persistent-connection-max-age`, as a fraction
of the configured value: with the
+ # default of 0.1 each connection is retired after between 90% and 110%
of the configured age, so that
+ # connections that were opened together are not all retired at the same
time.
+ # Set to 0 to disable jitter. Must be >= 0 and < 1.
+ persistent-connection-max-age-jitter = 0.1
+
# Starting backoff before reconnecting when a persistent HTTP/2 client
connection fails
# see `pekko.http.host-connection-pool.base-connection-backoff` for
details on the backoff mechanism.
base-connection-backoff =
${pekko.http.host-connection-pool.base-connection-backoff}
@@ -642,7 +662,8 @@ pekko.http {
max-connection-backoff =
${pekko.http.host-connection-pool.max-connection-backoff}
# When gracefully closing the HTTP/2 client, await at most
`completion-timeout` for in-flight
- # requests to complete.
+ # requests to complete. When the timeout expires, remaining in-flight
requests are terminated.
+ # This also bounds the drain of a connection retired after
`persistent-connection-max-age`.
completion-timeout = 3s
}
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
index 4e3cd4854..cee6b45d6 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/Http2Demux.scala
@@ -72,7 +72,8 @@ private[http2] class Http2ClientDemux(http2Settings:
Http2ClientSettings, master
override def completionTimeout: FiniteDuration =
http2Settings.completionTimeout
- // a maximum connection age is not supported on the client side
+ // the client side limits the age of managed persistent connections in
PersistentConnection instead,
+ // see `persistent-connection-max-age`
def maxConnectionAge: Duration = Duration.Inf
def maxConnectionAgeGrace: Duration = Duration.Inf
def maxConnectionAgeJitter: Double = 0.0
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
index 6f8f8db39..1731ab6e8 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/http2/client/PersistentConnection.scala
@@ -35,6 +35,7 @@ import scala.util.{ Failure, Success }
private[http2] object PersistentConnection {
private case class EmbargoEnded(connectsLeft: Option[Int], embargo:
FiniteDuration)
+ private case object MaxConnectionAgeReached
/**
* Wraps a connection flow with transparent reconnection support.
@@ -59,7 +60,8 @@ private[http2] object PersistentConnection {
settings.maxPersistentAttempts match {
case 0 => None
case n => Some(n)
- }, settings.baseConnectionBackoff, settings.maxConnectionBackoff))
+ }, settings.baseConnectionBackoff, settings.maxConnectionBackoff,
+ settings.persistentConnectionMaxAge,
settings.persistentConnectionMaxAgeJitter))
private class AssociationTag extends RequestResponseAssociation
private val associationTagKey =
AttributeKey[AssociationTag]("PersistentConnection.associationTagKey")
@@ -69,7 +71,8 @@ private[http2] object PersistentConnection {
entity = "The server closed the connection before delivering a
response.")
private class Stage(connectionFlow: Flow[HttpRequest, HttpResponse,
Future[OutgoingConnection]],
- maxAttempts: Option[Int], baseEmbargo: FiniteDuration, _maxBackoff:
FiniteDuration)
+ maxAttempts: Option[Int], baseEmbargo: FiniteDuration, _maxBackoff:
FiniteDuration,
+ persistentConnectionMaxAge: Duration, persistentConnectionMaxAgeJitter:
Double)
extends GraphStage[FlowShape[HttpRequest, HttpResponse]] {
val requestIn = Inlet[HttpRequest]("PersistentConnection.requestIn")
val responseOut = Outlet[HttpResponse]("PersistentConnection.responseOut")
@@ -78,6 +81,8 @@ private[http2] object PersistentConnection {
val shape: FlowShape[HttpRequest, HttpResponse] = FlowShape(requestIn,
responseOut)
override def createLogic(inheritedAttributes: Attributes): GraphStageLogic
=
new TimerGraphStageLogic(shape) with StageLogging {
+ private var currentConnection: Option[Connected] = None
+
become(Unconnected)
def become(state: State): Unit = setHandlers(requestIn, responseOut,
state)
@@ -86,7 +91,9 @@ private[http2] object PersistentConnection {
object Unconnected extends State {
override def onPush(): Unit = connect(maxAttempts, Duration.Zero)
override def onPull(): Unit =
- if (!isAvailable(requestIn) && !hasBeenPulled(requestIn)) //
requestIn might already have been pulled when we failed and went back to
Unconnected
+ if (isAvailable(requestIn)) connect(maxAttempts, Duration.Zero)
+ else if (isClosed(requestIn)) completeStage()
+ else if (!hasBeenPulled(requestIn)) // requestIn might already
have been pulled when we failed and went back to Unconnected
pull(requestIn)
}
@@ -131,12 +138,13 @@ private[http2] object PersistentConnection {
override def onPush(): Unit = () // Pull might have happened before
the connection failed. Element is kept in slot.
override def onPull(): Unit = {
- if (!isAvailable(requestIn) && !hasBeenPulled(requestIn)) //
requestIn might already have been pulled when we failed and went back to
Unconnected
+ if (!isAvailable(requestIn) && !isClosed(requestIn) &&
!hasBeenPulled(requestIn)) // requestIn might already have been pulled when we
failed and went back to Unconnected
pull(requestIn)
}
val onConnected = getAsyncCallback[Unit] { _ =>
val newState = new Connected(requestOut, responseIn)
+ currentConnection = Some(newState)
become(newState)
if (requestOutPulled) {
if (isAvailable(requestIn))
newState.dispatchRequest(grab(requestIn))
@@ -181,6 +189,8 @@ private[http2] object PersistentConnection {
case EmbargoEnded(connectsLeft, nextEmbargo) =>
log.debug("Reconnecting after backoff")
connect(connectsLeft, nextEmbargo)
+ case MaxConnectionAgeReached =>
+ currentConnection.foreach(_.retire())
}
}
@@ -188,12 +198,27 @@ private[http2] object PersistentConnection {
requestOut: SubSourceOutlet[HttpRequest],
responseIn: SubSinkInlet[HttpResponse]) extends State {
private var ongoingRequests: Map[AssociationTag,
Map[AttributeKey[?], RequestResponseAssociation]] = Map.empty
+ private var retiring = false
+
+ persistentConnectionMaxAge match {
+ case age: FiniteDuration =>
+ // The age of each connection is jittered so that connections
that were opened together are not
+ // all retired at the same time, see
`persistent-connection-max-age-jitter` in the configuration.
+ val jitterFactor =
+ 1.0 + persistentConnectionMaxAgeJitter * (2 *
ThreadLocalRandom.current().nextDouble() - 1)
+ scheduleOnce(MaxConnectionAgeReached, (age.toMillis *
jitterFactor).toLong.max(1L).millis)
+ case _ => // no maximum connection age configured
+ }
+
responseIn.pull()
requestOut.setHandler(new OutHandler {
override def onPull(): Unit =
- if (!isAvailable(requestIn)) pull(requestIn)
- else dispatchRequest(grab(requestIn))
+ if (isAvailable(requestIn)) {
+ dispatchRequest(grab(requestIn))
+ if (isClosed(requestIn)) requestOut.complete()
+ } else if (isClosed(requestIn)) requestOut.complete()
+ else if (!hasBeenPulled(requestIn)) pull(requestIn)
override def onDownstreamFinish(cause: Throwable): Unit =
onDisconnected()
})
@@ -210,23 +235,45 @@ private[http2] object PersistentConnection {
override def onUpstreamFailure(ex: Throwable): Unit =
onDisconnected() // FIXME: log error
})
def onDisconnected(): Unit = {
+ cancelTimer(MaxConnectionAgeReached)
+ currentConnection = None
+
emitMultiple[HttpResponse](responseOut,
ongoingRequests.values.map(errorResponse.withAttributes(_)).toVector,
() => setHandler(responseOut, Unconnected))
responseIn.cancel()
requestOut.fail(new RuntimeException("connection broken"))
- if (isClosed(requestIn)) {
- // user closed PersistentConnection before and we were waiting
for remaining responses
+ if (isClosed(requestIn) && !isAvailable(requestIn)) {
+ // user closed PersistentConnection before and there is no
buffered request left to dispatch
completeStage()
} else {
// become(Unconnected) doesn't work because of using emit
// so we need to do it more carefully here
setHandler(requestIn, Unconnected)
- if (isAvailable(responseOut) && !hasBeenPulled(requestIn))
pull(requestIn)
+ if (isAvailable(responseOut)) {
+ if (isAvailable(requestIn)) connect(maxAttempts, Duration.Zero)
+ else if (!hasBeenPulled(requestIn)) pull(requestIn)
+ }
}
}
+ def retire(): Unit =
+ if (!retiring) {
+ retiring = true
+ log.debug("Persistent HTTP/2 connection reached its configured
maximum age, retiring it")
+ setHandler(requestIn,
+ new InHandler {
+ override def onPush(): Unit = () // keep at most one next
request in the inlet slot
+ override def onUpstreamFinish(): Unit = ()
+ override def onUpstreamFailure(ex: Throwable): Unit = {
+ responseIn.cancel()
+ failStage(ex)
+ }
+ })
+ requestOut.complete()
+ }
+
def dispatchRequest(req: HttpRequest): Unit = {
val tag = new AssociationTag
// Some cross-compilation woes here:
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
b/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
index 230413739..f7e6e5fec 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/javadsl/settings/Http2ClientSettings.scala
@@ -15,7 +15,9 @@ package org.apache.pekko.http.javadsl.settings
import java.time.Duration
-import org.apache.pekko.http.scaladsl
+import org.apache.pekko
+import pekko.http.impl.util.JavaDurationConverter
+import pekko.http.scaladsl
import scala.concurrent.duration.DurationLong
@@ -79,6 +81,40 @@ trait Http2ClientSettings { self:
scaladsl.settings.Http2ClientSettings.Http2Cli
def getMaxPersistentAttempts: Int = maxPersistentAttempts
def withMaxPersistentAttempts(max: Int): Http2ClientSettings =
copy(maxPersistentAttempts = max)
+ /**
+ * The maximum age of a connection created by `managedPersistentHttp2` or
+ * `managedPersistentHttp2WithPriorKnowledge`. When the age of a connection
exceeds this value, the connection
+ * stops accepting new requests, lets requests that are already in flight
complete within
+ * [[getCompletionTimeout]], and is then closed. Later requests are sent on
a new connection. The value
+ * `ChronoUnit.FOREVER.getDuration` represents an infinite age, which
disables this mechanism and is the
+ * default.
+ *
+ * @since 2.0.0
+ */
+ def getPersistentConnectionMaxAge: Duration =
JavaDurationConverter.toJava(persistentConnectionMaxAge)
+
+ /**
+ * Pass `ChronoUnit.FOREVER.getDuration` to disable the maximum connection
age.
+ *
+ * @since 2.0.0
+ */
+ def withPersistentConnectionMaxAge(maxAge: Duration): Http2ClientSettings =
+ self.withPersistentConnectionMaxAge(JavaDurationConverter.toScala(maxAge))
+
+ /**
+ * The jitter applied to the persistent connection maximum age per
connection, as a fraction of the configured
+ * age: with the default of 0.1 each connection is retired after between 90%
and 110% of the configured age, so
+ * that connections that were opened together are not all retired at the
same time. 0 disables jitter.
+ *
+ * @since 2.0.0
+ */
+ def getPersistentConnectionMaxAgeJitter: Double =
persistentConnectionMaxAgeJitter
+
+ /**
+ * @since 2.0.0
+ */
+ def withPersistentConnectionMaxAgeJitter(jitter: Double): Http2ClientSettings
+
def getCompletionTimeout: Duration =
Duration.ofMillis(completionTimeout.toMillis)
def withCompletionTimeout(timeout: Duration): Http2ClientSettings =
copy(completionTimeout = timeout.toMillis.millis)
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
b/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
index 7d5b42bd3..e2616650d 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/scaladsl/settings/Http2ServerSettings.scala
@@ -347,6 +347,38 @@ trait Http2ClientSettings extends
javadsl.settings.Http2ClientSettings with Http
def maxPersistentAttempts: Int
override def withMaxPersistentAttempts(max: Int): Http2ClientSettings =
copy(maxPersistentAttempts = max)
+ /**
+ * The maximum age of a connection created by `managedPersistentHttp2` or
+ * `managedPersistentHttp2WithPriorKnowledge`. When the age of a connection
exceeds this value, the connection
+ * stops accepting new requests, lets requests that are already in flight
complete within [[completionTimeout]],
+ * and is then closed. Later requests are sent on a new connection. The
value `Duration.Inf` disables this
+ * mechanism and is the default.
+ *
+ * @since 2.0.0
+ */
+ def persistentConnectionMaxAge: Duration
+
+ /**
+ * @since 2.0.0
+ */
+ def withPersistentConnectionMaxAge(maxAge: Duration): Http2ClientSettings =
+ copy(persistentConnectionMaxAge = maxAge)
+
+ /**
+ * The jitter applied to [[persistentConnectionMaxAge]] per connection, as a
fraction of the configured age:
+ * with the default of 0.1 each connection is retired after between 90% and
110% of the configured age, so
+ * that connections that were opened together are not all retired at the
same time. 0 disables jitter.
+ *
+ * @since 2.0.0
+ */
+ def persistentConnectionMaxAgeJitter: Double
+
+ /**
+ * @since 2.0.0
+ */
+ override def withPersistentConnectionMaxAgeJitter(jitter: Double):
Http2ClientSettings =
+ copy(persistentConnectionMaxAgeJitter = jitter)
+
def completionTimeout: FiniteDuration
def withCompletionTimeout(timeout: FiniteDuration): Http2ClientSettings =
copy(completionTimeout = timeout)
@@ -380,6 +412,8 @@ object Http2ClientSettings extends
SettingsCompanion[Http2ClientSettings] {
pingInterval: FiniteDuration,
pingTimeout: FiniteDuration,
maxPersistentAttempts: Int,
+ persistentConnectionMaxAge: Duration,
+ persistentConnectionMaxAgeJitter: Double,
completionTimeout: FiniteDuration,
baseConnectionBackoff: FiniteDuration,
maxConnectionBackoff: FiniteDuration,
@@ -395,6 +429,10 @@ object Http2ClientSettings extends
SettingsCompanion[Http2ClientSettings] {
require(incomingStreamLevelBufferSize > 0,
"incoming-stream-level-buffer-size must be > 0")
require(outgoingControlFrameBufferSize > 0,
"outgoing-control-frame-buffer-size must be > 0")
require(maxPersistentAttempts >= 0, "max-persistent-attempts must be >= 0")
+ require(persistentConnectionMaxAge > Duration.Zero,
+ "persistent-connection-max-age must be > 0 or 'infinite' to disable")
+ require(persistentConnectionMaxAgeJitter >= 0 &&
persistentConnectionMaxAgeJitter < 1,
+ "persistent-connection-max-age-jitter must be >= 0 and < 1")
require(completionTimeout > Duration.Zero, "completion-timeout must be >
0")
require(baseConnectionBackoff <= maxConnectionBackoff,
"base-connection-backoff must be <= max-connection-backoff")
Http2CommonSettings.validate(this)
@@ -414,6 +452,8 @@ object Http2ClientSettings extends
SettingsCompanion[Http2ClientSettings] {
pingInterval = c.getFiniteDuration("ping-interval"),
pingTimeout = c.getFiniteDuration("ping-timeout"),
maxPersistentAttempts = c.getInt("max-persistent-attempts"),
+ persistentConnectionMaxAge =
c.getPotentiallyInfiniteDuration("persistent-connection-max-age"),
+ persistentConnectionMaxAgeJitter =
c.getDouble("persistent-connection-max-age-jitter"),
completionTimeout = c.getFiniteDuration("completion-timeout"),
baseConnectionBackoff = c.getFiniteDuration("base-connection-backoff"),
maxConnectionBackoff = c.getFiniteDuration("max-connection-backoff"),
diff --git
a/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ClientSettingsSpec.scala
b/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ClientSettingsSpec.scala
new file mode 100644
index 000000000..e4222e028
--- /dev/null
+++
b/http-core/src/test/scala/org/apache/pekko/http/scaladsl/settings/Http2ClientSettingsSpec.scala
@@ -0,0 +1,97 @@
+/*
+ * 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.http.scaladsl.settings
+
+import java.time.temporal.ChronoUnit
+
+import org.apache.pekko
+import pekko.testkit.PekkoSpec
+
+import scala.concurrent.duration._
+
+class Http2ClientSettingsSpec extends PekkoSpec {
+
+ "Http2ClientSettings persistent-connection-max-age" should {
+
+ "be disabled by default" in {
+ val settings = Http2ClientSettings(system)
+ settings.persistentConnectionMaxAge should ===(Duration.Inf)
+ settings.getPersistentConnectionMaxAge should
===(ChronoUnit.FOREVER.getDuration)
+ settings.internalSettings should ===(None)
+ }
+
+ "accept a finite value from config" in {
+ val settings =
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age =
2s")
+ settings.persistentConnectionMaxAge should ===(2.seconds)
+ settings.getPersistentConnectionMaxAge should
===(java.time.Duration.ofSeconds(2))
+ }
+
+ "round-trip an infinite value through the Java API" in {
+ val settings =
Http2ClientSettings(system).withPersistentConnectionMaxAge(Duration.Inf)
+ settings.getPersistentConnectionMaxAge should
===(ChronoUnit.FOREVER.getDuration)
+ val roundTripped =
settings.withPersistentConnectionMaxAge(settings.getPersistentConnectionMaxAge)
+ roundTripped.getPersistentConnectionMaxAge should
===(ChronoUnit.FOREVER.getDuration)
+ }
+
+ "round-trip a finite value through the Java API" in {
+ val settings =
Http2ClientSettings(system).withPersistentConnectionMaxAge(2.minutes)
+ settings.getPersistentConnectionMaxAge should
===(java.time.Duration.ofMinutes(2))
+ val roundTripped =
settings.withPersistentConnectionMaxAge(settings.getPersistentConnectionMaxAge)
+ roundTripped.getPersistentConnectionMaxAge should
===(java.time.Duration.ofMinutes(2))
+ }
+
+ "keep sub-millisecond precision through the Java API" in {
+ val age = java.time.Duration.ofNanos(1)
+
Http2ClientSettings(system).withPersistentConnectionMaxAge(age).getPersistentConnectionMaxAge
should ===(age)
+ }
+
+ "reject a zero or negative value" in {
+ intercept[IllegalArgumentException] {
+
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age =
0s")
+ }
+ intercept[IllegalArgumentException] {
+
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age =
-1s")
+ }
+ }
+ }
+
+ "Http2ClientSettings persistent-connection-max-age-jitter" should {
+
+ "default to 0.1" in {
+ val settings = Http2ClientSettings(system)
+ settings.persistentConnectionMaxAgeJitter should ===(0.1)
+ settings.getPersistentConnectionMaxAgeJitter should ===(0.1)
+ }
+
+ "accept a value from config and programmatically" in {
+
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age-jitter
= 0.25")
+ .persistentConnectionMaxAgeJitter should ===(0.25)
+ Http2ClientSettings(system).withPersistentConnectionMaxAgeJitter(0.2)
+ .persistentConnectionMaxAgeJitter should ===(0.2)
+ }
+
+ "reject a value outside of [0, 1)" in {
+ intercept[IllegalArgumentException] {
+
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age-jitter
= -0.1")
+ }
+ intercept[IllegalArgumentException] {
+
Http2ClientSettings("pekko.http.client.http2.persistent-connection-max-age-jitter
= 1")
+ }
+ }
+ }
+}
diff --git
a/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
b/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
index 7378681de..14bb3899f 100644
---
a/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
+++
b/http2-tests/src/test/scala/org/apache/pekko/http/impl/engine/http2/Http2PersistentClientSpec.scala
@@ -102,6 +102,73 @@ abstract class Http2PersistentClientSpec(tls: Boolean)
extends PekkoSpecWithMate
response.attribute(requestIdAttr).get.id shouldBe "request-1"
})
+ "retire an aged connection after in-flight requests
complete".inAssertAllStagesStopped(new TestSetup(tls) {
+ override def clientSettings: ClientConnectionSettings =
+ super.clientSettings.withHttp2Settings(Http2ClientSettings(
+ """
+ pekko.http.client.http2.persistent-connection-max-age = 500ms
+ pekko.http.client.http2.persistent-connection-max-age-jitter = 0
+ pekko.http.client.http2.completion-timeout = 2s
+ """))
+
+ client.responsesIn.request(2)
+ client.sendRequest(HttpRequest(uri =
"/first").addAttribute(requestIdAttr, RequestId("request-1")))
+
+ val first = server.expectRequest()
+ val firstClientPort = first.clientPort
+ killProbe.expectMsgType[UniqueKillSwitch]
+
+ // Let the configured age elapse while the first request is still in
flight. Use a wide margin so this
+ // remains stable on slow CI workers.
+ server.expectNoRequest(1500.millis)
+
+ client.sendRequest(HttpRequest(uri =
"/second").addAttribute(requestIdAttr, RequestId("request-2")))
+ // Closing the request source with a buffered request must not drop that
request during retirement.
+ client.requestsOut.sendComplete()
+ // A request arriving after retirement starts must wait for the current
in-flight request.
+ server.expectNoRequest(300.millis)
+
+ server.sendResponseFor(first, HttpResponse(entity = "first-response"))
+ val firstResponse = client.expectResponse()
+ Unmarshal(firstResponse.entity).to[String].futureValue shouldBe
"first-response"
+ firstResponse.attribute(requestIdAttr).get.id shouldBe "request-1"
+
+ val second = server.expectRequest()
+ second.clientPort should not be firstClientPort
+ server.sendResponseFor(second, HttpResponse(entity = "second-response"))
+
+ val secondResponse = client.expectResponse()
+ Unmarshal(secondResponse.entity).to[String].futureValue shouldBe
"second-response"
+ secondResponse.attribute(requestIdAttr).get.id shouldBe "request-2"
+ })
+
+ "retire an idle aged connection and reconnect lazily for the next
request".inAssertAllStagesStopped(
+ new TestSetup(tls) {
+ override def clientSettings: ClientConnectionSettings =
+ super.clientSettings.withHttp2Settings(Http2ClientSettings(
+ """
+ pekko.http.client.http2.persistent-connection-max-age = 500ms
+ pekko.http.client.http2.persistent-connection-max-age-jitter = 0
+ """))
+
+ client.responsesIn.request(2)
+ client.sendRequest(HttpRequest(uri = "/first"))
+ val first = server.expectRequest()
+ val firstClientPort = first.clientPort
+ server.sendResponseFor(first, HttpResponse())
+ client.expectResponse()
+
+ // Allow the idle connection to reach its configured age with a
generous scheduling margin.
+ server.expectNoRequest(1500.millis)
+
+ client.sendRequest(HttpRequest(uri = "/second"))
+ val second = server.expectRequest()
+ second.clientPort should not be firstClientPort
+ server.sendResponseFor(second, HttpResponse())
+ client.expectResponse()
+ client.requestsOut.sendComplete()
+ })
+
def reconnectionTests(withBackoff: Boolean): Unit = {
val changeSettings: Http2ClientSettings => Http2ClientSettings =
if (withBackoff) s =>
s.withBaseConnectionBackoff(300.millis).withMaxConnectionBackoff(800.millis)
@@ -396,6 +463,7 @@ abstract class Http2PersistentClientSpec(tls: Boolean)
extends PekkoSpecWithMate
}
def expectRequest(): ServerRequest =
requestProbe.expectMsgType[ServerRequest]
+ def expectNoRequest(duration: FiniteDuration): Unit =
requestProbe.expectNoMessage(duration)
def sendResponseFor(request: ServerRequest, response: HttpResponse):
Unit =
request.sendResponse(response)
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]