Kafka-certificate-signature-interceptor#

Реализует интерфейсы ProducerInterceptor и ConsumerInterceptor, позволяет подписывать сообщения (record.value()) с помощью X509 сертификата.

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

Интерсептор работает с типами String или Array<Byte> и в двух режимах: public-key и certificate-serial.

В случае, если тип сообщения отличается от String или Array<Byte>, можно использовать Kafka-certificate-signature-serde – аналогичная функциональность в виде де/сериализаторов.

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

public-key#

  1. Producer-signature-plugin (interceptor) получает приватный ключ и сертификат из vault или локального JKS.

  2. Producer-signature-plugin (interceptor) подписывает value сообщения и добавляет подпись (signature), алгоритм подписи (algorithm) и публичный ключ (public key) сертификата в заголовки сообщения.

  3. Consumer-signature-plugin (interceptor) проверяет подпись сообщения (signature), используя алгоритм (algorithm) и публичный ключ (public key) из заголовков сообщения.

certificate-serial#

  1. Producer-signature-plugin (interceptor) получает приватный ключ и сертификат из vault или локального JKS.

  2. Producer-signature-plugin (interceptor) подписывает value сообщения и добавляет подпись (signature), алгоритм подписи (algorithm) и серийный номер (certificate serial) сертификата в заголовки сообщения.

  3. Consumer-signature-plugin (interceptor) получает публичный ключ сертификата по его серийному номеру, указанному в заголовке сообщения из локального хранилища (JKS или текстового файла) или vault. При запросе сертификата из vault серийный номер дополнительно проверяется по настраиваемому списку разрешенных серийных номеров.

  4. Consumer-signature-plugin (interceptor) проверяет подпись сообщения (signature), используя алгоритм (algorithm) и полученный публичный ключ.

ВАЖНО! При использовании kafka-certificate-signature-interceptor возможно значительное падение производительности до 150 раз в связи с высокой математической сложностью алгоритма подписи сертификата х509. Для установления конкретных показателей для определенного взаимодействия необходимо проведение НТ.

Подключение к Kafka-клиентам#

  1. Добавить актуальную версию интерсептора в зависимости проекта.

  2. Сконфигурировать kafka-клиенты в соответствии с примерами.

Использование в cloud-среде#

  1. Добавить зависимость интерсептора в сборку образа приложения

  2. Сконфигурировать kafka-клиенты в соответствии с разделом Развертывание в облачной среде.

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

Пример конфигурации в режиме public-key#

Пример конфигурации Kafka Producer#
# 1) Подключить interceptor
interceptor.classes = ru.sbt.ss.kafka.interceptors.CertificateSignatureProducerInterceptor

# 2) Настроить режим работы interceptor
interceptor.signature.certificate.mode = public-key

# 3) Настроить алгоритм подписи
# Алгоритм должен поддерживаться используемым java security провайдером.
# Алгоритм подписи (в данном случае RSA) должен совпадать с алгоритмом приватного ключа.
interceptor.signature.certificate.algorithm = SHA256WithRSA

# 4) Настроить хранилище сертификатов в зависимости от типа (jks, pem или vault):
# ДОБАВИТЬ НАСТРОЙКИ ХРАНИЛИЩА СЕРТИФИКАТОВ ИЗ "Примеры конфигурации хранилища сертификатов для kafka-producer"
interceptor.signature.certificate.keystore.type = change-me

# 5) ОПЦИОНАЛЬНО Настроить имена системных заголовков
# Имя заголовка с подписью сообщения
# interceptor.signature.header = signature

# Префикс для имен системных заголовков
# interceptor.signature.attributes.prefix = signature.

# Имя загловка с используемым алгоритмом подписи
# interceptor.signature.certificate.algorithm.attribute=algorithm

# Имя заголовка с публичным ключом для проверки подписи
# interceptor.signature.certificate.public.key.attribute=public.key

# 6) ОПЦИОНАЛЬНО
# Включить логирование interceptor при подписи сообщений
# interceptor.trace.logger.enabled = true

# Задать имя логгера.
# По умолчанию генерируется автоматически в виде <Имя_класса_интерсептора>[<client.id_из_настроек_kafka_клиента>], например:
# ru.sbt.ss.kafka.interceptors.CertificateSignatureProducerInterceptor[Producer-1]
# Если client.id не задан - он генерируется автоматически
# interceptor.trace.logger.name = LoggerName

# 7) ОПЦИОНАЛЬНО
# Включить публикацию jmx-метрик количества успешно и неудачно обработанных сообщений
# interceptor.jmx.metrics.enabled = true

# 8) ОПЦИОНАЛЬНО
# Включить публикацию телеметрии
# interceptor.telemetry.enabled = true
# Задать путь к файлу конфигурации телеметрии
# interceptor.telemetry.config.path = /path/to/config/file
# Включить отображение спанов подписи сообщений (по-умолчанию false)
# interceptor.telemetry.callback.enabled = true
Пример конфигурации Kafka Consumer#
# 1) Подключить interceptor
interceptor.classes = ru.sbt.ss.kafka.interceptors.CertificateSignatureConsumerInterceptor

# 2) Настроить режим работы interceptor
interceptor.signature.certificate.mode = public-key

# 3) ОПЦИОНАЛЬНО Настроить удаление системных заголовков из сообщения
# interceptor.signature.remove.headers = false

# 4) ОПЦИОНАЛЬНО Настроить режим работы консьюмера при ошибках (по умолчанию failOnValue)
# interceptor.signature.mode = failOnValue

# 5) ОПЦИОНАЛЬНО Настроить имена системных заголовков

# Имя заголовка с подписью сообщения
# interceptor.signature.header = signature

# Префикс для имен системных заголовков
# interceptor.signature.attributes.prefix = signature.

# Имя заголовка с используемым алгоритмом подписи
# interceptor.signature.certificate.algorithm.attribute = algorithm

# Имя заголовка с публичным ключом для проверки подписи
# interceptor.signature.certificate.public.key.attribute = public.key

# 6) ОПЦИОНАЛЬНО
# Включить логирование interceptor при подписи сообщений
# interceptor.trace.logger.enabled = true

# Задать имя логгера.
# По умолчанию генерируется автоматически в виде <Имя_класса_интерсептора>[<client.id_из_настроек_kafka_клиента>], например:
# ru.sbt.ss.kafka.interceptors.CertificateSignatureConsumerInterceptor[Consumer-1]
# Если client.id не задан - он генерируется автоматически
# interceptor.trace.logger.name = LoggerName

# Включить публикацию jmx-метрик количества успешно и неудачно обработанных сообщений
# interceptor.jmx.metrics.enabled = true

# 7) ОПЦИОНАЛЬНО

# Включить публикацию телеметрии
# interceptor.telemetry.enabled = true

# Задать путь к файлу конфигурации телеметрии
# interceptor.telemetry.config.path = /path/to/config/file

# Включить отображение спанов подписи сообщений (по-умолчанию false)
# interceptor.telemetry.callback.enabled = true

Пример конфигурации в режиме certificate-serial#

Пример конфигурации Kafka Producer#
# 1) Подключить interceptor
interceptor.classes = ru.sbt.ss.kafka.interceptors.CertificateSignatureProducerInterceptor

# 2) Настроить режим работы interceptor
interceptor.signature.certificate.mode = certificate-serial

# 3) Настроить алгоритм подписи
# Алгоритм должен поддерживаться используемым java security провайдером.
# Алгоритм подписи (в данном случае RSA) должен совпадать с алгоритмом приватного ключа.
interceptor.signature.certificate.algorithm = SHA256WithRSA

# 4) Настроить хранилище сертификатов в зависимости от типа (jks, pem или vault):
# ДОБАВИТЬ НАСТРОЙКИ ХРАНИЛИЩА СЕРТИФИКАТОВ ИЗ "Примеры конфигурации хранилища сертификатов для kafka-producer"
interceptor.signature.certificate.keystore.type = change-me

# 5) ОПЦИОНАЛЬНО Настроить имена системных заголовков
# Имя заголовка с подписью сообщения
# interceptor.signature.header = signature

# Префикс для имен системных заголовков
# interceptor.signature.attributes.prefix = signature.

# Имя загловка с используемым алгоритмом подписи
# interceptor.signature.certificate.algorithm.attribute=algorithm

# Имя заголовка с серийным номером сертификата для проверки подписи
# interceptor.signature.certificate.serial.attribute=certificate.serial

# 6) ОПЦИОНАЛЬНО
# Включить логирование интерсептора при подписи сообщений
# interceptor.trace.logger.enabled = true

# Задать имя логгера.
# По умолчанию генерируется автоматически в виде <Имя_класса_интерсептора>[<client.id_из_настроек_kafka_клиента>], например:
# ru.sbt.ss.kafka.interceptors.CertificateSignatureProducerInterceptor[Producer-1]
# Если client.id не задан - он генерируется автоматически
# interceptor.trace.logger.name = LoggerName

# Включить публикацию jmx-метрик количества успешно и неудачно обработанных сообщений
# interceptor.jmx.metrics.enabled = true
Пример конфигурации Kafka Consumer#
# 1) Подключить interceptor
interceptor.classes = ru.sbt.ss.kafka.interceptors.CertificateSignatureConsumerInterceptor

# 2) Настроить режим работы interceptor
interceptor.signature.certificate.mode = certificate-serial

# 3) ОПЦИОНАЛЬНО Настроить удаление системных заголовков из сообщения
# interceptor.signature.remove.headers = false

# 4) ОПЦИОНАЛЬНО Настроить режим работы консьюмера при ошибках (по умолчанию failOnValue)
# interceptor.signature.mode = failOnValue

# 5) Настроить хранилище сертификатов в зависимости от типа (jks, pem, file или vault):
# ДОБАВИТЬ НАСТРОЙКИ ХРАНИЛИЩА СЕРТИФИКАТОВ ИЗ "Примеры конфигурации хранилища сертификатов для kafka-consumer"
interceptor.signature.certificate.truststore.type = change-me

# 6) ОПЦИОНАЛЬНО Настроить имена системных заголовков

# Имя заголовка с подписью сообщения
# interceptor.signature.header = signature

# Префикс для имен системных заголовков
# interceptor.signature.attributes.prefix = signature.

# Имя загловка с используемым алгоритмом подписи
# interceptor.signature.certificate.algorithm.attribute=algorithm

# Имя заголовка с серийным номером сертификата для проверки подписи
# interceptor.signature.certificate.serial.attribute=certificate.serial

# 7) ОПЦИОНАЛЬНО
# Включить логирование интерсептора при подписи сообщений
# interceptor.trace.logger.enabled = true

# Задать имя логгера.
# По умолчанию генерируется автоматически в виде <Имя_класса_интерсептора>[<client.id_из_настроек_kafka_клиента>], например:
# ru.sbt.ss.kafka.interceptors.CertificateSignatureConsumerInterceptor[Consumer-1]
# Если client.id не задан - он генерируется автоматически
# interceptor.trace.logger.name = LoggerName

# Включить публикацию jmx-метрик количества успешно и неудачно обработанных сообщений
# interceptor.jmx.metrics.enabled = true

Примеры конфигурации хранилища сертификатов для Kafka Producer#

Пример конфигурации хранилища типа jks для Kafka Producer#
# 4.1) Настроить хранилище сертификатов типа jks
interceptor.signature.certificate.keystore.type = jks
# По умолчанию используется хранилище сертификатов из настроек kafka-клиента:
ssl.keystore.location = ssl/keystore.jks
ssl.keystore.password = password
ssl.key.password = password

# Перехватчик может использовать отдельное хранилище сертификатов:
# interceptor.signature.certificate.ssl.keystore.location = ssl/keystore.jks
# interceptor.signature.certificate.ssl.keystore.password = password
# interceptor.signature.certificate.ssl.key.password = password

# Выбрать сертификат из хранилища, указав его alias (по умолчанию используется первый)
# interceptor.signature.certificate.key.alias =
Пример конфигурации хранилища типа pem для Kafka Producer#
# 4.1) Настроить хранилище сертификатов типа pem
# interceptor.signature.certificate.keystore.type = pem
# interceptor.signature.certificate.keystore.crt.location = certificate.pem
# interceptor.signature.certificate.keystore.private.key.location= privateKey.key
Пример конфигурации хранилища типа vault для Kafka Producer#
# 4.1) Настроить подключение к vault
interceptor.signature.certificate.keystore.type = vault

# ВСЕ ПАРАМЕТРЫ НИЖЕ ЯВЛЯЮТСЯ ПАРАМЕТРАМИ ПЛАГИНА ssl-context-builder С ПРЕФИКСОМ `interceptor.signature.certificate.`
# Полный список возможных параметров есть в документации плагина ssl-context-builder

# Адрес vault
interceptor.signature.certificate.ssl.vault.address = https://host:port

# ОПЦИОНАЛЬНО Namespace vault
# interceptor.signature.certificate.ssl.vault.namespace = namespace

# ОПЦИОНАЛЬНО Настройки повторной отправки запросов к vault
# Кол-во попыток переотправки запроса
# interceptor.signature.certificate.ssl.vault.retries = 5

# Тайм-аут отправки запроса, в секундах
# interceptor.signature.certificate.ssl.vault.timeout = 3

# Интервал между повторными попытками переотправки запроса, мс
# interceptor.signature.certificate.ssl.vault.retry.interval = 500

# Настройки ssl для подключения к vault
interceptor.signature.certificate.ssl.vault.tls.enable = true
interceptor.signature.certificate.ssl.vault.tls.keystore.location = vault-keystore.jks
interceptor.signature.certificate.ssl.vault.tls.keystore.password = password
interceptor.signature.certificate.ssl.vault.tls.key.password = password
interceptor.signature.certificate.ssl.vault.tls.truststore.location = vault-keystore.jks
interceptor.signature.certificate.ssl.vault.tls.truststore.password = password
interceptor.signature.certificate.ssl.endpoint.identification.algorithm =

# Настройки авторизации vault
# Пример авторизации с помощью логина и пароля
# Тип авторизации (approle, certificate, password или token)
interceptor.signature.certificate.ssl.vault.auth.type = password
interceptor.signature.certificate.ssl.vault.auth.username = test
interceptor.signature.certificate.ssl.vault.auth.password = password

# Пример авторизации с помощью approle
# interceptor.signature.certificate.ssl.vault.auth.type = approle
# interceptor.signature.certificate.ssl.vault.auth.role.id = role
# interceptor.signature.certificate.ssl.vault.auth.secret.id = secret

# 4.2 Настроить получение сертификата из vault
# Получить сертификат можно двумя способами:
# a) Сгенерировать новый сертификат с помощью pki engine
# b) Получить заранее загруженный сертификат в формате pem с помощью kv engine

# 4.2.a) Настроить генерацию сертификата с помощью vault (pki engine)
# interceptor.signature.certificate.ssl.vault.pki.mode = pki
#
# ОПЦИОНАЛЬНО Метод получения сертификата:
# * issue (по умолчанию) - выпускает новый сертификат при каждом обращении
# * fetch (движок SberCA) - выпускает новый сертификат только если сертификат с указанными параметрами не существует или истек
# При использовании движка SberCA параметры `ssl.vault.pki.ttl`, `ssl.vault.pki.csr`, `ssl.vault.pki.csr.path` и `ssl.vault.pki.not.after` не поддерживаются и будут проигнорированы
# interceptor.signature.certificate.ssl.vault.pki.method = issue
#
# ОПЦИОНАЛЬНО Путь до pki engine
# interceptor.signature.certificate.ssl.vault.pki.mount = pki
#
# Имя роли для выпуска сертификата
interceptor.signature.certificate.ssl.vault.pki.role.name = role
#
# Common name сертификата (CN)
interceptor.signature.certificate.ssl.vault.pki.common.name = INTERCEPTOR-TEST
#
# ОПЦИОНАЛЬНО Электронный адрес владельца сертификата, задается при использовании Secret Manager
# interceptor.signature.certificate.ssl.vault.pki.email = email@example.com
#
# ОПЦИОНАЛЬНО Alternative names сертификата
# interceptor.signature.certificate.ssl.vault.pki.alt.names = alt-name
#
# ОПЦИОНАЛЬНО Alternative ip сертификата
# interceptor.signature.certificate.ssl.vault.pki.alt.ip = 127.0.0.1
#
# ОПЦИОНАЛЬНО TTL (time-to-live) сертификата
# interceptor.signature.certificate.ssl.vault.pki.ttl =
#
# ОПЦИОНАЛЬНО Запрос на создание сертификата csr (Certificate signing request)
# interceptor.signature.certificate.ssl.vault.pki.csr =
#
# ОПЦИОНАЛЬНО Путь до файла с запросом на создание сертификата csr (Certificate signing request)
# interceptor.signature.certificate.ssl.vault.pki.csr.path =

# 4.2.b) Настроить получение заранее загруженного сертификата с помощью kv engine
# interceptor.signature.certificate.ssl.vault.pki.mode = kv
#
# Указать путь до key-value хранилища с приватным ключем и сертификатом в формате pem
# interceptor.signature.certificate.ssl.vault.pem.path  = kv1/certificate
#
# Указать имя секрета в key-value хранилище, который содержит сертификат в формате pem
# interceptor.signature.certificate.ssl.vault.pem.name = cert
#
# Указать имя секрета в key-value хранилище, который содержит приватный ключ в формате pem
# interceptor.signature.certificate.ssl.vault.pem.key = key

# 4.3) ОПЦИОНАЛЬНО Настроить хранилище сертификатов (локальный кэш)
# Пути до хранилищ сертификатов, сгенерированных vault
# Если не указаны - сертификаты будут генерироваться каждый раз при старте приложения
# interceptor.signature.certificate.ssl.keystore.location = producer-vault-keystore.jks
# interceptor.signature.certificate.ssl.truststore.location = producer-vault-truststore.jks

# alias клиентского сертификата в keystore
# interceptor.signature.certificate.ssl.vault.alias.key = key
# alias ca сертификата в truststore
# interceptor.signature.certificate.ssl.vault.alias.ca = ca

# 4.4) ОПЦИОНАЛЬНО Настроить получение паролей для хранилища сертификатов из vault
# Путь до хранилища секретов в vault
# interceptor.signature.certificate.ssl.vault.secret.path = kv1/interceptor
# Версия secret engine
# interceptor.signature.certificate.ssl.vault.engine.version = 1

# Имя секрета, содержащего пароль для private key
# interceptor.signature.certificate.ssl.vault.secret.key = key
# Имя секрета, содержащего пароль для keystore
# interceptor.signature.certificate.ssl.vault.secret.keystore = keystore
# Имя секрета, содержащего пароль для truststore
# interceptor.signature.certificate.ssl.vault.secret.truststore = truststore

Примеры конфигурации хранилища сертификатов для Kafka Consumer#

Пример конфигурации хранилища типа jks для Kafka Consumer#
interceptor.signature.certificate.truststore.type = jks

# По умолчанию используется хранилище сертификатов из настроек kafka-клиента:
ssl.truststore.location = ssl/truststore.jks
ssl.truststore.password = password

# Перехватчик может использовать отдельное хранилище сертификатов:
interceptor.signature.certificate.ssl.truststore.location = truststore.jks
interceptor.signature.certificate.ssl.truststore.password = password
Пример конфигурации хранилища типа pem для Kafka Consumer#
interceptor.signature.certificate.truststore.type = pem

# Указать сертификат или директорию с сертификатами
interceptor.signature.truststore.crt.location = certificate.pem

# ОПЦИОНАЛЬНО Указать расширение сертификатов для фильтрации файлов в директории. По умолчанию ".pem"
# interceptor.signature.truststore.crt.location.filter = .cer
Пример конфигурации хранилища типа file для Kafka Consumer#
interceptor.signature.certificate.truststore.type = file
# Файл с публичными ключами в формате "серийный номер сертификата" = "публичный ключ в base64"
interceptor.signature.certificate.truststore.file = truststore.txt
Пример конфигурации хранилища типа vault для Kafka Consumer#
# 4.1) Настроить подключение к vault
interceptor.signature.certificate.keystore.type = vault

# ВСЕ ПАРАМЕТРЫ НИЖЕ ЯВЛЯЮТСЯ ПАРАМЕТРАМИ ПЛАГИНА ssl-context-builder С ПРЕФИКСОМ `interceptor.signature.certificate.`
# Полный список возможных параметров есть в документации плагина ssl-context-builder

# Адрес vault
interceptor.signature.certificate.ssl.vault.address = https://host:port

# ОПЦИОНАЛЬНО Namespace vault
# interceptor.signature.certificate.ssl.vault.namespace = namespace

# ОПЦИОНАЛЬНО Настройки повторной отправки запросов к vault
# Кол-во попыток переотправки запроса
# interceptor.signature.certificate.ssl.vault.retries = 5

# Тайм-аут отправки запроса, с
# interceptor.signature.certificate.ssl.vault.timeout = 3

# Интервал между повторными попытками переотправки запроса, мс
# interceptor.signature.certificate.ssl.vault.retry.interval = 500

# Настройки ssl для подключения к vault
interceptor.signature.certificate.ssl.vault.tls.enable = true
interceptor.signature.certificate.ssl.vault.tls.keystore.location = vault-keystore.jks
interceptor.signature.certificate.ssl.vault.tls.keystore.password = password
interceptor.signature.certificate.ssl.vault.tls.key.password = password
interceptor.signature.certificate.ssl.vault.tls.truststore.location = vault-keystore.jks
interceptor.signature.certificate.ssl.vault.tls.truststore.password = password
interceptor.signature.certificate.ssl.endpoint.identification.algorithm =

# Настройки авторизации vault
# Пример авторизации с помощью логина и пароля
# Тип авторизации (approle, certificate, password или token)
interceptor.signature.certificate.ssl.vault.auth.type = password
interceptor.signature.certificate.ssl.vault.auth.username = test
interceptor.signature.certificate.ssl.vault.auth.password = password

# Пример авторизации с помощью approle
# interceptor.signature.certificate.ssl.vault.auth.type = approle
# interceptor.signature.certificate.ssl.vault.auth.role.id = role
# interceptor.signature.certificate.ssl.vault.auth.secret.id = secret

# 4.2) ОПЦИОНАЛЬНО Настроить хранилище сертификатов (локальный кэш)
# Путь до хранилища сертификатов, полученных из vault
# interceptor.signature.certificate.ssl.truststore.location = consumer-vault-truststore.jks

# 4.3) ОПЦИОНАЛЬНО Настроить получение паролей для хранилища сертификатов из vault
# Путь до хранилища секретов в vault
# interceptor.signature.certificate.ssl.vault.secret.path = kv1/interceptor
# Версия secret engine
# interceptor.signature.certificate.ssl.vault.engine.version = 1

# Имя секрета, содержащего пароль для truststore
# interceptor.signature.certificate.ssl.vault.secret.truststore = truststore

# 4.4) ОПЦИОНАЛЬНО Настроить список доверенных сертификатов
# По умолчанию все сертификаты являются доверенными
# Если сертификат не входит в этот список - будет получена ошибка проверки подписи.
#
# Можно использовать один из трех вариантов проверки сертификата:
# 4.4.1) Проверка по серийному номеру сертификата
# Атрибут сертификата, который будет проверяться по списку доверенных сертификатов (по умолчанию serial):
# interceptor.signature.certificate.allowed.certificate.list.type = serial
#
# Список серийных номеров строкой через запятую:
# interceptor.signature.certificate.allowed.certificate.list = 3ca7322c, serial
# Путь до файла, в файле каждая строка является отдельным серийным номером:
# interceptor.signature.certificate.allowed.certificate.list.file = allowed-serials.txt
#
# DEPRECATED Также для типа serial существуют устаревшие настройки, аналогичные настройкам выше:
# Настройки allowed.certificate.list.* имеют наивысший приоритет
# interceptor.signature.certificate.allowed.serial = 3ca7322c, serial
# interceptor.signature.certificate.allowed.serial.file = allowed-serials.txt
#
# 4.4.2) Проверка по Common Name (CN) сертификата
# Атрибут сертификата, который будет проверяться по списку доверенных сертификатов:
# interceptor.signature.certificate.allowed.certificate.list.type = cn
#
# Список CN сертификатов строкой через запятую
# interceptor.signature.certificate.allowed.certificate.list = CN=test, CN=test2
# Путь до файла, в файле каждая строка является отдельным CN сертификата:
# interceptor.signature.certificate.allowed.certificate.list.file = allowed-cn.txt
#
# 4.4.3) Проверка по Distinguished Name (DN) сертификата
# Атрибут сертификата, который будет проверяться по списку доверенных сертификатов:
# interceptor.signature.certificate.allowed.certificate.list.type = dn
#
# Список DN сертификатов строкой через точку с запятой
# interceptor.signature.certificate.allowed.certificate.list = CN=test, OU=FPSS, O=SBT, ST=Moscow, C=RU; CN=test2, C=RU
# Путь до файла, в файле каждая строка является отдельным CN сертификата:
# interceptor.signature.certificate.allowed.certificate.list.file = allowed-dn.txt

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

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

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

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

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

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

В случае, если сообщение не прошло проверку подписи при получении (consumer.poll()), поведение настраивается параметром interceptor.signature.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, клиент не получит ни одного сообщения из пачки.

Загрузка хранилища сертификатов из classpath#

Для загрузки хранилища сертификатов из classpath необходимо указать протокол classpath:// в пути до файла, например:

interceptor.signature.certificate.ssl.keystore.location = classpath://ssl/keystore.jks

Не работает для стандартных ssl-настроек kafka-client типа ssl.keystore.location.

Не работает при использовании pem сертификатов.

Использование Conscrypt#

Conscrypt – open-source java security provider, его реализация алгоритмов подписи RSA в тестах показывает двухкратный прирост производительности по сравнению со стандартным провайдером JVM.

Подключение к Kafka Producer#

Библиотека Conscrypt уже включена в транзитивные зависимости.

Есть два способа подключить Conscrypt:

  1. С помощью параметра conscrypt.enabled (по умолчанию данный параметр уже активирован):

interceptor.signature.certificate.conscrypt.enabled = true
  1. (Устаревший способ) Инициализировать провайдер перед запуском kafka-producer:

import org.conscrypt.Conscrypt;
Security.insertProviderAt(Conscrypt.newProvider(), Security.getProviders.length);

2.1 Указать провайдер в конфигурации producer:

# Имя java security provider'а, используемого для подписи сообщений
interceptor.signature.certificate.ssl.provider = Conscrypt

Использование Conscrypt для проверки подписи на консьюмер не дает значительного прироста производительности

JMX метрики#

При включении публикации JMX-метрик с помощью настройки interceptor.jmx.metrics.enabled = true в JMX будут добавлены метрики с количеством успешно и ошибочно обработанных сообщений и информацией о подключенном перехватчике:

kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=OttSignatureInterceptor,name=FailedProcessedMessage
kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=OttSignatureInterceptor,name=SuccessfulProcessedMessage
kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=OttSignatureInterceptor,name=Info

Или:

kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=OttSignatureInterceptor,name=FailedProcessedMessage
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=OttSignatureInterceptor,name=SuccessfulProcessedMessage
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=OttSignatureInterceptor,name=Info

, где client-id берется из конфигурации client.id, при отсутствии в конфигурации генерируется автоматически.

Результат#

Выполнено подключение перехватчика Kafka-certificate-signature-interceptor.