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]