From 2f7f4f0c0b00475d867ff19a8c94acf90c82647a Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Sun, 14 Feb 2021 20:56:07 +0000 Subject: [PATCH 01/16] Create KafkaCredentialStore --- .../kafka/security/ClientCertificate.scala | 27 ++++ .../fs2/kafka/security/ClientPrivateKey.scala | 47 ++++++ .../kafka/security/KafkaCredentialStore.scala | 141 ++++++++++++++++++ .../fs2/kafka/security/KeyStoreFile.scala | 28 ++++ .../fs2/kafka/security/KeyStorePassword.scala | 24 +++ .../kafka/security/ServiceCertificate.scala | 27 ++++ .../fs2/kafka/security/TrustStoreFile.scala | 28 ++++ .../kafka/security/TrustStorePassword.scala | 24 +++ .../security/internal/CertificateOps.scala | 19 +++ .../fs2/kafka/security/internal/FileOps.scala | 14 ++ .../scala/fs2/kafka/security/package.scala | 35 +++++ 11 files changed, 414 insertions(+) create mode 100644 modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala create mode 100644 modules/core/src/main/scala/fs2/kafka/security/package.scala diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala new file mode 100644 index 000000000..ae4a61172 --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala @@ -0,0 +1,27 @@ +package fs2.kafka.security + +import java.security.cert.{Certificate, CertificateException} + +sealed abstract class ClientCertificate { + def value: Certificate +} + +object ClientCertificate { + def fromString(clientCertificate: String): Either[CertificateException, ClientCertificate] = { + internal.CertificateOps.loadFromString(clientCertificate).map { certificate => + new ClientCertificate { + override final val value: Certificate = + certificate + + override final def toString: String = + s"ClientCertificate(${clientCertificate.valueShortHash})" + } + } + } + + def fromCertificate(certificate: Certificate): ClientCertificate = { + new ClientCertificate { + override def value: Certificate = certificate + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala new file mode 100644 index 000000000..e1663d172 --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala @@ -0,0 +1,47 @@ +package fs2.kafka.security + +import cats.syntax.all._ + +import java.nio.charset.StandardCharsets +import java.security.{GeneralSecurityException, KeyFactory, PrivateKey} +import java.security.spec.PKCS8EncodedKeySpec +import java.util.Base64 + +sealed abstract class ClientPrivateKey { + def value: PrivateKey +} + +object ClientPrivateKey { + def fromString(clientPrivateKey: String): Either[GeneralSecurityException, ClientPrivateKey] = { + Either.catchOnly[GeneralSecurityException] { + new ClientPrivateKey { + override final val value: PrivateKey = + KeyFactory + .getInstance("RSA") + .generatePrivate { + new PKCS8EncodedKeySpec( + Base64.getDecoder.decode { + clientPrivateKey + .replace("-----BEGIN PRIVATE KEY-----", "") + .replace("-----END PRIVATE KEY-----", "") + .filterNot(_.isWhitespace) + .getBytes(StandardCharsets.UTF_8) + } + ) + } + + override final def toString: String = + s"ClientPrivateKey(${clientPrivateKey.valueShortHash})" + } + } + } + + def fromPrivateKey(privateKey: PrivateKey): ClientPrivateKey = { + new ClientPrivateKey { + override def value: PrivateKey = privateKey + + override final def toString: String = + s"ClientPrivateKey(${privateKey.toString.valueShortHash})" + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala new file mode 100644 index 000000000..931e4649f --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala @@ -0,0 +1,141 @@ +package fs2.kafka.security + +import cats.effect._ +import cats.syntax.all._ + +import java.nio.file.{Files, Path} +import java.security.KeyStore + +sealed abstract class KafkaCredentialStore { + def keyStoreFile: KeyStoreFile + + def keyStorePassword: KeyStorePassword + + def trustStoreFile: TrustStoreFile + + def trustStorePassword: TrustStorePassword + + def properties: Map[String, String] +} + +object KafkaCredentialStore { + final def apply[F[_]: ContextShift]( + clientPrivateKey: ClientPrivateKey, + clientCertificate: ClientCertificate, + serviceCertificate: ServiceCertificate, + blocker: Blocker + )(implicit F: Sync[F]): F[KafkaCredentialStore] = + for { + setupDetails <- KafkaCredentialStore.createTemporary[F] + _ <- setupKeyStore( + clientPrivateKey = clientPrivateKey, + clientCertificate = clientCertificate, + keyStoreFile = setupDetails.keyStoreFile, + keyStorePassword = setupDetails.keyStorePassword, + blocker = blocker + ) + _ <- setupTrustStore( + serviceCertificate = serviceCertificate, + trustStoreFile = setupDetails.trustStoreFile, + trustStorePassword = setupDetails.trustStorePassword, + blocker = blocker + ) + } yield setupDetails + + private final def setupStore[F[_]: ContextShift]( + storeType: String, + storePath: Path, + storePasswordChars: Array[Char], + setupStore: KeyStore => Unit, + blocker: Blocker + )(implicit F: Sync[F]): F[Unit] = + blocker.delay { + val keyStore = KeyStore.getInstance(storeType) + keyStore.load(null, storePasswordChars) + setupStore(keyStore) + + val outputStream = Files.newOutputStream(storePath) + + try { + keyStore.store(outputStream, storePasswordChars) + } finally { + outputStream.close() + } + } + + private final def setupKeyStore[F[_]: ContextShift]( + clientPrivateKey: ClientPrivateKey, + clientCertificate: ClientCertificate, + keyStoreFile: KeyStoreFile, + keyStorePassword: KeyStorePassword, + blocker: Blocker + )(implicit F: Sync[F]): F[Unit] = { + val keyStorePasswordChars = + keyStorePassword.value.toCharArray + + setupStore( + storeType = "PKCS12", + storePath = keyStoreFile.path, + storePasswordChars = keyStorePasswordChars, + setupStore = _.setEntry( + "service_key", + new KeyStore.PrivateKeyEntry( + clientPrivateKey.value, + Array(clientCertificate.value) + ), + new KeyStore.PasswordProtection(keyStorePasswordChars) + ), + blocker = blocker + ) + } + + private final def setupTrustStore[F[_]: ContextShift]( + serviceCertificate: ServiceCertificate, + trustStoreFile: TrustStoreFile, + trustStorePassword: TrustStorePassword, + blocker: Blocker + )(implicit F: Sync[F]): F[Unit] = + setupStore( + storeType = "JKS", + storePath = trustStoreFile.path, + storePasswordChars = trustStorePassword.value.toCharArray, + setupStore = _.setCertificateEntry("CA", serviceCertificate.value), + blocker = blocker + ) + + private final def createTemporary[F[_]](implicit F: Sync[F]): F[KafkaCredentialStore] = + for { + _keyStoreFile <- KeyStoreFile.createTemporary[F] + _keyStorePassword <- KeyStorePassword.createTemporary[F] + _trustStoreFile <- TrustStoreFile.createTemporary[F] + _trustStorePassword <- TrustStorePassword.createTemporary[F] + } yield { + new KafkaCredentialStore { + override final val keyStoreFile: KeyStoreFile = + _keyStoreFile + + override final val keyStorePassword: KeyStorePassword = + _keyStorePassword + + override final val trustStoreFile: TrustStoreFile = + _trustStoreFile + + override final val trustStorePassword: TrustStorePassword = + _trustStorePassword + + override final val properties: Map[String, String] = + Map( + "security.protocol" -> "SSL", + "ssl.truststore.location" -> trustStoreFile.pathAsString, + "ssl.truststore.password" -> trustStorePassword.value, + "ssl.keystore.type" -> "PKCS12", + "ssl.keystore.location" -> keyStoreFile.pathAsString, + "ssl.keystore.password" -> keyStorePassword.value, + "ssl.key.password" -> keyStorePassword.value + ) + + override final def toString: String = + s"KafkaCredentialStore($keyStoreFile, $keyStorePassword, $trustStoreFile, $trustStorePassword)" + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala new file mode 100644 index 000000000..c34b8b66e --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala @@ -0,0 +1,28 @@ +package fs2.kafka.security + +import cats.effect.Sync +import cats.syntax.all._ + +import java.nio.file.Path + +sealed abstract class KeyStoreFile { + def path: Path + + def pathAsString: String +} + +private[security] object KeyStoreFile { + final def createTemporary[F[_]](implicit F: Sync[F]): F[KeyStoreFile] = + internal.FileOps.createTemp("client.keystore-", ".p12").map { _path => + new KeyStoreFile { + override final val path: Path = + _path + + override final def pathAsString: String = + path.toString + + override final def toString: String = + s"KeyStoreFile($pathAsString)" + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala new file mode 100644 index 000000000..0927ec7f1 --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala @@ -0,0 +1,24 @@ +package fs2.kafka.security + +import cats.effect.Sync + +import java.util.UUID + +sealed abstract class KeyStorePassword { + def value: String +} + +private[security] object KeyStorePassword { + def createTemporary[F[_]](implicit F: Sync[F]): F[KeyStorePassword] = + F.delay { + val _value = UUID.randomUUID().toString + + new KeyStorePassword { + override final val value: String = + _value + + override final def toString: String = + s"KeyStorePassword(${value.valueShortHash})" + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala new file mode 100644 index 000000000..ad4f7201f --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala @@ -0,0 +1,27 @@ +package fs2.kafka.security + +import java.security.cert.{Certificate, CertificateException} + +sealed abstract class ServiceCertificate { + def value: Certificate +} + +object ServiceCertificate { + def apply(serviceCertificate: String): Either[CertificateException, ServiceCertificate] = { + internal.CertificateOps.loadFromString(serviceCertificate).map { certificate => + new ServiceCertificate { + override final val value: Certificate = + certificate + + override final def toString: String = + s"ServiceCertificate(${serviceCertificate.valueShortHash})" + } + } + } + + def fromCertificate(certificate: Certificate): ServiceCertificate = { + new ServiceCertificate { + override def value: Certificate = certificate + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala new file mode 100644 index 000000000..9758d9ff5 --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala @@ -0,0 +1,28 @@ +package fs2.kafka.security + +import cats.effect.Sync +import cats.syntax.all._ + +import java.nio.file.Path + +sealed abstract class TrustStoreFile { + def path: Path + + def pathAsString: String +} + +private[security] object TrustStoreFile { + def createTemporary[F[_]](implicit F: Sync[F]): F[TrustStoreFile] = + internal.FileOps.createTemp[F]("client.truststore-", ".jks").map { _path => + new TrustStoreFile { + override final val path: Path = + _path + + override final def pathAsString: String = + path.toString + + override final def toString: String = + s"TrustStoreFile($pathAsString)" + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala new file mode 100644 index 000000000..2524a0d3e --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala @@ -0,0 +1,24 @@ +package fs2.kafka.security + +import cats.effect.Sync + +import java.util.UUID + +sealed abstract class TrustStorePassword { + def value: String +} + +private[security] object TrustStorePassword { + def createTemporary[F[_]](implicit F: Sync[F]): F[TrustStorePassword] = + F.delay { + val _value = UUID.randomUUID().toString + + new TrustStorePassword { + override final val value: String = + _value + + override final def toString: String = + s"TrustStorePassword(${value.valueShortHash})" + } + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala b/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala new file mode 100644 index 000000000..6dd8e6e0a --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala @@ -0,0 +1,19 @@ +package fs2.kafka.security.internal + +import cats.syntax.all._ + +import java.io.ByteArrayInputStream +import java.nio.charset.StandardCharsets +import java.security.cert.{Certificate, CertificateException, CertificateFactory} + +private[security] object CertificateOps { + def loadFromString(certificate: String): Either[CertificateException, Certificate] = + loadFromBytes(certificate.getBytes(StandardCharsets.UTF_8)) + + def loadFromBytes(certificateBytes: Array[Byte]): Either[CertificateException, Certificate] = + Either.catchOnly[CertificateException] { + CertificateFactory + .getInstance("X.509") + .generateCertificate(new ByteArrayInputStream(certificateBytes)) + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala b/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala new file mode 100644 index 000000000..a8580cd78 --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala @@ -0,0 +1,14 @@ +package fs2.kafka.security.internal + +import cats.effect.Sync + +import java.nio.file.{Files, Path} + +private[security] object FileOps { + def createTemp[F[_]](prefix: String, suffix: String)(implicit F: Sync[F]): F[Path] = F.delay { + val path = Files.createTempFile(prefix, suffix) + path.toFile.deleteOnExit() + Files.delete(path) + path + } +} diff --git a/modules/core/src/main/scala/fs2/kafka/security/package.scala b/modules/core/src/main/scala/fs2/kafka/security/package.scala new file mode 100644 index 000000000..efadf05a8 --- /dev/null +++ b/modules/core/src/main/scala/fs2/kafka/security/package.scala @@ -0,0 +1,35 @@ +package fs2.kafka + +import java.nio.charset.StandardCharsets +import java.security.MessageDigest +import scala.annotation.tailrec + +package object security { + private[this] final val hexChars: Array[Char] = + Array('0', '1', '2', '3', '4', '5', '6', '7', '8', '9', 'a', 'b', 'c', 'd', 'e', 'f') + + implicit class StringOps(val str: String) extends AnyVal { + private[this] final def hex(in: Array[Byte]): Array[Char] = { + val length = in.length + + @tailrec def encode(out: Array[Char], i: Int, j: Int): Array[Char] = { + if (i < length) { + out(j) = hexChars((0xf0 & in(i)) >>> 4) + out(j + 1) = hexChars(0x0f & in(i)) + encode(out, i + 1, j + 2) + } else out + } + + encode(new Array(length << 1), 0, 0) + } + + private[this] final def sha1(bytes: Array[Byte]): Array[Byte] = + MessageDigest.getInstance("SHA-1").digest(bytes) + + def sha1Hex: String = + new String(hex(sha1(str.getBytes(StandardCharsets.UTF_8)))) + + final def valueShortHash: String = + sha1Hex.take(7) + } +} From e206794aa2562f38e1fbda46d92bdc6d22bbeef8 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Sun, 14 Feb 2021 21:44:35 +0000 Subject: [PATCH 02/16] Add certificates.md --- docs/src/main/mdoc/certificates.md | 29 +++++++++++++++++++++++++++++ 1 file changed, 29 insertions(+) create mode 100644 docs/src/main/mdoc/certificates.md diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md new file mode 100644 index 000000000..d3fda556d --- /dev/null +++ b/docs/src/main/mdoc/certificates.md @@ -0,0 +1,29 @@ +## Security: certificates, trust stores, and passwords + +The `KafkaCredentialStore` can be used to create the necessary. + +The parameters passed in are string representations of the client private key, client certificate +and service certificate. the `properties` field in `KafkaCredentialStore` can then be applied to +any of the `*Settings` classes by using the `withProperties(kafkaCredentialStore.properties)`. + +```scala mdoc +import cats.effect._ +import cats.syntax.all._ +import fs2.kafka.security._ + +def loadKafkaSetup[F[_]: Async: ContextShift]( + clientPrivateKey: String, + clientCertificate: String, + serviceCertificate: String, +): Resource[F, KafkaCredentialStore] = + Blocker[F].evalMap { blocker => + ( + ClientPrivateKey(clientPrivateKey).liftTo[F], + ClientCertificate(clientCertificate).liftTo[F], + ServiceCertificate(serviceCertificate).liftTo[F], + ).tupled.flatMap { + case (clientPrivateKey, clientCertificate, serviceCertificate) => + KafkaCredentialStore[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) + } + } +``` \ No newline at end of file From 39e0330d5e3affe050a4343af952ae81b03a9319 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Sun, 14 Feb 2021 21:48:39 +0000 Subject: [PATCH 03/16] certificates.md: add title and id --- docs/src/main/mdoc/certificates.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index d3fda556d..e3dde5f44 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -1,3 +1,8 @@ +--- +id: security +title: Security & Certificates +--- + ## Security: certificates, trust stores, and passwords The `KafkaCredentialStore` can be used to create the necessary. From 8894492a11adac0b19635ec116d31139877be220 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Sun, 14 Feb 2021 22:04:50 +0000 Subject: [PATCH 04/16] certificates.md: fix docs using fromString methods --- docs/src/main/mdoc/certificates.md | 6 +++--- .../main/scala/fs2/kafka/security/ServiceCertificate.scala | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index e3dde5f44..05f0ab5bc 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -23,9 +23,9 @@ def loadKafkaSetup[F[_]: Async: ContextShift]( ): Resource[F, KafkaCredentialStore] = Blocker[F].evalMap { blocker => ( - ClientPrivateKey(clientPrivateKey).liftTo[F], - ClientCertificate(clientCertificate).liftTo[F], - ServiceCertificate(serviceCertificate).liftTo[F], + ClientPrivateKey.fromString(clientPrivateKey).liftTo[F], + ClientCertificate.fromString(clientCertificate).liftTo[F], + ServiceCertificate.fromString(serviceCertificate).liftTo[F], ).tupled.flatMap { case (clientPrivateKey, clientCertificate, serviceCertificate) => KafkaCredentialStore[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) diff --git a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala index ad4f7201f..5e0802aad 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala @@ -7,7 +7,7 @@ sealed abstract class ServiceCertificate { } object ServiceCertificate { - def apply(serviceCertificate: String): Either[CertificateException, ServiceCertificate] = { + def fromString(serviceCertificate: String): Either[CertificateException, ServiceCertificate] = { internal.CertificateOps.loadFromString(serviceCertificate).map { certificate => new ServiceCertificate { override final val value: Certificate = From ffe9c4b26ea1469c1c6df5856591fa0c2a685b5a Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Sun, 14 Feb 2021 22:17:55 +0000 Subject: [PATCH 05/16] Add headers --- .../main/scala/fs2/kafka/security/ClientCertificate.scala | 6 ++++++ .../main/scala/fs2/kafka/security/ClientPrivateKey.scala | 6 ++++++ .../scala/fs2/kafka/security/KafkaCredentialStore.scala | 6 ++++++ .../src/main/scala/fs2/kafka/security/KeyStoreFile.scala | 6 ++++++ .../main/scala/fs2/kafka/security/KeyStorePassword.scala | 6 ++++++ .../main/scala/fs2/kafka/security/ServiceCertificate.scala | 6 ++++++ .../src/main/scala/fs2/kafka/security/TrustStoreFile.scala | 6 ++++++ .../main/scala/fs2/kafka/security/TrustStorePassword.scala | 6 ++++++ .../scala/fs2/kafka/security/internal/CertificateOps.scala | 6 ++++++ .../main/scala/fs2/kafka/security/internal/FileOps.scala | 6 ++++++ .../core/src/main/scala/fs2/kafka/security/package.scala | 6 ++++++ 11 files changed, 66 insertions(+) diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala index ae4a61172..656faca32 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import java.security.cert.{Certificate, CertificateException} diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala index e1663d172..e48e1808f 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import cats.syntax.all._ diff --git a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala index 931e4649f..9fa6e5ba2 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import cats.effect._ diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala index c34b8b66e..d7044b676 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import cats.effect.Sync diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala index 0927ec7f1..ad46805ca 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import cats.effect.Sync diff --git a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala index 5e0802aad..d34d71836 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import java.security.cert.{Certificate, CertificateException} diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala index 9758d9ff5..1e99343cd 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import cats.effect.Sync diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala index 2524a0d3e..8b702e286 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security import cats.effect.Sync diff --git a/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala b/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala index 6dd8e6e0a..1a3fcc3a3 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security.internal import cats.syntax.all._ diff --git a/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala b/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala index a8580cd78..8b2cab7ed 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka.security.internal import cats.effect.Sync diff --git a/modules/core/src/main/scala/fs2/kafka/security/package.scala b/modules/core/src/main/scala/fs2/kafka/security/package.scala index efadf05a8..4103b750c 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/package.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/package.scala @@ -1,3 +1,9 @@ +/* + * Copyright 2018-2021 OVO Energy Limited + * + * SPDX-License-Identifier: Apache-2.0 + */ + package fs2.kafka import java.nio.charset.StandardCharsets From ee2a05229ff015eb928b2d94077dc29a14f7d0c3 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Sun, 14 Feb 2021 23:12:44 +0000 Subject: [PATCH 06/16] scalafmt --- .../main/scala/fs2/kafka/security/ClientCertificate.scala | 6 ++---- .../main/scala/fs2/kafka/security/ClientPrivateKey.scala | 6 ++---- .../main/scala/fs2/kafka/security/ServiceCertificate.scala | 6 ++---- .../core/src/main/scala/fs2/kafka/security/package.scala | 3 +-- 4 files changed, 7 insertions(+), 14 deletions(-) diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala index 656faca32..7b9a9fef6 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala @@ -13,7 +13,7 @@ sealed abstract class ClientCertificate { } object ClientCertificate { - def fromString(clientCertificate: String): Either[CertificateException, ClientCertificate] = { + def fromString(clientCertificate: String): Either[CertificateException, ClientCertificate] = internal.CertificateOps.loadFromString(clientCertificate).map { certificate => new ClientCertificate { override final val value: Certificate = @@ -23,11 +23,9 @@ object ClientCertificate { s"ClientCertificate(${clientCertificate.valueShortHash})" } } - } - def fromCertificate(certificate: Certificate): ClientCertificate = { + def fromCertificate(certificate: Certificate): ClientCertificate = new ClientCertificate { override def value: Certificate = certificate } - } } diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala index e48e1808f..fe4396baa 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala @@ -18,7 +18,7 @@ sealed abstract class ClientPrivateKey { } object ClientPrivateKey { - def fromString(clientPrivateKey: String): Either[GeneralSecurityException, ClientPrivateKey] = { + def fromString(clientPrivateKey: String): Either[GeneralSecurityException, ClientPrivateKey] = Either.catchOnly[GeneralSecurityException] { new ClientPrivateKey { override final val value: PrivateKey = @@ -40,14 +40,12 @@ object ClientPrivateKey { s"ClientPrivateKey(${clientPrivateKey.valueShortHash})" } } - } - def fromPrivateKey(privateKey: PrivateKey): ClientPrivateKey = { + def fromPrivateKey(privateKey: PrivateKey): ClientPrivateKey = new ClientPrivateKey { override def value: PrivateKey = privateKey override final def toString: String = s"ClientPrivateKey(${privateKey.toString.valueShortHash})" } - } } diff --git a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala index d34d71836..d1c352ad4 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala @@ -13,7 +13,7 @@ sealed abstract class ServiceCertificate { } object ServiceCertificate { - def fromString(serviceCertificate: String): Either[CertificateException, ServiceCertificate] = { + def fromString(serviceCertificate: String): Either[CertificateException, ServiceCertificate] = internal.CertificateOps.loadFromString(serviceCertificate).map { certificate => new ServiceCertificate { override final val value: Certificate = @@ -23,11 +23,9 @@ object ServiceCertificate { s"ServiceCertificate(${serviceCertificate.valueShortHash})" } } - } - def fromCertificate(certificate: Certificate): ServiceCertificate = { + def fromCertificate(certificate: Certificate): ServiceCertificate = new ServiceCertificate { override def value: Certificate = certificate } - } } diff --git a/modules/core/src/main/scala/fs2/kafka/security/package.scala b/modules/core/src/main/scala/fs2/kafka/security/package.scala index 4103b750c..34b1db5d6 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/package.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/package.scala @@ -18,13 +18,12 @@ package object security { private[this] final def hex(in: Array[Byte]): Array[Char] = { val length = in.length - @tailrec def encode(out: Array[Char], i: Int, j: Int): Array[Char] = { + @tailrec def encode(out: Array[Char], i: Int, j: Int): Array[Char] = if (i < length) { out(j) = hexChars((0xf0 & in(i)) >>> 4) out(j + 1) = hexChars(0x0f & in(i)) encode(out, i + 1, j + 2) } else out - } encode(new Array(length << 1), 0, 0) } From bcaca0a1ce796ab9ac912691341dda97b12af5c7 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 15 Feb 2021 17:15:29 +0000 Subject: [PATCH 07/16] Update docs/src/main/mdoc/certificates.md Co-authored-by: Ben Plommer --- docs/src/main/mdoc/certificates.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index 05f0ab5bc..71b834778 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -20,8 +20,8 @@ def loadKafkaSetup[F[_]: Async: ContextShift]( clientPrivateKey: String, clientCertificate: String, serviceCertificate: String, -): Resource[F, KafkaCredentialStore] = - Blocker[F].evalMap { blocker => +): F[KafkaCredentialStore] = + Blocker[F].use { blocker => ( ClientPrivateKey.fromString(clientPrivateKey).liftTo[F], ClientCertificate.fromString(clientCertificate).liftTo[F], @@ -31,4 +31,4 @@ def loadKafkaSetup[F[_]: Async: ContextShift]( KafkaCredentialStore[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) } } -``` \ No newline at end of file +``` From 88cfa0ba0f306714c3a603e434c77b3620b9822a Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 15 Feb 2021 17:42:35 +0000 Subject: [PATCH 08/16] Update certificates.md --- docs/src/main/mdoc/certificates.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index 71b834778..21a5788d5 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -5,7 +5,7 @@ title: Security & Certificates ## Security: certificates, trust stores, and passwords -The `KafkaCredentialStore` can be used to create the necessary. +The `KafkaCredentialStore` can be used to create the necessary trust stores and passwords to access kafka. The parameters passed in are string representations of the client private key, client certificate and service certificate. the `properties` field in `KafkaCredentialStore` can then be applied to @@ -16,7 +16,7 @@ import cats.effect._ import cats.syntax.all._ import fs2.kafka.security._ -def loadKafkaSetup[F[_]: Async: ContextShift]( +def loadKafkaSetup[F[_]: Sync: ContextShift]( clientPrivateKey: String, clientCertificate: String, serviceCertificate: String, From 7af23c04192993d9f9e93772d9a107269bd36127 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 15 Feb 2021 17:44:17 +0000 Subject: [PATCH 09/16] Code review comments --- .../fs2/kafka/security/KafkaCredentialStore.scala | 4 ++-- .../main/scala/fs2/kafka/security/KeyStoreFile.scala | 10 ++-------- .../main/scala/fs2/kafka/security/TrustStoreFile.scala | 10 ++-------- 3 files changed, 6 insertions(+), 18 deletions(-) diff --git a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala index 9fa6e5ba2..db1c9c2dd 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala @@ -132,10 +132,10 @@ object KafkaCredentialStore { override final val properties: Map[String, String] = Map( "security.protocol" -> "SSL", - "ssl.truststore.location" -> trustStoreFile.pathAsString, + "ssl.truststore.location" -> trustStoreFile.path.toString, "ssl.truststore.password" -> trustStorePassword.value, "ssl.keystore.type" -> "PKCS12", - "ssl.keystore.location" -> keyStoreFile.pathAsString, + "ssl.keystore.location" -> keyStoreFile.path.toString, "ssl.keystore.password" -> keyStorePassword.value, "ssl.key.password" -> keyStorePassword.value ) diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala index d7044b676..c7facc1f5 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala @@ -13,22 +13,16 @@ import java.nio.file.Path sealed abstract class KeyStoreFile { def path: Path - - def pathAsString: String } private[security] object KeyStoreFile { final def createTemporary[F[_]](implicit F: Sync[F]): F[KeyStoreFile] = internal.FileOps.createTemp("client.keystore-", ".p12").map { _path => new KeyStoreFile { - override final val path: Path = - _path - - override final def pathAsString: String = - path.toString + override final val path: Path = _path override final def toString: String = - s"KeyStoreFile($pathAsString)" + s"KeyStoreFile(${path.toString})" } } } diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala index 1e99343cd..94e49524d 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala @@ -13,22 +13,16 @@ import java.nio.file.Path sealed abstract class TrustStoreFile { def path: Path - - def pathAsString: String } private[security] object TrustStoreFile { def createTemporary[F[_]](implicit F: Sync[F]): F[TrustStoreFile] = internal.FileOps.createTemp[F]("client.truststore-", ".jks").map { _path => new TrustStoreFile { - override final val path: Path = - _path - - override final def pathAsString: String = - path.toString + override final val path: Path = _path override final def toString: String = - s"TrustStoreFile($pathAsString)" + s"TrustStoreFile(${path.toString})" } } } From 7e8c6f718f8967a1bbdad360dbff952bfb26e58d Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 15 Feb 2021 20:32:48 +0000 Subject: [PATCH 10/16] Add KafkaCredentialStore.createFromStrings --- docs/src/main/mdoc/certificates.md | 16 ++++++---------- .../kafka/security/KafkaCredentialStore.scala | 15 +++++++++++++++ 2 files changed, 21 insertions(+), 10 deletions(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index 21a5788d5..851ef1c26 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -16,19 +16,15 @@ import cats.effect._ import cats.syntax.all._ import fs2.kafka.security._ -def loadKafkaSetup[F[_]: Sync: ContextShift]( +def createKafkaProducer[F[_]: Sync: ContextShift]( clientPrivateKey: String, clientCertificate: String, serviceCertificate: String, -): F[KafkaCredentialStore] = +): F[ProducerSettings[F, UUID, String]] = Blocker[F].use { blocker => - ( - ClientPrivateKey.fromString(clientPrivateKey).liftTo[F], - ClientCertificate.fromString(clientCertificate).liftTo[F], - ServiceCertificate.fromString(serviceCertificate).liftTo[F], - ).tupled.flatMap { - case (clientPrivateKey, clientCertificate, serviceCertificate) => - KafkaCredentialStore[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) - } + KafkaCredentialStore.createFromStrings[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) + }.map { credentialStore => + ProducerSettings(Serializer.uuid[F], Serializer.string[F]) + .withProperties(credentialStore.properties) } ``` diff --git a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala index db1c9c2dd..18b3c7a70 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala @@ -48,6 +48,21 @@ object KafkaCredentialStore { ) } yield setupDetails + final def createFromStrings[F[_]: Sync: ContextShift]( + clientPrivateKey: String, + clientCertificate: String, + serviceCertificate: String, + blocker: Blocker + ): F[KafkaCredentialStore] = + ( + ClientPrivateKey.fromString(clientPrivateKey).liftTo[F], + ClientCertificate.fromString(clientCertificate).liftTo[F], + ServiceCertificate.fromString(serviceCertificate).liftTo[F] + ).tupled.flatMap { + case (clientPrivateKey, clientCertificate, serviceCertificate) => + apply[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) + } + private final def setupStore[F[_]: ContextShift]( storeType: String, storePath: Path, From fc794758e38e63600131da3e32d542b4b28fb21a Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 15 Feb 2021 23:00:44 +0000 Subject: [PATCH 11/16] Add withCredentials(credentialsStore: KafkaCredentialStore) method to *Settings --- docs/src/main/mdoc/certificates.md | 2 +- .../src/main/scala/fs2/kafka/AdminClientSettings.scala | 8 ++++++++ .../core/src/main/scala/fs2/kafka/ConsumerSettings.scala | 8 ++++++++ .../core/src/main/scala/fs2/kafka/ProducerSettings.scala | 8 ++++++++ .../src/main/scala/fs2/kafka/vulcan/AvroSettings.scala | 7 +++++++ .../fs2/kafka/vulcan/SchemaRegistryClientSettings.scala | 7 +++++++ 6 files changed, 39 insertions(+), 1 deletion(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index 851ef1c26..c8b341152 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -25,6 +25,6 @@ def createKafkaProducer[F[_]: Sync: ContextShift]( KafkaCredentialStore.createFromStrings[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) }.map { credentialStore => ProducerSettings(Serializer.uuid[F], Serializer.string[F]) - .withProperties(credentialStore.properties) + .withCredentials(credentialStore) } ``` diff --git a/modules/core/src/main/scala/fs2/kafka/AdminClientSettings.scala b/modules/core/src/main/scala/fs2/kafka/AdminClientSettings.scala index 95c2a9dc5..ae98ba09e 100644 --- a/modules/core/src/main/scala/fs2/kafka/AdminClientSettings.scala +++ b/modules/core/src/main/scala/fs2/kafka/AdminClientSettings.scala @@ -9,7 +9,9 @@ package fs2.kafka import cats.effect.{Blocker, Sync} import cats.Show import fs2.kafka.internal.converters.collection._ +import fs2.kafka.security.KafkaCredentialStore import org.apache.kafka.clients.admin.{AdminClient, AdminClientConfig} + import scala.concurrent.duration._ /** @@ -201,6 +203,12 @@ sealed abstract class AdminClientSettings[F[_]] { def withCreateAdminClient( createAdminClient: Map[String, String] => F[AdminClient] ): AdminClientSettings[F] + + /** + * Includes the credentials properties from the provided [[KafkaCredentialStore]] + */ + def withCredentials(credentialsStore: KafkaCredentialStore): AdminClientSettings[F] = + withProperties(credentialsStore.properties) } object AdminClientSettings { diff --git a/modules/core/src/main/scala/fs2/kafka/ConsumerSettings.scala b/modules/core/src/main/scala/fs2/kafka/ConsumerSettings.scala index 2a97780e7..1f9f37ca1 100644 --- a/modules/core/src/main/scala/fs2/kafka/ConsumerSettings.scala +++ b/modules/core/src/main/scala/fs2/kafka/ConsumerSettings.scala @@ -9,9 +9,11 @@ package fs2.kafka import cats.effect.{Blocker, Sync} import cats.Show import fs2.kafka.internal.converters.collection._ +import fs2.kafka.security.KafkaCredentialStore import org.apache.kafka.clients.consumer.ConsumerConfig import org.apache.kafka.common.requests.OffsetFetchResponse import org.apache.kafka.common.serialization.ByteArrayDeserializer + import scala.concurrent.duration._ /** @@ -397,6 +399,12 @@ sealed abstract class ConsumerSettings[F[_], K, V] { * instead be set to `2` and not the specified value. */ def withMaxPrefetchBatches(maxPrefetchBatches: Int): ConsumerSettings[F, K, V] + + /** + * Includes the credentials properties from the provided [[KafkaCredentialStore]] + */ + def withCredentials(credentialsStore: KafkaCredentialStore): ConsumerSettings[F, K, V] = + withProperties(credentialsStore.properties) } object ConsumerSettings { diff --git a/modules/core/src/main/scala/fs2/kafka/ProducerSettings.scala b/modules/core/src/main/scala/fs2/kafka/ProducerSettings.scala index 6e30a874d..9ad9dcc4b 100644 --- a/modules/core/src/main/scala/fs2/kafka/ProducerSettings.scala +++ b/modules/core/src/main/scala/fs2/kafka/ProducerSettings.scala @@ -9,8 +9,10 @@ package fs2.kafka import cats.effect.{Blocker, Sync} import cats.Show import fs2.kafka.internal.converters.collection._ +import fs2.kafka.security.KafkaCredentialStore import org.apache.kafka.clients.producer.ProducerConfig import org.apache.kafka.common.serialization.ByteArraySerializer + import scala.concurrent.duration._ /** @@ -239,6 +241,12 @@ sealed abstract class ProducerSettings[F[_], K, V] { def withCreateProducer( createProducer: Map[String, String] => F[KafkaByteProducer] ): ProducerSettings[F, K, V] + + /** + * Includes the credentials properties from the provided [[KafkaCredentialStore]] + */ + def withCredentials(credentialsStore: KafkaCredentialStore): ProducerSettings[F, K, V] = + withProperties(credentialsStore.properties) } object ProducerSettings { diff --git a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala index 59a2bde0f..45828b2de 100644 --- a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala +++ b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala @@ -10,6 +10,7 @@ import cats.effect.Sync import cats.implicits._ import fs2.kafka.internal.converters.collection._ import fs2.kafka.internal.syntax._ +import fs2.kafka.security.KafkaCredentialStore /** * Describes how to create a `KafkaAvroDeserializer` and a @@ -118,6 +119,12 @@ sealed abstract class AvroSettings[F[_]] { createAvroSerializerWith: (F[SchemaRegistryClient], Boolean, Map[String, String]) => F[(KafkaAvroSerializer, SchemaRegistryClient)] // format: on ): AvroSettings[F] + + /** + * Includes the credentials properties from the provided [[KafkaCredentialStore]] + */ + def withCredentials(credentialsStore: KafkaCredentialStore): AvroSettings[F] = + withProperties(credentialsStore.properties) } object AvroSettings { diff --git a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala index e7365e0ee..f47926caa 100644 --- a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala +++ b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala @@ -9,6 +9,7 @@ package fs2.kafka.vulcan import cats.effect.Sync import cats.Show import fs2.kafka.internal.converters.collection._ +import fs2.kafka.security.KafkaCredentialStore /** * Describes how to create a `SchemaRegistryClient` and which @@ -83,6 +84,12 @@ sealed abstract class SchemaRegistryClientSettings[F[_]] { def withCreateSchemaRegistryClient( createSchemaRegistryClientWith: (String, Int, Map[String, String]) => F[SchemaRegistryClient] ): SchemaRegistryClientSettings[F] + + /** + * Includes the credentials properties from the provided [[KafkaCredentialStore]] + */ + def withCredentials(credentialsStore: KafkaCredentialStore): SchemaRegistryClientSettings[F] = + withProperties(credentialsStore.properties) } object SchemaRegistryClientSettings { From 83aaa678af84820fca8fff829d231940abefe229 Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 15 Feb 2021 23:52:09 +0000 Subject: [PATCH 12/16] Fix certificates.md --- docs/src/main/mdoc/certificates.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index c8b341152..c2ecedc85 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -14,17 +14,18 @@ any of the `*Settings` classes by using the `withProperties(kafkaCredentialStore ```scala mdoc import cats.effect._ import cats.syntax.all._ +import fs2.kafka._ import fs2.kafka.security._ -def createKafkaProducer[F[_]: Sync: ContextShift]( +def createKafkaProducer[F[_]: Sync: ContextShift, K, V]( clientPrivateKey: String, clientCertificate: String, serviceCertificate: String, -): F[ProducerSettings[F, UUID, String]] = +)(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): F[ProducerSettings[F, K, V]] = Blocker[F].use { blocker => KafkaCredentialStore.createFromStrings[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) }.map { credentialStore => - ProducerSettings(Serializer.uuid[F], Serializer.string[F]) + ProducerSettings(keySer, valSer) .withCredentials(credentialStore) } ``` From eef320890d29a575c7d98503072ca5b5b485d54a Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 8 Mar 2021 10:23:57 +0000 Subject: [PATCH 13/16] Revert withCredentials method on AvroSettings and SchemaRegistryClientSettings --- .../src/main/scala/fs2/kafka/vulcan/AvroSettings.scala | 7 ------- .../fs2/kafka/vulcan/SchemaRegistryClientSettings.scala | 7 ------- 2 files changed, 14 deletions(-) diff --git a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala index 45828b2de..59a2bde0f 100644 --- a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala +++ b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/AvroSettings.scala @@ -10,7 +10,6 @@ import cats.effect.Sync import cats.implicits._ import fs2.kafka.internal.converters.collection._ import fs2.kafka.internal.syntax._ -import fs2.kafka.security.KafkaCredentialStore /** * Describes how to create a `KafkaAvroDeserializer` and a @@ -119,12 +118,6 @@ sealed abstract class AvroSettings[F[_]] { createAvroSerializerWith: (F[SchemaRegistryClient], Boolean, Map[String, String]) => F[(KafkaAvroSerializer, SchemaRegistryClient)] // format: on ): AvroSettings[F] - - /** - * Includes the credentials properties from the provided [[KafkaCredentialStore]] - */ - def withCredentials(credentialsStore: KafkaCredentialStore): AvroSettings[F] = - withProperties(credentialsStore.properties) } object AvroSettings { diff --git a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala index f47926caa..e7365e0ee 100644 --- a/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala +++ b/modules/vulcan/src/main/scala/fs2/kafka/vulcan/SchemaRegistryClientSettings.scala @@ -9,7 +9,6 @@ package fs2.kafka.vulcan import cats.effect.Sync import cats.Show import fs2.kafka.internal.converters.collection._ -import fs2.kafka.security.KafkaCredentialStore /** * Describes how to create a `SchemaRegistryClient` and which @@ -84,12 +83,6 @@ sealed abstract class SchemaRegistryClientSettings[F[_]] { def withCreateSchemaRegistryClient( createSchemaRegistryClientWith: (String, Int, Map[String, String]) => F[SchemaRegistryClient] ): SchemaRegistryClientSettings[F] - - /** - * Includes the credentials properties from the provided [[KafkaCredentialStore]] - */ - def withCredentials(credentialsStore: KafkaCredentialStore): SchemaRegistryClientSettings[F] = - withProperties(credentialsStore.properties) } object SchemaRegistryClientSettings { From 9d3695573713717e1994493e35094f81a695944c Mon Sep 17 00:00:00 2001 From: Andrew Gustafson Date: Mon, 8 Mar 2021 12:11:11 +0000 Subject: [PATCH 14/16] Create PemKafkaCredentialStore --- docs/src/main/mdoc/certificates.md | 30 +++++++++++++- .../kafka/security/KafkaCredentialStore.scala | 41 +++++++++++++++---- 2 files changed, 60 insertions(+), 11 deletions(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index c2ecedc85..e85794eee 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -17,15 +17,41 @@ import cats.syntax.all._ import fs2.kafka._ import fs2.kafka.security._ -def createKafkaProducer[F[_]: Sync: ContextShift, K, V]( +def createKafkaProducerUsingPkcs12[F[_]: Sync: ContextShift, K, V]( clientPrivateKey: String, clientCertificate: String, serviceCertificate: String, )(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): F[ProducerSettings[F, K, V]] = Blocker[F].use { blocker => - KafkaCredentialStore.createFromStrings[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) + Pkcs12KafkaCredentialStore.createFromStrings[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) }.map { credentialStore => ProducerSettings(keySer, valSer) .withCredentials(credentialStore) } + +def createKafkaProducerUsingPem[F[_]: Sync: ContextShift, K, V]( + caCertificate: String, + accessKey: String, + accessCertificate: String, +)(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): ProducerSettings[F, K, V] = { + val credentialStore = PemKafkaCredentialStore.createFromStrings[F]( + s""" + -----BEGIN CERTIFICATE-----" + $caCertificate + -----END CERTIFICATE----- + """.replace('\n', ' '), + s""" + -----BEGIN PRIVATE KEY-----" + $accessKey + -----END PRIVATE KEY----- + """.replace('\n', ' '), + s""" + -----BEGIN CERTIFICATE-----" + $accessCertificate + -----END CERTIFICATE----- + """.replace('\n', ' ') + ) + ProducerSettings(keySer, valSer) + .withCredentials(credentialStore) + } ``` diff --git a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala index 18b3c7a70..99a3fd360 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala @@ -12,7 +12,11 @@ import cats.syntax.all._ import java.nio.file.{Files, Path} import java.security.KeyStore -sealed abstract class KafkaCredentialStore { +sealed trait KafkaCredentialStore { + def properties: Map[String, String] +} + +sealed abstract class Pkcs12KafkaCredentialStore extends KafkaCredentialStore { def keyStoreFile: KeyStoreFile def keyStorePassword: KeyStorePassword @@ -20,19 +24,17 @@ sealed abstract class KafkaCredentialStore { def trustStoreFile: TrustStoreFile def trustStorePassword: TrustStorePassword - - def properties: Map[String, String] } -object KafkaCredentialStore { +object Pkcs12KafkaCredentialStore { final def apply[F[_]: ContextShift]( clientPrivateKey: ClientPrivateKey, clientCertificate: ClientCertificate, serviceCertificate: ServiceCertificate, blocker: Blocker - )(implicit F: Sync[F]): F[KafkaCredentialStore] = + )(implicit F: Sync[F]): F[Pkcs12KafkaCredentialStore] = for { - setupDetails <- KafkaCredentialStore.createTemporary[F] + setupDetails <- Pkcs12KafkaCredentialStore.createTemporary[F] _ <- setupKeyStore( clientPrivateKey = clientPrivateKey, clientCertificate = clientCertificate, @@ -53,7 +55,7 @@ object KafkaCredentialStore { clientCertificate: String, serviceCertificate: String, blocker: Blocker - ): F[KafkaCredentialStore] = + ): F[Pkcs12KafkaCredentialStore] = ( ClientPrivateKey.fromString(clientPrivateKey).liftTo[F], ClientCertificate.fromString(clientCertificate).liftTo[F], @@ -124,14 +126,14 @@ object KafkaCredentialStore { blocker = blocker ) - private final def createTemporary[F[_]](implicit F: Sync[F]): F[KafkaCredentialStore] = + private final def createTemporary[F[_]](implicit F: Sync[F]): F[Pkcs12KafkaCredentialStore] = for { _keyStoreFile <- KeyStoreFile.createTemporary[F] _keyStorePassword <- KeyStorePassword.createTemporary[F] _trustStoreFile <- TrustStoreFile.createTemporary[F] _trustStorePassword <- TrustStorePassword.createTemporary[F] } yield { - new KafkaCredentialStore { + new Pkcs12KafkaCredentialStore { override final val keyStoreFile: KeyStoreFile = _keyStoreFile @@ -160,3 +162,24 @@ object KafkaCredentialStore { } } } + +sealed abstract class PemKafkaCredentialStore extends KafkaCredentialStore + +object PemKafkaCredentialStore { + final def createFromStrings[F[_]: Sync: ContextShift]( + caCertificate: String, + accessKey: String, + accessCertificate: String + ): PemKafkaCredentialStore = + new PemKafkaCredentialStore { + override def properties: Map[String, String] = + Map( + "security.protocol" -> "SSL", + "ssl.truststore.type" -> "PEM", + "ssl.truststore.certificates" -> caCertificate, + "ssl.keystore.type" -> "PEM", + "ssl.keystore.key" -> accessKey, + "ssl.keystore.certificate.chain" -> accessCertificate + ) + } +} From 266409cb60818539fd77124472187b7792a4c080 Mon Sep 17 00:00:00 2001 From: Ben Plommer Date: Mon, 5 Apr 2021 13:07:01 +0100 Subject: [PATCH 15/16] Simplify + add tests --- docs/src/main/mdoc/certificates.md | 49 ++--- .../kafka/security/ClientCertificate.scala | 31 ---- .../fs2/kafka/security/ClientPrivateKey.scala | 51 ----- .../kafka/security/KafkaCredentialStore.scala | 175 +----------------- .../fs2/kafka/security/KeyStoreFile.scala | 28 --- .../fs2/kafka/security/KeyStorePassword.scala | 30 --- .../kafka/security/ServiceCertificate.scala | 31 ---- .../fs2/kafka/security/TrustStoreFile.scala | 28 --- .../kafka/security/TrustStorePassword.scala | 30 --- .../security/internal/CertificateOps.scala | 25 --- .../fs2/kafka/security/internal/FileOps.scala | 20 -- .../scala/fs2/kafka/security/package.scala | 40 ---- .../security/KafkaCredentialStoreSpec.scala | 48 +++++ 13 files changed, 70 insertions(+), 516 deletions(-) delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala delete mode 100644 modules/core/src/main/scala/fs2/kafka/security/package.scala create mode 100644 modules/core/src/test/scala/fs2/kafka/security/KafkaCredentialStoreSpec.scala diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index e85794eee..b8ec0b994 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -13,45 +13,20 @@ any of the `*Settings` classes by using the `withProperties(kafkaCredentialStore ```scala mdoc import cats.effect._ -import cats.syntax.all._ import fs2.kafka._ import fs2.kafka.security._ -def createKafkaProducerUsingPkcs12[F[_]: Sync: ContextShift, K, V]( - clientPrivateKey: String, - clientCertificate: String, - serviceCertificate: String, -)(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): F[ProducerSettings[F, K, V]] = - Blocker[F].use { blocker => - Pkcs12KafkaCredentialStore.createFromStrings[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) - }.map { credentialStore => - ProducerSettings(keySer, valSer) - .withCredentials(credentialStore) - } - -def createKafkaProducerUsingPem[F[_]: Sync: ContextShift, K, V]( - caCertificate: String, - accessKey: String, - accessCertificate: String, -)(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): ProducerSettings[F, K, V] = { - val credentialStore = PemKafkaCredentialStore.createFromStrings[F]( - s""" - -----BEGIN CERTIFICATE-----" - $caCertificate - -----END CERTIFICATE----- - """.replace('\n', ' '), - s""" - -----BEGIN PRIVATE KEY-----" - $accessKey - -----END PRIVATE KEY----- - """.replace('\n', ' '), - s""" - -----BEGIN CERTIFICATE-----" - $accessCertificate - -----END CERTIFICATE----- - """.replace('\n', ' ') +def createKafkaProducerUsingPem[F[_]: Sync, K, V]( + caCertificate: String, + accessKey: String, + accessCertificate: String +)(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): ProducerSettings[F, K, V] = + ProducerSettings[F, K, V] + .withCredentials( + PemKafkaCredentialStore.fromPemStrings( + caCertificate, + accessKey, + accessCertificate + ) ) - ProducerSettings(keySer, valSer) - .withCredentials(credentialStore) - } ``` diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala deleted file mode 100644 index 7b9a9fef6..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/ClientCertificate.scala +++ /dev/null @@ -1,31 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import java.security.cert.{Certificate, CertificateException} - -sealed abstract class ClientCertificate { - def value: Certificate -} - -object ClientCertificate { - def fromString(clientCertificate: String): Either[CertificateException, ClientCertificate] = - internal.CertificateOps.loadFromString(clientCertificate).map { certificate => - new ClientCertificate { - override final val value: Certificate = - certificate - - override final def toString: String = - s"ClientCertificate(${clientCertificate.valueShortHash})" - } - } - - def fromCertificate(certificate: Certificate): ClientCertificate = - new ClientCertificate { - override def value: Certificate = certificate - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala b/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala deleted file mode 100644 index fe4396baa..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/ClientPrivateKey.scala +++ /dev/null @@ -1,51 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import cats.syntax.all._ - -import java.nio.charset.StandardCharsets -import java.security.{GeneralSecurityException, KeyFactory, PrivateKey} -import java.security.spec.PKCS8EncodedKeySpec -import java.util.Base64 - -sealed abstract class ClientPrivateKey { - def value: PrivateKey -} - -object ClientPrivateKey { - def fromString(clientPrivateKey: String): Either[GeneralSecurityException, ClientPrivateKey] = - Either.catchOnly[GeneralSecurityException] { - new ClientPrivateKey { - override final val value: PrivateKey = - KeyFactory - .getInstance("RSA") - .generatePrivate { - new PKCS8EncodedKeySpec( - Base64.getDecoder.decode { - clientPrivateKey - .replace("-----BEGIN PRIVATE KEY-----", "") - .replace("-----END PRIVATE KEY-----", "") - .filterNot(_.isWhitespace) - .getBytes(StandardCharsets.UTF_8) - } - ) - } - - override final def toString: String = - s"ClientPrivateKey(${clientPrivateKey.valueShortHash})" - } - } - - def fromPrivateKey(privateKey: PrivateKey): ClientPrivateKey = - new ClientPrivateKey { - override def value: PrivateKey = privateKey - - override final def toString: String = - s"ClientPrivateKey(${privateKey.toString.valueShortHash})" - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala index 99a3fd360..b3ffc67b4 100644 --- a/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala +++ b/modules/core/src/main/scala/fs2/kafka/security/KafkaCredentialStore.scala @@ -6,180 +6,25 @@ package fs2.kafka.security -import cats.effect._ -import cats.syntax.all._ - -import java.nio.file.{Files, Path} -import java.security.KeyStore - sealed trait KafkaCredentialStore { def properties: Map[String, String] } -sealed abstract class Pkcs12KafkaCredentialStore extends KafkaCredentialStore { - def keyStoreFile: KeyStoreFile - - def keyStorePassword: KeyStorePassword - - def trustStoreFile: TrustStoreFile - - def trustStorePassword: TrustStorePassword -} - -object Pkcs12KafkaCredentialStore { - final def apply[F[_]: ContextShift]( - clientPrivateKey: ClientPrivateKey, - clientCertificate: ClientCertificate, - serviceCertificate: ServiceCertificate, - blocker: Blocker - )(implicit F: Sync[F]): F[Pkcs12KafkaCredentialStore] = - for { - setupDetails <- Pkcs12KafkaCredentialStore.createTemporary[F] - _ <- setupKeyStore( - clientPrivateKey = clientPrivateKey, - clientCertificate = clientCertificate, - keyStoreFile = setupDetails.keyStoreFile, - keyStorePassword = setupDetails.keyStorePassword, - blocker = blocker - ) - _ <- setupTrustStore( - serviceCertificate = serviceCertificate, - trustStoreFile = setupDetails.trustStoreFile, - trustStorePassword = setupDetails.trustStorePassword, - blocker = blocker - ) - } yield setupDetails - - final def createFromStrings[F[_]: Sync: ContextShift]( - clientPrivateKey: String, - clientCertificate: String, - serviceCertificate: String, - blocker: Blocker - ): F[Pkcs12KafkaCredentialStore] = - ( - ClientPrivateKey.fromString(clientPrivateKey).liftTo[F], - ClientCertificate.fromString(clientCertificate).liftTo[F], - ServiceCertificate.fromString(serviceCertificate).liftTo[F] - ).tupled.flatMap { - case (clientPrivateKey, clientCertificate, serviceCertificate) => - apply[F](clientPrivateKey, clientCertificate, serviceCertificate, blocker) - } - - private final def setupStore[F[_]: ContextShift]( - storeType: String, - storePath: Path, - storePasswordChars: Array[Char], - setupStore: KeyStore => Unit, - blocker: Blocker - )(implicit F: Sync[F]): F[Unit] = - blocker.delay { - val keyStore = KeyStore.getInstance(storeType) - keyStore.load(null, storePasswordChars) - setupStore(keyStore) - - val outputStream = Files.newOutputStream(storePath) - - try { - keyStore.store(outputStream, storePasswordChars) - } finally { - outputStream.close() - } - } - - private final def setupKeyStore[F[_]: ContextShift]( - clientPrivateKey: ClientPrivateKey, - clientCertificate: ClientCertificate, - keyStoreFile: KeyStoreFile, - keyStorePassword: KeyStorePassword, - blocker: Blocker - )(implicit F: Sync[F]): F[Unit] = { - val keyStorePasswordChars = - keyStorePassword.value.toCharArray - - setupStore( - storeType = "PKCS12", - storePath = keyStoreFile.path, - storePasswordChars = keyStorePasswordChars, - setupStore = _.setEntry( - "service_key", - new KeyStore.PrivateKeyEntry( - clientPrivateKey.value, - Array(clientCertificate.value) - ), - new KeyStore.PasswordProtection(keyStorePasswordChars) - ), - blocker = blocker - ) - } - - private final def setupTrustStore[F[_]: ContextShift]( - serviceCertificate: ServiceCertificate, - trustStoreFile: TrustStoreFile, - trustStorePassword: TrustStorePassword, - blocker: Blocker - )(implicit F: Sync[F]): F[Unit] = - setupStore( - storeType = "JKS", - storePath = trustStoreFile.path, - storePasswordChars = trustStorePassword.value.toCharArray, - setupStore = _.setCertificateEntry("CA", serviceCertificate.value), - blocker = blocker - ) - - private final def createTemporary[F[_]](implicit F: Sync[F]): F[Pkcs12KafkaCredentialStore] = - for { - _keyStoreFile <- KeyStoreFile.createTemporary[F] - _keyStorePassword <- KeyStorePassword.createTemporary[F] - _trustStoreFile <- TrustStoreFile.createTemporary[F] - _trustStorePassword <- TrustStorePassword.createTemporary[F] - } yield { - new Pkcs12KafkaCredentialStore { - override final val keyStoreFile: KeyStoreFile = - _keyStoreFile - - override final val keyStorePassword: KeyStorePassword = - _keyStorePassword - - override final val trustStoreFile: TrustStoreFile = - _trustStoreFile - - override final val trustStorePassword: TrustStorePassword = - _trustStorePassword - - override final val properties: Map[String, String] = - Map( - "security.protocol" -> "SSL", - "ssl.truststore.location" -> trustStoreFile.path.toString, - "ssl.truststore.password" -> trustStorePassword.value, - "ssl.keystore.type" -> "PKCS12", - "ssl.keystore.location" -> keyStoreFile.path.toString, - "ssl.keystore.password" -> keyStorePassword.value, - "ssl.key.password" -> keyStorePassword.value - ) - - override final def toString: String = - s"KafkaCredentialStore($keyStoreFile, $keyStorePassword, $trustStoreFile, $trustStorePassword)" - } - } -} - -sealed abstract class PemKafkaCredentialStore extends KafkaCredentialStore - -object PemKafkaCredentialStore { - final def createFromStrings[F[_]: Sync: ContextShift]( +object KafkaCredentialStore { + final def fromPemStrings( caCertificate: String, - accessKey: String, - accessCertificate: String - ): PemKafkaCredentialStore = - new PemKafkaCredentialStore { - override def properties: Map[String, String] = + clientPrivateKey: String, + clientCertificate: String + ): KafkaCredentialStore = + new KafkaCredentialStore { + override val properties: Map[String, String] = Map( "security.protocol" -> "SSL", "ssl.truststore.type" -> "PEM", - "ssl.truststore.certificates" -> caCertificate, + "ssl.truststore.certificates" -> caCertificate.replace("\n", ""), "ssl.keystore.type" -> "PEM", - "ssl.keystore.key" -> accessKey, - "ssl.keystore.certificate.chain" -> accessCertificate + "ssl.keystore.key" -> clientPrivateKey.replace("\n", ""), + "ssl.keystore.certificate.chain" -> clientCertificate.replace("\n", "") ) } } diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala deleted file mode 100644 index c7facc1f5..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala +++ /dev/null @@ -1,28 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import cats.effect.Sync -import cats.syntax.all._ - -import java.nio.file.Path - -sealed abstract class KeyStoreFile { - def path: Path -} - -private[security] object KeyStoreFile { - final def createTemporary[F[_]](implicit F: Sync[F]): F[KeyStoreFile] = - internal.FileOps.createTemp("client.keystore-", ".p12").map { _path => - new KeyStoreFile { - override final val path: Path = _path - - override final def toString: String = - s"KeyStoreFile(${path.toString})" - } - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala b/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala deleted file mode 100644 index ad46805ca..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/KeyStorePassword.scala +++ /dev/null @@ -1,30 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import cats.effect.Sync - -import java.util.UUID - -sealed abstract class KeyStorePassword { - def value: String -} - -private[security] object KeyStorePassword { - def createTemporary[F[_]](implicit F: Sync[F]): F[KeyStorePassword] = - F.delay { - val _value = UUID.randomUUID().toString - - new KeyStorePassword { - override final val value: String = - _value - - override final def toString: String = - s"KeyStorePassword(${value.valueShortHash})" - } - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala b/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala deleted file mode 100644 index d1c352ad4..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/ServiceCertificate.scala +++ /dev/null @@ -1,31 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import java.security.cert.{Certificate, CertificateException} - -sealed abstract class ServiceCertificate { - def value: Certificate -} - -object ServiceCertificate { - def fromString(serviceCertificate: String): Either[CertificateException, ServiceCertificate] = - internal.CertificateOps.loadFromString(serviceCertificate).map { certificate => - new ServiceCertificate { - override final val value: Certificate = - certificate - - override final def toString: String = - s"ServiceCertificate(${serviceCertificate.valueShortHash})" - } - } - - def fromCertificate(certificate: Certificate): ServiceCertificate = - new ServiceCertificate { - override def value: Certificate = certificate - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala deleted file mode 100644 index 94e49524d..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/TrustStoreFile.scala +++ /dev/null @@ -1,28 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import cats.effect.Sync -import cats.syntax.all._ - -import java.nio.file.Path - -sealed abstract class TrustStoreFile { - def path: Path -} - -private[security] object TrustStoreFile { - def createTemporary[F[_]](implicit F: Sync[F]): F[TrustStoreFile] = - internal.FileOps.createTemp[F]("client.truststore-", ".jks").map { _path => - new TrustStoreFile { - override final val path: Path = _path - - override final def toString: String = - s"TrustStoreFile(${path.toString})" - } - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala b/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala deleted file mode 100644 index 8b702e286..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/TrustStorePassword.scala +++ /dev/null @@ -1,30 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security - -import cats.effect.Sync - -import java.util.UUID - -sealed abstract class TrustStorePassword { - def value: String -} - -private[security] object TrustStorePassword { - def createTemporary[F[_]](implicit F: Sync[F]): F[TrustStorePassword] = - F.delay { - val _value = UUID.randomUUID().toString - - new TrustStorePassword { - override final val value: String = - _value - - override final def toString: String = - s"TrustStorePassword(${value.valueShortHash})" - } - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala b/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala deleted file mode 100644 index 1a3fcc3a3..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/internal/CertificateOps.scala +++ /dev/null @@ -1,25 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security.internal - -import cats.syntax.all._ - -import java.io.ByteArrayInputStream -import java.nio.charset.StandardCharsets -import java.security.cert.{Certificate, CertificateException, CertificateFactory} - -private[security] object CertificateOps { - def loadFromString(certificate: String): Either[CertificateException, Certificate] = - loadFromBytes(certificate.getBytes(StandardCharsets.UTF_8)) - - def loadFromBytes(certificateBytes: Array[Byte]): Either[CertificateException, Certificate] = - Either.catchOnly[CertificateException] { - CertificateFactory - .getInstance("X.509") - .generateCertificate(new ByteArrayInputStream(certificateBytes)) - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala b/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala deleted file mode 100644 index 8b2cab7ed..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/internal/FileOps.scala +++ /dev/null @@ -1,20 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka.security.internal - -import cats.effect.Sync - -import java.nio.file.{Files, Path} - -private[security] object FileOps { - def createTemp[F[_]](prefix: String, suffix: String)(implicit F: Sync[F]): F[Path] = F.delay { - val path = Files.createTempFile(prefix, suffix) - path.toFile.deleteOnExit() - Files.delete(path) - path - } -} diff --git a/modules/core/src/main/scala/fs2/kafka/security/package.scala b/modules/core/src/main/scala/fs2/kafka/security/package.scala deleted file mode 100644 index 34b1db5d6..000000000 --- a/modules/core/src/main/scala/fs2/kafka/security/package.scala +++ /dev/null @@ -1,40 +0,0 @@ -/* - * Copyright 2018-2021 OVO Energy Limited - * - * SPDX-License-Identifier: Apache-2.0 - */ - -package fs2.kafka - -import java.nio.charset.StandardCharsets -import java.security.MessageDigest -import scala.annotation.tailrec - -package object security { - private[this] final val hexChars: Array[Char] = - Array('0', '1', '2', '3', '4', '5', '6', '7', '8', '9', 'a', 'b', 'c', 'd', 'e', 'f') - - implicit class StringOps(val str: String) extends AnyVal { - private[this] final def hex(in: Array[Byte]): Array[Char] = { - val length = in.length - - @tailrec def encode(out: Array[Char], i: Int, j: Int): Array[Char] = - if (i < length) { - out(j) = hexChars((0xf0 & in(i)) >>> 4) - out(j + 1) = hexChars(0x0f & in(i)) - encode(out, i + 1, j + 2) - } else out - - encode(new Array(length << 1), 0, 0) - } - - private[this] final def sha1(bytes: Array[Byte]): Array[Byte] = - MessageDigest.getInstance("SHA-1").digest(bytes) - - def sha1Hex: String = - new String(hex(sha1(str.getBytes(StandardCharsets.UTF_8)))) - - final def valueShortHash: String = - sha1Hex.take(7) - } -} diff --git a/modules/core/src/test/scala/fs2/kafka/security/KafkaCredentialStoreSpec.scala b/modules/core/src/test/scala/fs2/kafka/security/KafkaCredentialStoreSpec.scala new file mode 100644 index 000000000..cd93999d6 --- /dev/null +++ b/modules/core/src/test/scala/fs2/kafka/security/KafkaCredentialStoreSpec.scala @@ -0,0 +1,48 @@ +package fs2.kafka.security + +import fs2.kafka.BaseSpec + +final class KafkaCredentialStoreSpec extends BaseSpec { + describe("KafkaCredentialStore") { + describe("fromPemStrigs") { + it("should create a KafkaCredentialStore with the expected properties") { + val caCert = + """ + |-----BEGIN CERTIFICATE----- + |RmFrZSBDQSBjZXJ0aWZpY2F0ZSBGYWtlIENBIGNlcnRpZmljYXRlIEZha2UgQ0EgY2VydGlmaWNh + |dGUgRmFrZSBDQSBjZXJ0aWZpY2F0ZQ== + |-----END CERTIFICATE----- + |""".stripMargin + + val privateKey = + """ + |-----BEGIN PRIVATE KEY----- + |RmFrZSBwcml2YXRlIGtleSBGYWtlIHByaXZhdGUga2V5IEZha2UgcHJpdmF0ZSBrZXkgRmFrZSBw + |cml2YXRlIGtleSBGYWtlIHByaXZhdGUga2V5IA== + |-----END PRIVATE KEY----- + |""".stripMargin + + val clientCert = + """ + |-----BEGIN CERTIFICATE----- + |RmFrZSBjbGllbnQgY2VydCBGYWtlIGNsaWVudCBjZXJ0IEZha2UgY2xpZW50IGNlcnQgRmFrZSBj + |bGllbnQgY2VydCBGYWtlIGNsaWVudCBjZXJ0IA== + |-----END CERTIFICATE----- + |""".stripMargin + + val store = KafkaCredentialStore.fromPemStrings(caCert, privateKey, clientCert) + + assert( + store.properties === Map( + "security.protocol" -> "SSL", + "ssl.truststore.type" -> "PEM", + "ssl.truststore.certificates" -> "-----BEGIN CERTIFICATE-----RmFrZSBDQSBjZXJ0aWZpY2F0ZSBGYWtlIENBIGNlcnRpZmljYXRlIEZha2UgQ0EgY2VydGlmaWNhdGUgRmFrZSBDQSBjZXJ0aWZpY2F0ZQ==-----END CERTIFICATE-----", + "ssl.keystore.type" -> "PEM", + "ssl.keystore.key" -> "-----BEGIN PRIVATE KEY-----RmFrZSBwcml2YXRlIGtleSBGYWtlIHByaXZhdGUga2V5IEZha2UgcHJpdmF0ZSBrZXkgRmFrZSBwcml2YXRlIGtleSBGYWtlIHByaXZhdGUga2V5IA==-----END PRIVATE KEY-----", + "ssl.keystore.certificate.chain" -> "-----BEGIN CERTIFICATE-----RmFrZSBjbGllbnQgY2VydCBGYWtlIGNsaWVudCBjZXJ0IEZha2UgY2xpZW50IGNlcnQgRmFrZSBjbGllbnQgY2VydCBGYWtlIGNsaWVudCBjZXJ0IA==-----END CERTIFICATE-----" + ) + ) + } + } + } +} From 53cf695fc5b24f54b648d7b790d515cdfb21e63f Mon Sep 17 00:00:00 2001 From: Ben Plommer Date: Mon, 5 Apr 2021 16:37:05 +0100 Subject: [PATCH 16/16] Update certificates.md --- docs/src/main/mdoc/certificates.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/src/main/mdoc/certificates.md b/docs/src/main/mdoc/certificates.md index b8ec0b994..b2bd9520f 100644 --- a/docs/src/main/mdoc/certificates.md +++ b/docs/src/main/mdoc/certificates.md @@ -23,7 +23,7 @@ def createKafkaProducerUsingPem[F[_]: Sync, K, V]( )(implicit keySer: Serializer[F, K], valSer: Serializer[F, V]): ProducerSettings[F, K, V] = ProducerSettings[F, K, V] .withCredentials( - PemKafkaCredentialStore.fromPemStrings( + KafkaCredentialStore.fromPemStrings( caCertificate, accessKey, accessCertificate