Skip to content
Merged
Show file tree
Hide file tree
Changes from 14 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 31 additions & 0 deletions docs/src/main/mdoc/certificates.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
---
id: security
title: Security & Certificates
---

## Security: certificates, trust stores, and passwords

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
any of the `*Settings` classes by using the `withProperties(kafkaCredentialStore.properties)`.
Comment thread
agustafson marked this conversation as resolved.

```scala mdoc
import cats.effect._
import cats.syntax.all._
import fs2.kafka._
import fs2.kafka.security._

def createKafkaProducer[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)
}.map { credentialStore =>
ProducerSettings(keySer, valSer)
.withCredentials(credentialStore)
}
```
Original file line number Diff line number Diff line change
Expand Up @@ -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._

/**
Expand Down Expand Up @@ -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 {
Expand Down
8 changes: 8 additions & 0 deletions modules/core/src/main/scala/fs2/kafka/ConsumerSettings.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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._

/**
Expand Down Expand Up @@ -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 {
Expand Down
8 changes: 8 additions & 0 deletions modules/core/src/main/scala/fs2/kafka/ProducerSettings.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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._

/**
Expand Down Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
/*
* 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
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
/*
* 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})"
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
/*
* Copyright 2018-2021 OVO Energy Limited
*
* SPDX-License-Identifier: Apache-2.0
*/

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

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,
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.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)"
}
}
}
28 changes: 28 additions & 0 deletions modules/core/src/main/scala/fs2/kafka/security/KeyStoreFile.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
/*
* 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})"
}
}
}
Loading