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-management.git
The following commit(s) were added to refs/heads/main by this push:
new 2fd1f0f9 No Request-Level Timeout on K8s API singleRequest (#915)
2fd1f0f9 is described below
commit 2fd1f0f98085ac5354ed93458f5ae7a53edf4fcc
Author: PJ Fanning <[email protected]>
AuthorDate: Mon Aug 10 09:50:33 2026 +0100
No Request-Level Timeout on K8s API singleRequest (#915)
---
.../kubernetes/KubernetesApiServiceDiscovery.scala | 21 ++++++++++++++++++++-
1 file changed, 20 insertions(+), 1 deletion(-)
diff --git
a/discovery-kubernetes-api/src/main/scala/org/apache/pekko/discovery/kubernetes/KubernetesApiServiceDiscovery.scala
b/discovery-kubernetes-api/src/main/scala/org/apache/pekko/discovery/kubernetes/KubernetesApiServiceDiscovery.scala
index 943af7b5..36dcd766 100644
---
a/discovery-kubernetes-api/src/main/scala/org/apache/pekko/discovery/kubernetes/KubernetesApiServiceDiscovery.scala
+++
b/discovery-kubernetes-api/src/main/scala/org/apache/pekko/discovery/kubernetes/KubernetesApiServiceDiscovery.scala
@@ -15,12 +15,14 @@ package org.apache.pekko.discovery.kubernetes
import java.net.InetAddress
import java.nio.charset.StandardCharsets
+import java.util.concurrent.TimeoutException
import java.nio.file.{ Files, Paths }
import scala.collection.immutable
import scala.collection.immutable.Seq
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
+import scala.concurrent.Promise
import scala.concurrent.duration.FiniteDuration
import scala.util.Try
import scala.util.control.{ NoStackTrace, NonFatal }
@@ -154,7 +156,24 @@ class KubernetesApiServiceDiscovery(settings: Settings)(
)
}
- response <- http.singleRequest(request,
setup.clientHttpsConnectionContext).map(decodeResponse)
+ response <- {
+ val rawResponse = http.singleRequest(request,
setup.clientHttpsConnectionContext)
+ val promise = Promise[HttpResponse]()
+ val timeoutCancellable = system.scheduler.scheduleOnce(resolveTimeout)
{
+ promise.tryFailure(new TimeoutException(s"Kubernetes API request
timed out after $resolveTimeout"))
+ }
+ rawResponse.onComplete {
+ case scala.util.Success(resp) =>
+ timeoutCancellable.cancel()
+ if (!promise.trySuccess(resp)) {
+ resp.discardEntityBytes()
+ }
+ case scala.util.Failure(ex) =>
+ timeoutCancellable.cancel()
+ promise.tryFailure(ex)
+ }(system.dispatcher)
+ promise.future.map(decodeResponse)
+ }
entity <- response.entity.toStrict(resolveTimeout)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]