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: <>"
Таблица параметров#
Ключ |
Дефолтное значение |
Описание параметра |
|---|---|---|
|
|
Режим обработки ошибок ( |
|
|
Длина тега |
|
для режима подписи: |
Алгоритм шифрования сообщения |
|
- |
Имя ключа для подписания/шифрования сообщений. Параметр является обязательным для продюсера, для консьюмера – необязательным |
|
- |
Тип хранилища ключей ( |
|
- |
Путь до json-файла с ключами |
|
|
Частота обновление ключей в локальном хранилище |
|
|
Время повторного обновления ключей при запросе ключа отсутствующего в хранилище |
|
|
Включает создание ключей на основе паролей |
|
- |
Режим работы перехватчика ( |
Параметр interceptor.encryption.performance.logging является устаревшим и более не поддерживается.
Поведение при ошибках#
Поведение при ошибках валидации настраивается с помощью параметра interceptor.encryption.mode:
В случае, если сообщение не прошло валидацию при отправке producer.send() возможны следующие режимы:
failOnValue(используется по умолчанию) – клиент получит сообщение-заглушку вместо невалидного сообщения, которое содержитnullвместоvalueи выбросит исключениеru.sbt.ss.kafka.interceptors.ProducerInterceptorExceptionпри вызовеrecord.value(). Данный способ позволяет использовать логику обработки ошибок kafka-клиента, в том числе вызывать методы перехватчиков (но не callback).failOnSend– будет выброшено исключениеru.sbt.ss.kafka.interceptors.ProducerInterceptorError (extends Throwable)с сообщением, игнорируя логику обработки ошибок kafka-клиента.addErrorHeader– в сообщение будет добавлен заголовок с сообщением об ошибке валидации (по умолчанию__interceptor.errorнастраивается с помощью параметраinterceptor.encryption.error.header.name).
Сообщения исключений имеют формат Error while processing record(topic: topic): ошибка валидации в зависимости от формата сообщения.
В случае, если сообщение не прошло валидацию при получении consumer.poll() поведение настраивается параметром interceptor.encryption.mode:
failOnValue(используется по умолчанию) – клиент получит сообщение-заглушку вместо невалидного сообщения, которое содержитnullвместоvalueи выбросит исключениеru.sbt.ss.kafka.interceptors.ConsumerInterceptorExceptionпри вызовеrecord.value().filter– сообщение об ошибке будет залогировано в error, клиент не получит невалидное сообщение.failOnConsume– методconsumer.poll()выбросит исключениеru.sbt.ss.kafka.interceptors.ConsumerInterceptorError, клиент не получит ни одного сообщения из пачки.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.