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]

Reply via email to