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 8665f626 improve Consul support (#906)
8665f626 is described below

commit 8665f626c88f9f620b59e63f53ed616533f92e04
Author: PJ Fanning <[email protected]>
AuthorDate: Tue Aug 4 22:16:30 2026 +0100

    improve Consul support (#906)
    
    * improve Consul support
    
    * Update ConsulSettings.scala
    
    * tls client support
    
    * review comments
---
 discovery-consul/src/main/resources/reference.conf | 22 ++++++++
 .../discovery/consul/ConsulServiceDiscovery.scala  | 64 +++++++++++++++++++---
 .../pekko/discovery/consul/ConsulSettings.scala    | 54 ++++++++++++++++++
 3 files changed, 133 insertions(+), 7 deletions(-)

diff --git a/discovery-consul/src/main/resources/reference.conf 
b/discovery-consul/src/main/resources/reference.conf
index 733c0d73..fc9f9491 100644
--- a/discovery-consul/src/main/resources/reference.conf
+++ b/discovery-consul/src/main/resources/reference.conf
@@ -23,5 +23,27 @@ pekko.discovery {
     # Prefix for tag containing port number where Pekko management is set up 
so that
     # the seed nodes can be found, an example value for the tag would be 
`pekko-management-port:19999`
     application-pekko-management-port-tag-prefix = "pekko-management-port:"
+
+    # Connection timeout for the Consul HTTP client.
+    connect-timeout = 10s
+
+    # Read timeout for the Consul HTTP client.
+    read-timeout = 10s
+
+    # Write timeout for the Consul HTTP client.
+    write-timeout = 10s
+
+    # Maximum number of concurrent requests when looking up multiple service 
IDs.
+    lookup-parallelism = 8
+
+    # ACL token for Consul API authentication. Empty means no authentication.
+    consul-token = ""
+
+    # Whether to use HTTPS when connecting to Consul.
+    tls-enabled = false
+
+    # Path to a PEM-encoded CA certificate file for TLS verification.
+    # Only used when tls-enabled = true. If empty, the default JVM trust store 
is used.
+    ca-path = ""
   }
 }
diff --git 
a/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulServiceDiscovery.scala
 
b/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulServiceDiscovery.scala
index 0cb26734..ade1ed7b 100644
--- 
a/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulServiceDiscovery.scala
+++ 
b/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulServiceDiscovery.scala
@@ -15,20 +15,25 @@ package org.apache.pekko.discovery.consul
 
 import com.google.common.net.HostAndPort
 import org.apache.pekko
-import pekko.actor.ActorSystem
+import pekko.actor.{ ActorSystem, CoordinatedShutdown }
 import pekko.annotation.ApiMayChange
 import pekko.discovery.ServiceDiscovery.{ Resolved, ResolvedTarget }
 import pekko.discovery.consul.ConsulServiceDiscovery._
 import pekko.discovery.{ Lookup, ServiceDiscovery }
+import pekko.dispatch.Dispatchers.DefaultBlockingDispatcherId
 import org.kiwiproject.consul.Consul
 import org.kiwiproject.consul.async.ConsulResponseCallback
 import org.kiwiproject.consul.model.ConsulResponse
 import org.kiwiproject.consul.model.catalog.CatalogService
 import org.kiwiproject.consul.option.Options
 
+import java.io.FileInputStream
 import java.net.InetAddress
+import java.security.KeyStore
+import java.security.cert.CertificateFactory
 import java.util
 import java.util.concurrent.TimeoutException
+import javax.net.ssl.{ SSLContext, TrustManagerFactory }
 import scala.collection.immutable.Seq
 import scala.concurrent.duration.FiniteDuration
 import scala.concurrent.{ ExecutionContext, Future, Promise }
@@ -39,8 +44,39 @@ import scala.util.Try
 class ConsulServiceDiscovery(system: ActorSystem) extends ServiceDiscovery {
 
   private val settings = ConsulSettings.get(system)
-  private val consul =
-    
Consul.builder().withHostAndPort(HostAndPort.fromParts(settings.consulHost, 
settings.consulPort)).build()
+  private val consul = {
+    val builder = Consul
+      .builder()
+      .withHostAndPort(HostAndPort.fromParts(settings.consulHost, 
settings.consulPort))
+      .withConnectTimeoutMillis(settings.connectTimeout.toMillis)
+      .withReadTimeoutMillis(settings.readTimeout.toMillis)
+      .withWriteTimeoutMillis(settings.writeTimeout.toMillis)
+    settings.consulToken.foreach(builder.withTokenAuth)
+    if (settings.tlsEnabled) {
+      builder.withHttps(true)
+      settings.caPath.foreach { caPath =>
+        val cf = CertificateFactory.getInstance("X.509")
+        val caCert = cf.generateCertificate(new FileInputStream(caPath))
+        val ks = KeyStore.getInstance(KeyStore.getDefaultType)
+        ks.load(null)
+        ks.setCertificateEntry("ca", caCert)
+        val tmf = 
TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm)
+        tmf.init(ks)
+        val sslContext = SSLContext.getInstance("TLS")
+        sslContext.init(null, tmf.getTrustManagers, null)
+        builder.withSslContext(sslContext)
+      }
+    }
+    builder.build()
+  }
+  private val blockingEc: ExecutionContext = 
system.dispatchers.lookup(DefaultBlockingDispatcherId)
+
+  CoordinatedShutdown(system).addTask(CoordinatedShutdown.PhaseServiceUnbind, 
"consul-close") { () =>
+    Future {
+      consul.destroy()
+      pekko.Done
+    }(system.dispatcher)
+  }
 
   override def lookup(lookup: Lookup, resolveTimeout: FiniteDuration): 
Future[Resolved] = {
     implicit val ec: ExecutionContext = system.dispatcher
@@ -60,14 +96,16 @@ class ConsulServiceDiscovery(system: ActorSystem) extends 
ServiceDiscovery {
   private def lookupInConsul(name: String)(implicit executionContext: 
ExecutionContext): Future[Resolved] = {
     val consulResult = for {
       servicesWithTags <- getServicesWithTags
+      nameTag = settings.applicationNameTagPrefix + name
       serviceIds = servicesWithTags.getResponse
         .entrySet()
         .asScala
-        .filter(e => e.getValue.contains(settings.applicationNameTagPrefix + 
name))
+        .filter(e => e.getValue.contains(nameTag))
         .map(_.getKey)
-      catalogServices <- Future.sequence(serviceIds.map(id => 
getService(id).map(_.getResponse.asScala.toList)))
-      resolvedTargets = catalogServices.flatten.toSeq.map(catalogService =>
-        extractResolvedTargetFromCatalogService(catalogService))
+      catalogServices <- boundedTraverse(serviceIds.toSeq)(id => 
getService(id).map(_.getResponse.asScala.toList))
+      resolvedTargets <- Future.traverse(catalogServices.flatten.toSeq) { 
catalogService =>
+        
Future(extractResolvedTargetFromCatalogService(catalogService))(blockingEc)
+      }
     } yield resolvedTargets
     consulResult.map(targets => Resolved(name, 
scala.collection.immutable.Seq(targets: _*)))
   }
@@ -86,6 +124,18 @@ class ConsulServiceDiscovery(system: ActorSystem) extends 
ServiceDiscovery {
       address = Try(InetAddress.getByName(address)).toOption)
   }
 
+  private def boundedTraverse[A, B](items: Seq[A])(f: A => Future[B])(
+      implicit ec: ExecutionContext): Future[Seq[B]] = {
+    def loop(remaining: Seq[A], acc: Seq[B]): Future[Seq[B]] = {
+      if (remaining.isEmpty) Future.successful(acc.reverse)
+      else {
+        val (batch, rest) = remaining.splitAt(settings.parallelism)
+        Future.traverse(batch)(f).flatMap(results => loop(rest, 
results.reverse ++ acc))
+      }
+    }
+    loop(items, Seq.empty)
+  }
+
   private def getServicesWithTags: Future[ConsulResponse[util.Map[String, 
util.List[String]]]] = {
     ((callback: ConsulResponseCallback[util.Map[String, util.List[String]]]) =>
           consul.catalogClient().getServices(callback)).asFuture
diff --git 
a/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulSettings.scala
 
b/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulSettings.scala
index 62404878..baabcb8a 100644
--- 
a/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulSettings.scala
+++ 
b/discovery-consul/src/main/scala/org/apache/pekko/discovery/consul/ConsulSettings.scala
@@ -17,6 +17,8 @@ import org.apache.pekko
 import pekko.actor.ClassicActorSystemProvider
 import pekko.actor.{ ActorSystem, ExtendedActorSystem, Extension, ExtensionId, 
ExtensionIdProvider }
 import pekko.annotation.ApiMayChange
+import scala.concurrent.duration._
+import scala.jdk.DurationConverters._
 
 @ApiMayChange
 final class ConsulSettings(system: ExtendedActorSystem) extends Extension {
@@ -29,6 +31,58 @@ final class ConsulSettings(system: ExtendedActorSystem) 
extends Extension {
   val applicationNameTagPrefix: String = 
consulConfig.getString("application-name-tag-prefix")
   val applicationPekkoManagementPortTagPrefix: String =
     consulConfig.getString("application-pekko-management-port-tag-prefix")
+
+  /**
+   * Connection timeout for the Consul HTTP client.
+   * @since 2.0.0
+   */
+  val connectTimeout: FiniteDuration =
+    consulConfig.getDuration("connect-timeout").toScala
+
+  /**
+   * Read timeout for the Consul HTTP client.
+   * @since 2.0.0
+   */
+  val readTimeout: FiniteDuration =
+    consulConfig.getDuration("read-timeout").toScala
+
+  /**
+   * Write timeout for the Consul HTTP client.
+   * @since 2.0.0
+   */
+  val writeTimeout: FiniteDuration =
+    consulConfig.getDuration("write-timeout").toScala
+
+  /**
+   * Maximum number of concurrent requests when looking up multiple service 
IDs.
+   * @since 2.0.0
+   */
+  val parallelism: Int = consulConfig.getInt("lookup-parallelism")
+
+  /**
+   * ACL token for Consul API authentication. Empty means no authentication.
+   * @since 2.0.0
+   */
+  val consulToken: Option[String] = consulConfig.getString("consul-token") 
match {
+    case ""  => None
+    case tok => Some(tok)
+  }
+
+  /**
+   * Whether to use HTTPS when connecting to Consul.
+   * @since 2.0.0
+   */
+  val tlsEnabled: Boolean = consulConfig.getBoolean("tls-enabled")
+
+  /**
+   * Path to a PEM-encoded CA certificate file for TLS verification.
+   * Only used when tls-enabled = true. If empty, the default JVM trust store 
is used.
+   * @since 2.0.0
+   */
+  val caPath: Option[String] = consulConfig.getString("ca-path") match {
+    case ""   => None
+    case path => Some(path)
+  }
 }
 
 @ApiMayChange


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to