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 5b9a02ee7 fix: cancel open request timeouts when the server connection
stage stops (#1284)
5b9a02ee7 is described below
commit 5b9a02ee7b30a391aa941f7563e3cc389530a0a6
Author: PJ Fanning <[email protected]>
AuthorDate: Fri Sep 11 09:43:10 2026 +0100
fix: cancel open request timeouts when the server connection stage stops
(#1284)
Motivation:
`RequestTimeoutSupport` schedules a request's timeout through
`materializer.scheduleOnce` and keeps the `TimeoutAccess` in
`openTimeouts` until the response for that request comes back, which is
when the timeout is cancelled.
The stage had no `postStop`, so a connection that goes away with
requests still open - the client disconnecting, the connection being
aborted, or the server being torn down - left one scheduled task per
open request behind. Those tasks keep the request, the timeout handler
and the stage's async callback reachable for as long as the configured
`request-timeout` (20 seconds by default), and firing them achieves
nothing: they only invoke a callback on a stage that is already gone.
Modification:
Cancel everything still in `openTimeouts` in `postStop`, using the same
best-effort `clear()` the regular response path uses.
Result:
Connections that end before their requests are answered no longer leave
scheduled timeouts behind.
Tests:
- sbt "http-core/testOnly
org.apache.pekko.http.impl.engine.server.HttpServerSpec
org.apache.pekko.http.impl.engine.server.HttpServerWithExplicitSchedulerSpec" -
pass (66 tests). Reverting only the `postStop` addition makes the new test
fail, as the timeout still fires after the stage stopped.
- sbt "http-core/mimaReportBinaryIssues" - pass.
- scalafmt --mode diff-ref=upstream/main --test - pass.
References:
None - found while auditing `src/main` for resource leaks
---
.../impl/engine/server/HttpServerBluePrint.scala | 7 ++++
.../HttpServerWithExplicitSchedulerSpec.scala | 38 ++++++++++++++++++++++
2 files changed, 45 insertions(+)
diff --git
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/server/HttpServerBluePrint.scala
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/server/HttpServerBluePrint.scala
index e08901e36..2948ceaed 100644
---
a/http-core/src/main/scala/org/apache/pekko/http/impl/engine/server/HttpServerBluePrint.scala
+++
b/http-core/src/main/scala/org/apache/pekko/http/impl/engine/server/HttpServerBluePrint.scala
@@ -349,6 +349,13 @@ private[http] object HttpServerBluePrint {
def onPull(): Unit = pull(responseIn)
override def onDownstreamFinish(cause: Throwable): Unit =
cancel(responseIn)
})
+
+ override def postStop(): Unit = {
+ // the timeouts are scheduled on the materializer, so they outlive
this stage unless they are cancelled here.
+ // Whatever is still open now belongs to a connection that is gone, so
firing them cannot do anything useful.
+ openTimeouts.foreach(_.clear())
+ openTimeouts = immutable.Queue.empty
+ }
}
}
diff --git
a/http-core/src/test/scala/org/apache/pekko/http/impl/engine/server/HttpServerWithExplicitSchedulerSpec.scala
b/http-core/src/test/scala/org/apache/pekko/http/impl/engine/server/HttpServerWithExplicitSchedulerSpec.scala
index c575c7ac0..54fea9b17 100644
---
a/http-core/src/test/scala/org/apache/pekko/http/impl/engine/server/HttpServerWithExplicitSchedulerSpec.scala
+++
b/http-core/src/test/scala/org/apache/pekko/http/impl/engine/server/HttpServerWithExplicitSchedulerSpec.scala
@@ -19,10 +19,13 @@ import pekko.http.impl.util._
import pekko.http.scaladsl.model._
import pekko.http.scaladsl.model.headers.{ `Timeout-Access`, Connection }
import pekko.stream.Materializer
+import pekko.stream.scaladsl.{ BidiFlow, Flow, Keep }
+import pekko.stream.testkit.scaladsl.{ TestSink, TestSource }
import pekko.stream.testkit.Utils.assertAllStagesStopped
import pekko.testkit.ExplicitlyTriggeredScheduler
import org.scalatest.Inside
+import scala.concurrent.Promise
import scala.concurrent.duration._
/** Tests similar to HttpServerSpec that need ExplicitlyTriggeredScheduler */
@@ -171,6 +174,41 @@ class HttpServerWithExplicitSchedulerSpec extends
PekkoSpecWithMaterializer(
netOut.expectComplete()
netIn.sendComplete()
})
+
+ // the timeouts are scheduled on the materializer, so they are not tied
to the lifetime of the stage that
+ // created them and have to be cancelled when a connection goes away
with a request still open
+ "are cancelled when the stage that scheduled them stops" in
assertAllStagesStopped {
+ val scheduler =
system.scheduler.asInstanceOf[ExplicitlyTriggeredScheduler]
+ val timeoutHandlerCalled = Promise[Unit]()
+
+ val timeoutSupport =
+ BidiFlow.fromGraph(new
HttpServerBluePrint.RequestTimeoutSupport(400.millis, system.log))
+ val applicationSide =
+ Flow.fromSinkAndSourceMat(TestSink[HttpRequest](),
TestSource[HttpResponse]())(Keep.both)
+ val ((requestsIn, (applicationRequests, applicationResponses)),
responsesOut) =
+ TestSource[HttpRequest]()
+
.viaMat(timeoutSupport.joinMat(applicationSide)(Keep.right))(Keep.both)
+ .toMat(TestSink[HttpResponse]())(Keep.both)
+ .run()
+
+ responsesOut.request(1)
+ requestsIn.sendNext(HttpRequest())
+ applicationRequests.requestNext().header[`Timeout-Access`].foreach(
+ _.timeoutAccess.updateHandler { (_: HttpRequest) =>
+ timeoutHandlerCalled.trySuccess(())
+ HttpResponse(StatusCodes.InternalServerError)
+ })
+
+ // the connection goes away while the application is still working on
the response
+ requestsIn.sendComplete()
+ applicationRequests.expectComplete()
+ applicationResponses.sendComplete()
+ responsesOut.expectComplete()
+
+ // nothing is left to run the timeout of a request that can no longer
be answered
+ scheduler.timePasses(500.millis)
+ timeoutHandlerCalled.isCompleted shouldBe false
+ }
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]