Kafka-encryption-interceptor#

Представляет собой интерсептор, использующий симметричные ключи. Имеет два режима работы: шифрование сообщений (encrypt) и создание к ним подписей (signature).

Предусловия#

Не требуются.

Последовательность выполнения#

Поддерживаемые алгоритмы#

Алгоритмы для шифрования#

Для шифрования сообщений доступны следующие алгоритмы:

  • AES_128/GCM/NoPadding

  • AES_256/GCM/NoPadding

  • ChaCha20-Poly1305-NoPadding

Алгоритмы для подписей#

Поддерживает следующие алгоритмы для создания подписей:

  • AES_128/GCM/NoPadding

  • AES_256/GCM/NoPadding

  • ChaCha20-Poly1305-NoPadding

  • HmacSHA256

Хранение и обновление ключей#

Все ключи хранятся локально в кеше. Обновление ключей происходит с заданным таймингом в параметре interceptor.encryption.storage.update.time. Ключи могут быть получены из различных источников, таких как JSON-файл или Vault.

Перед выгрузкой ключей из JSON-файла происходит сверка даты последнего изменения файла и, если эта дата различается, ключи будут загружены.

Ключи из хранилища Vault выгружаются без проверок.

При запросе ключа, которого нет в локальном хранилище, будет произведена загрузка ключей. Если ключ с таким же именем будет повторно запрошен и не найден в течение времени, установленного с помощью параметра interceptor.encryption.key.reload.time, загрузка не будет произведена.

Преобразование пароля в ключ#

Перехватчик поддерживает режим преобразования пароля в ключ. Для этого в параметрах хранилища должен быть установлен соответствующий префикс – это позволяет использовать пароли в качестве исходных данных для генерации симметричных ключей.

Пример JSON-файла:

{
  "aes256": "generate_key: password: password, algorithm: AES, iterations: 60000, length: 256",
  "aes128": "generate_key: password: password2, algorithm: AES, iterations: 60000, length: 128",
  "chacha20": "generate_key: password: password3, algorithm: ChaCha20, iterations: 60000, length: 256"
}

Маска для генерации ключа:

"aes256": "generate_key: password: <пароль>, algorithm: <>, iterations: <>, length: <>"

Таблица параметров#

Ключ

Дефолтное значение

Описание параметра

interceptor.encryption.mode

failOnValue

Режим обработки ошибок ( failOnValue, failOnSend, addErrorHeader, filter, failOnConsume), описание режимов ниже

interceptor.encryption.authentication.tag.length

128

Длина тега

interceptor.encryption.algorithm

для режима подписи: HmacSHA256, для режима шифрования: AES_256/GCM/NoPadding

Алгоритм шифрования сообщения

interceptor.encryption.key.alias

-

Имя ключа для подписания/шифрования сообщений. Параметр является обязательным для продюсера, для консьюмера – необязательным

interceptor.encryption.key.store.type

-

Тип хранилища ключей (json, vault)

interceptor.encryption.json.storage.location

-

Путь до json-файла с ключами

interceptor.encryption.storage.update.time

60 (значение в секундах)

Частота обновление ключей в локальном хранилище

interceptor.encryption.key.reload.time

5 (значение в секундах)

Время повторного обновления ключей при запросе ключа отсутствующего в хранилище

interceptor.encryption.generate.key.password

true

Включает создание ключей на основе паролей

interceptor.encryption.mode.operation

-

Режим работы перехватчика (encrypt, signature)

Параметр interceptor.encryption.performance.logging является устаревшим и более не поддерживается.

Поведение при ошибках#

Поведение при ошибках валидации настраивается с помощью параметра interceptor.encryption.mode:

В случае, если сообщение не прошло валидацию при отправке producer.send() возможны следующие режимы:

  1. failOnValue (используется по умолчанию) – клиент получит сообщение-заглушку вместо невалидного сообщения, которое содержит null вместо value и выбросит исключение ru.sbt.ss.kafka.interceptors.ProducerInterceptorException при вызове record.value(). Данный способ позволяет использовать логику обработки ошибок kafka-клиента, в том числе вызывать методы перехватчиков (но не callback).

  2. failOnSend – будет выброшено исключение ru.sbt.ss.kafka.interceptors.ProducerInterceptorError (extends Throwable) с сообщением, игнорируя логику обработки ошибок kafka-клиента.

  3. addErrorHeader – в сообщение будет добавлен заголовок с сообщением об ошибке валидации (по умолчанию __interceptor.error настраивается с помощью параметра interceptor.encryption.error.header.name).

Сообщения исключений имеют формат Error while processing record(topic: topic): ошибка валидации в зависимости от формата сообщения.

В случае, если сообщение не прошло валидацию при получении consumer.poll() поведение настраивается параметром interceptor.encryption.mode:

  1. failOnValue (используется по умолчанию) – клиент получит сообщение-заглушку вместо невалидного сообщения, которое содержит null вместо value и выбросит исключение ru.sbt.ss.kafka.interceptors.ConsumerInterceptorException при вызове record.value().

  2. filter – сообщение об ошибке будет залогировано в error, клиент не получит невалидное сообщение.

  3. failOnConsume – метод consumer.poll() выбросит исключение ru.sbt.ss.kafka.interceptors.ConsumerInterceptorError, клиент не получит ни одного сообщения из пачки.

  4. addErrorHeader – в сообщение будет добавлен заголовок с сообщением об ошибке валидации (по умолчанию __interceptor.error, настраивается с помощью параметра interceptor.validator.error.header.name).

Пример конфигурации#

JSON#

Пример конфигурации Kafka Producer#
bootstrap.servers = localhost:9092
security.protocol = PLAINTEXT
group.id = test-group

key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer

interceptor.classes = ru.sbt.ss.kafka.encryption.interceptor.EncryptionProducerInterceptor

interceptor.encryption.algorithm = AES_128/GCM/NoPadding # Алгоритм шифрования сообщения
interceptor.encryption.key.alias = aes128 # Имя ключа для шифрования
interceptor.encryption.json.storage.location = /valid-password-and-keystore.json # Путь до json-файла, где лежат ключи
interceptor.encryption.key.store.type = json # Тип хранилища (json/vault)
interceptor.encryption.mode.operation = encrypt # Режим interceptor (encrypt/signature)
interceptor.encryption.generate.key.password = true # Включает преобразование паролей в ключи
Пример конфигурации Kafka Consumer#
bootstrap.servers = localhost:9092
security.protocol = PLAINTEXT
group.id = test-group

key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

interceptor.classes = ru.sbt.ss.kafka.encryption.interceptor.EncryptionConsumerInterceptor

interceptor.encryption.json.storage.location = /valid-password-and-keystore.json # Путь до json-файла, где лежат ключи
interceptor.encryption.key.store.type = json # Тип хранилища (json/vault)
interceptor.encryption.mode.operation = signature # Режим interceptor (encrypt/signature)
interceptor.encryption.generate.key.password = true # Включает преобразование паролей в ключи

Vault#

Пример конфигурации Kafka Producer#
ssl.vault.address = https://localhost:8200 # Адрес подключения к Vault
ssl.vault.auth.type = APPROLE # Тип алгоритма аутентификации в Vault
ssl.vault.auth.role.id = role-id # Идентификатор роли приложения при ssl.vault.auth.type=approle
ssl.vault.auth.secret.id = secret-id # Секрет роли приложения при ssl.vault.auth.type=approle

ssl.vault.tls.enable = true # Включение протокола TLS при подключении к Vault
ssl.vault.tls.keystore.location = /vault.jks # Путь до keystore
ssl.vault.tls.keystore.password = password # Пароль от сертификата
ssl.vault.tls.key.password = password  # Пароль от ключа
ssl.vault.tls.truststore.location = /vault.jks # Путь до truststore
ssl.vault.tls.truststore.password = password # Пароль от truststore

ssl.vault.pki.mode = kv # Режим выпуска сертификатов
ssl.vault.secret.path = /keys # Путь до ключей в vault хранилище
ssl.vault.disable.pem.certificate.generation = true # Отключает генерацию pem сертификата

bootstrap.servers = localhost:9092
security.protocol = PLAINTEXT
group.id = test-group

key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer

interceptor.classes = ru.sbt.ss.kafka.encryption.interceptor.EncryptionProducerInterceptor

interceptor.encryption.algorithm = AES_128/GCM/NoPadding # Алгоритм шифрования сообщения
interceptor.encryption.key.alias = aes128 # Имя ключа для шифрования
interceptor.encryption.key.store.type = vault # Тип хранилища (json/vault)
interceptor.encryption.mode.operation = signature # Режим interceptor (encrypt/signature)
interceptor.encryption.generate.key.password = true # Включает преобразование паролей в ключи
Пример конфигурации Kafka Consumer#
ssl.vault.address = https://localhost:8200 # Адрес подключения к Vault
ssl.vault.auth.type = APPROLE # Тип алгоритма аутентификации в Vault
ssl.vault.auth.role.id = role-id # Идентификатор роли приложения при ssl.vault.auth.type=approle
ssl.vault.auth.secret.id = secret-id # Секрет роли приложения при ssl.vault.auth.type=approle

ssl.vault.tls.enable = true # Включение протокола TLS при подключении к Vault
ssl.vault.tls.keystore.location = /vault.jks # Путь до keystore
ssl.vault.tls.keystore.password = password # Пароль от сертификата
ssl.vault.tls.key.password = password  # Пароль от ключа
ssl.vault.tls.truststore.location = /vault.jks # Путь до truststore
ssl.vault.tls.truststore.password = password # Пароль от truststore

ssl.vault.pki.mode = kv # Режим выпуска сертификатов
ssl.vault.secret.path = /keys # Путь до ключей в vault хранилище
ssl.vault.disable.pem.certificate.generation = true # Отключает генерацию pem сертификата

bootstrap.servers = localhost:9092
security.protocol = PLAINTEXT
group.id = test-group

key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

interceptor.classes = ru.sbt.ss.kafka.encryption.interceptor.EncryptionConsumerInterceptor

interceptor.encryption.key.store.type = vault # Тип хранилища (json/vault)
interceptor.encryption.mode.operation = signature # Режим interceptor (encrypt/signature)
interceptor.encryption.generate.key.password = true # Включает преобразование паролей в ключи

Результат#

Выполнено подключение интерсептора Kafka-encryption-interceptor.