Kora облачно ориентированный серверный фреймворк написанный на Java для написания Java / Kotlin приложений с упором на производительность, эффективность, прозрачность сделанный выходцами из Т-Банк / Тинькофф

Kora is a cloud-oriented server-side Java framework for writing Java / Kotlin applications with a focus on performance, efficiency and transparency

Перейти к содержанию
V1 V2

gRPC сервер

Модуль запускает gRPC-сервер на основе grpc-java и подключает к нему обработчики из графа приложения. Обработчик — это BindableService, обычно класс, который наследует сгенерированный ...ImplBase и реализует унарные или потоковые методы RPC.

Kora строит сервер на транспорте gRPC OkHttp, добавляет сервисы и реализации ServerInterceptor из графа вместе со своим перехватчиком телеметрии, управляет жизненным циклом сервера и участвует в проверках готовности приложения. Если параметров конфигурации недостаточно, итоговый билдер можно дополнительно настроить в коде через компонент Configurer.

Если нужен пошаговый разбор перед справочным описанием, смотрите gRPC-сервер и продвинутый gRPC-сервер.

Подключение

Зависимость build.gradle:

implementation "io.koraframework:grpc-server"
implementation "io.grpc:grpc-protobuf:1.83.1"
implementation "javax.annotation:javax.annotation-api:1.3.2"

Модуль:

@KoraApp
public interface Application extends GrpcServerModule { }

Зависимость build.gradle.kts:

implementation("io.koraframework:grpc-server")
implementation("io.grpc:grpc-protobuf:1.83.1")
implementation("javax.annotation:javax.annotation-api:1.3.2")

Модуль:

@KoraApp
interface Application : GrpcServerModule

Вместе с io.koraframework:grpc-server приходит рантайм gRPC версии 1.83.1. Все остальные артефакты io.grpcgrpc-protobuf, grpc-services и всё, что подключается в тестовой области, — должны быть той же версии, смотрите Тестирование.

Плагин

Код для gRPC-сервера генерируется с помощью gradle-плагина protobuf.

Плагин в build.gradle:

plugins {
    id "com.google.protobuf" version "0.10.0"
}

protobuf {
    protoc { artifact = "com.google.protobuf:protoc:4.35.1" }
    plugins {
        grpc { artifact = "io.grpc:protoc-gen-grpc-java:1.83.1" }
    }
    generateProtoTasks {
        all()*.plugins { grpc {} }
    }
}

sourceSets {
    main {
        java {
            srcDirs "build/generated/source/proto/main/grpc"
            srcDirs "build/generated/source/proto/main/java"
        }
    }
}

Плагин в build.gradle.kts:

import com.google.protobuf.gradle.id

plugins {
    id("com.google.protobuf") version ("0.10.0")
}

protobuf {
    protoc { artifact = "com.google.protobuf:protoc:4.35.1" }
    plugins {
        id("grpc") { artifact = "io.grpc:protoc-gen-grpc-java:1.83.1" }
    }
    generateProtoTasks {
        all().forEach { task -> task.plugins { id("grpc") } }
    }
}

sourceSets.main {
    java.srcDir(layout.buildDirectory.dir("generated/source/proto/main/grpc"))
    java.srcDir(layout.buildDirectory.dir("generated/source/proto/main/java"))
}

Плагин генерирует классы на Java, поэтому в проекте на Kotlin сгенерированные исходники всё равно подключаются к набору исходников java.

Конфигурация

Обычно нужно задать только port; все остальные параметры имеют значения по умолчанию. Минимальная конфигурация, которая привязывает порт и включает логирование:

grpcServer {
    port = 8090
    telemetry.logging.enabled = true
}
grpcServer:
  port: 8090
  telemetry:
    logging:
      enabled: true

Основные параметры конфигурации:

grpcServer {
    port = 8090 //(1)!
    maxMessageSize = "4MiB" //(2)!
    reflectionEnabled = false //(3)!
}
  1. Порт gRPC-сервера (по умолчанию: 8090).
  2. Максимальный размер входящего сообщения (по умолчанию: 4MiB).
  3. Включает сервис gRPC Server Reflection (по умолчанию: false).
grpcServer:
  port: 8090 #(1)!
  maxMessageSize: "4MiB" #(2)!
  reflectionEnabled: false #(3)!
  1. Порт gRPC-сервера (по умолчанию: 8090).
  2. Максимальный размер входящего сообщения (по умолчанию: 4MiB).
  3. Включает сервис gRPC Server Reflection (по умолчанию: false).
Полная конфигурация

Пример полной конфигурации, описанной в GrpcServerConfig:

grpcServer {
    port = 8090 //(1)!
    maxMessageSize = "4MiB" //(2)!
    reflectionEnabled = false //(3)!
    shutdownWait = "30s" //(4)!
    maxConnectionAge = "5m" //(5)!
    maxConnectionAgeGrace = "30s" //(6)!
    keepAliveTime = "30s" //(7)!
    keepAliveTimeout = "10s" //(8)!
    telemetry {
        logging {
            enabled = false //(9)!
        }
        metrics {
            enabled = false //(10)!
            slo = [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] //(11)!
            tags = { // (12)!
                "key1" = "value1"
                "key2" = "value2"
            }
        }
        tracing {
            enabled = true //(13)!
            attributes = { // (14)!
                "key1" = "value1"
                "key2" = "value2"
            }
        }
    }
}
  1. Порт gRPC-сервера (по умолчанию: 8090).
  2. Максимальный размер входящего сообщения (по умолчанию: 4MiB). Может быть указан в виде числа байт или как 4MiB, 4MB, 1000Kb и подобных значений.
  3. Включает сервис gRPC Server Reflection (по умолчанию: false).
  4. Время ожидания завершения выполняющихся вызовов перед выключением сервера при штатном завершении (по умолчанию: 30s).
  5. Максимальное время жизни соединения, после которого соединение штатно завершается (опционально, без значения по умолчанию). К значению добавляется случайное отклонение +/-10%.
  6. Дополнительное время на штатное завершение соединения после достижения максимального времени жизни (опционально, без значения по умолчанию). Вызовы RPC, которые не успевают завершиться, отменяются, чтобы соединение могло закрыться.
  7. Интервал между кадрами PING (опционально, без значения по умолчанию).
  8. Тайм-аут подтверждения кадра PING (опционально, без значения по умолчанию). Если подтверждение не получено за это время, соединение закрывается.
  9. Включает логирование модуля (по умолчанию: false).
  10. Включает метрики модуля (по умолчанию: false). Метрики пишутся только если модуль метрик также предоставляет MeterRegistry.
  11. Настройка SLO для метрики Timer (по умолчанию: io.koraframework.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO).
  12. Теги метрик (по умолчанию: {}).
  13. Включает трассировку модуля (по умолчанию: true). Спаны экспортируются только если модуль трассировки также предоставляет Tracer.
  14. Атрибуты трассировки (по умолчанию: {}).
grpcServer:
  port: 8090 #(1)!
  maxMessageSize: "4MiB" #(2)!
  reflectionEnabled: false #(3)!
  shutdownWait: "30s" #(4)!
  maxConnectionAge: "5m" #(5)!
  maxConnectionAgeGrace: "30s" #(6)!
  keepAliveTime: "30s" #(7)!
  keepAliveTimeout: "10s" #(8)!
  telemetry:
    logging:
      enabled: false #(9)!
    metrics:
      enabled: false #(10)!
      slo: [ 1, 10, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 30000, 60000, 90000 ] #(11)!
      tags: #(12)!
        key1: value1
        key2: value2
    tracing:
      enabled: true #(13)!
      attributes: #(14)!
        key1: value1
        key2: value2
  1. Порт gRPC-сервера (по умолчанию: 8090).
  2. Максимальный размер входящего сообщения (по умолчанию: 4MiB). Может быть указан в виде числа байт или как 4MiB, 4MB, 1000Kb и подобных значений.
  3. Включает сервис gRPC Server Reflection (по умолчанию: false).
  4. Время ожидания завершения выполняющихся вызовов перед выключением сервера при штатном завершении (по умолчанию: 30s).
  5. Максимальное время жизни соединения, после которого соединение штатно завершается (опционально, без значения по умолчанию). К значению добавляется случайное отклонение +/-10%.
  6. Дополнительное время на штатное завершение соединения после достижения максимального времени жизни (опционально, без значения по умолчанию). Вызовы RPC, которые не успевают завершиться, отменяются, чтобы соединение могло закрыться.
  7. Интервал между кадрами PING (опционально, без значения по умолчанию).
  8. Тайм-аут подтверждения кадра PING (опционально, без значения по умолчанию). Если подтверждение не получено за это время, соединение закрывается.
  9. Включает логирование модуля (по умолчанию: false).
  10. Включает метрики модуля (по умолчанию: false). Метрики пишутся только если модуль метрик также предоставляет MeterRegistry.
  11. Настройка SLO для метрики Timer (по умолчанию: io.koraframework.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO).
  12. Теги метрик (по умолчанию: {}).
  13. Включает трассировку модуля (по умолчанию: true). Спаны экспортируются только если модуль трассировки также предоставляет Tracer.
  14. Атрибуты трассировки (по умолчанию: {}).

Всё, что не покрыто конфигурацией, доступно через настройку в коде.

Настройка в коде

Если параметров конфигурации недостаточно, зарегистрируйте компонент Configurer<ForwardingServerBuilder<?>> и донастройте билдер сервера в коде. Компонент вызывается последним: после применения конфигурации и после того, как добавлены сервисы, пользовательские реализации ServerInterceptor и стандартный перехватчик. Сервер собирается из того билдера, который вернул компонент.

@Component
public final class MyGrpcServerConfigurer implements Configurer<ForwardingServerBuilder<?>> {

    @Override
    public ForwardingServerBuilder<?> configure(ForwardingServerBuilder<?> builder) {
        builder.maxInboundMetadataSize(16 * 1024); //(1)!
        builder.handshakeTimeout(10, TimeUnit.SECONDS);

        if (builder instanceof OkHttpServerBuilder okHttpBuilder) { //(2)!
            okHttpBuilder.permitKeepAliveWithoutCalls(true);
            okHttpBuilder.maxConcurrentCallsPerConnection(200);
        }

        return builder; //(3)!
    }
}
  1. Опции, объявленные в io.grpc.ServerBuilder, доступны прямо на делегирующем билдере
  2. Настройки конкретного транспорта требуют конкретного билдера: Kora поднимает сервер на транспорте gRPC OkHttp, поэтому это io.grpc.okhttp.OkHttpServerBuilder
  3. Билдеры gRPC мутируют себя и возвращают самих себя, поэтому наружу отдаётся тот же экземпляр
@Component
class MyGrpcServerConfigurer : Configurer<ForwardingServerBuilder<*>> {

    override fun configure(builder: ForwardingServerBuilder<*>): ForwardingServerBuilder<*> {
        builder.maxInboundMetadataSize(16 * 1024) //(1)!
        builder.handshakeTimeout(10, TimeUnit.SECONDS)

        if (builder is OkHttpServerBuilder) { //(2)!
            builder.permitKeepAliveWithoutCalls(true)
            builder.maxConcurrentCallsPerConnection(200)
        }

        return builder //(3)!
    }
}
  1. Опции, объявленные в io.grpc.ServerBuilder, доступны прямо на делегирующем билдере
  2. Настройки конкретного транспорта требуют конкретного билдера: Kora поднимает сервер на транспорте gRPC OkHttp, поэтому это io.grpc.okhttp.OkHttpServerBuilder
  3. Билдеры gRPC мутируют себя и возвращают самих себя, поэтому наружу отдаётся тот же экземпляр

Метрики модуля описаны в разделе Справочник метрик.

Защита транспорта

По умолчанию сервер принимает соединения без шифрования: если в графе нет io.grpc.ServerCredentials, Kora использует InsecureServerCredentials. Чтобы терминировать TLS на самом сервере, предоставьте ServerCredentials компонентом — например, через io.grpc.TlsServerCredentials:

@KoraApp
public interface Application extends GrpcServerModule {

    default ServerCredentials grpcServerCredentials() throws IOException { //(1)!
        return TlsServerCredentials.create(
            new File("/etc/certs/server.crt"), //(2)!
            new File("/etc/certs/server.key"));
    }
}
  1. Фабричный метод графа приложения: учётные данные подхватываются при создании билдера gRPC-сервера
  2. Цепочка сертификатов в формате PEM и незашифрованный приватный ключ PKCS#8
@KoraApp
interface Application : GrpcServerModule {

    fun grpcServerCredentials(): ServerCredentials { //(1)!
        return TlsServerCredentials.create(
            File("/etc/certs/server.crt"), //(2)!
            File("/etc/certs/server.key"))
    }
}
  1. Фабричный метод графа приложения: учётные данные подхватываются при создании билдера gRPC-сервера
  2. Цепочка сертификатов в формате PEM и незашифрованный приватный ключ PKCS#8

Для взаимного TLS и собственного хранилища доверенных сертификатов собирайте учётные данные через TlsServerCredentials.newBuilder().

Обработчики

Обработчик — это класс, который наследует сгенерированный ...ImplBase и регистрируется в графе приложения с помощью аннотации @Component. Класс ...ImplBase создается из контракта proto с помощью gradle-плагина protobuf; вы переопределяете его методы RPC, чтобы реализовать поведение сервера. Обычные компоненты Kora, такие как сервисы и репозитории, можно внедрить в обработчик через его конструктор.

Рассмотрим контракт proto с единственным унарным методом:

src/main/proto/message.proto
syntax = "proto3";

package io.koraframework.generated.grpc;

service UserService {
  rpc createUser(RequestEvent) returns (ResponseEvent) {} //(1)!
}

message RequestEvent {
  string name = 1;
  string code = 2;
}

message ResponseEvent {
  bytes id = 1;
}
  1. Унарный RPC: одно сообщение запроса порождает одно сообщение ответа.

Плагин генерирует UserServiceGrpc.UserServiceImplBase, а обработчик переопределяет сгенерированный метод. Сгенерированный метод получает сообщение запроса и StreamObserver, который используется для отправки ответов обратно клиенту:

@Component
public final class UserService extends UserServiceGrpc.UserServiceImplBase {

    @Override
    public void createUser(Message.RequestEvent request, StreamObserver<Message.ResponseEvent> responseObserver) { //(1)!
        var response = Message.ResponseEvent.newBuilder()
            .setId(ByteString.copyFromUtf8(UUID.randomUUID().toString()))
            .build();

        responseObserver.onNext(response); //(2)!
        responseObserver.onCompleted(); //(3)!
    }
}
  1. Сгенерированный метод получает сообщение запроса и StreamObserver для отправки ответа
  2. Отправляет клиенту одно сообщение ответа
  3. Сигнализирует о завершении вызова; для унарного метода вызывается ровно один раз, после единственного onNext
@Component
class UserService : UserServiceGrpc.UserServiceImplBase() {

    override fun createUser(request: Message.RequestEvent, responseObserver: StreamObserver<Message.ResponseEvent>) { //(1)!
        val response = Message.ResponseEvent.newBuilder()
            .setId(ByteString.copyFromUtf8(UUID.randomUUID().toString()))
            .build()

        responseObserver.onNext(response) //(2)!
        responseObserver.onCompleted() //(3)!
    }
}
  1. Сгенерированный метод получает сообщение запроса и StreamObserver для отправки ответа
  2. Отправляет клиенту одно сообщение ответа
  3. Сигнализирует о завершении вызова; для унарного метода вызывается ровно один раз, после единственного onNext

Серверная потоковая передача

Для серверного потокового RPC (returns (stream ...) в proto) клиент отправляет один запрос, а сервер возвращает много сообщений. Вызывайте onNext для каждого сообщения, а затем один раз onCompleted в конце:

@Override
public void getAllUsers(Message.RequestEvent request, StreamObserver<Message.ResponseEvent> responseObserver) {
    for (var user : userService.findAll()) {
        responseObserver.onNext(toResponse(user)); //(1)!
    }
    responseObserver.onCompleted(); //(2)!
}
  1. Отправляет одно из нескольких сообщений ответа
  2. Завершает поток ответа после последнего сообщения
override fun getAllUsers(request: Message.RequestEvent, responseObserver: StreamObserver<Message.ResponseEvent>) {
    userService.findAll().forEach { responseObserver.onNext(toResponse(it)) } //(1)!
    responseObserver.onCompleted() //(2)!
}
  1. Отправляет одно из нескольких сообщений ответа
  2. Завершает поток ответа после последнего сообщения

Клиентская потоковая передача

Для клиентского потокового RPC (rpc method(stream ...)) клиент отправляет много сообщений, а сервер отвечает один раз в конце. Сгенерированный метод возвращает StreamObserver, который получает входящие сообщения запроса; итоговый ответ формируется в onCompleted:

@Override
public StreamObserver<Message.RequestEvent> createUsers(StreamObserver<Message.ResponseEvent> responseObserver) {
    return new StreamObserver<>() {
        private final List<Message.RequestEvent> received = new ArrayList<>();

        @Override
        public void onNext(Message.RequestEvent value) {
            received.add(value); //(1)!
        }

        @Override
        public void onError(Throwable t) {
            responseObserver.onError(t); //(2)!
        }

        @Override
        public void onCompleted() {
            responseObserver.onNext(aggregate(received)); //(3)!
            responseObserver.onCompleted();
        }
    };
}
  1. Собирает каждое входящее сообщение запроса
  2. Пробрасывает ошибку потока со стороны клиента
  3. Формирует единственный агрегированный ответ после того, как клиент завершил отправку
override fun createUsers(responseObserver: StreamObserver<Message.ResponseEvent>): StreamObserver<Message.RequestEvent> {
    return object : StreamObserver<Message.RequestEvent> {
        private val received = mutableListOf<Message.RequestEvent>()

        override fun onNext(value: Message.RequestEvent) {
            received += value //(1)!
        }

        override fun onError(t: Throwable) {
            responseObserver.onError(t) //(2)!
        }

        override fun onCompleted() {
            responseObserver.onNext(aggregate(received)) //(3)!
            responseObserver.onCompleted()
        }
    }
}
  1. Собирает каждое входящее сообщение запроса
  2. Пробрасывает ошибку потока со стороны клиента
  3. Формирует единственный агрегированный ответ после того, как клиент завершил отправку

Двунаправленная потоковая передача

Для двунаправленного потокового RPC (rpc method(stream ...) returns (stream ...)) обе стороны обмениваются множеством сообщений в рамках одного вызова. Метод возвращает StreamObserver для входящих запросов и может отправлять ответы в любой момент через responseObserver:

@Override
public StreamObserver<Message.RequestEvent> updateUsers(StreamObserver<Message.ResponseEvent> responseObserver) {
    return new StreamObserver<>() {
        @Override
        public void onNext(Message.RequestEvent value) {
            responseObserver.onNext(process(value)); //(1)!
        }

        @Override
        public void onError(Throwable t) {
            responseObserver.onError(t);
        }

        @Override
        public void onCompleted() {
            responseObserver.onCompleted(); //(2)!
        }
    };
}
  1. Отвечает на каждое входящее сообщение по мере его поступления
  2. Завершает поток ответа, когда клиент прекращает отправку
override fun updateUsers(responseObserver: StreamObserver<Message.ResponseEvent>): StreamObserver<Message.RequestEvent> {
    return object : StreamObserver<Message.RequestEvent> {
        override fun onNext(value: Message.RequestEvent) {
            responseObserver.onNext(process(value)) //(1)!
        }

        override fun onError(t: Throwable) {
            responseObserver.onError(t)
        }

        override fun onCompleted() {
            responseObserver.onCompleted() //(2)!
        }
    }
}
  1. Отвечает на каждое входящее сообщение по мере его поступления
  2. Завершает поток ответа, когда клиент прекращает отправку

Обработка ошибок

Описание: gRPC представляет ошибки вызова с помощью кода io.grpc.Status и необязательного описания, а не с помощью кодов ответа HTTP. Чтобы завершить вызов с ошибкой, завершите observer ответа вызовом responseObserver.onError(status.asRuntimeException()) или выбросьте StatusRuntimeException из обработчика. Автоматически зарегистрированный TelemetryInterceptor наблюдает финальный Status при закрытии вызова и соответствующим образом записывает логирование, метрики и трассировку.

Причины: выбирайте код Status, соответствующий сбою, — например, Status.NOT_FOUND для отсутствующей сущности, Status.INVALID_ARGUMENT для некорректных входных данных, Status.UNAUTHENTICATED или Status.PERMISSION_DENIED для сбоев авторизации и Status.INTERNAL для непредвиденных ошибок сервера.

Рекомендации:

  • Прикрепляйте понятное человеку сообщение через withDescription(...) и сохраняйте исходное исключение через withCause(...), чтобы телеметрия могла его записать.
  • Завершайте вызов ровно один раз: никогда не вызывайте onError после onCompleted и не вызывайте ни один из них дважды.
  • Не раскрывайте клиентам внутренние детали исключений; сначала сопоставьте их с подходящим Status.

Пример обработки: унарный обработчик, который возвращает NOT_FOUND, когда сущность отсутствует, и сопоставляет непредвиденные сбои с INTERNAL:

@Override
public void getUser(Message.RequestEvent request, StreamObserver<Message.ResponseEvent> responseObserver) {
    try {
        var user = userService.getUser(request.getName())
            .orElseThrow(() -> Status.NOT_FOUND
                .withDescription("User not found: " + request.getName())
                .asRuntimeException()); //(1)!
        responseObserver.onNext(toResponse(user));
        responseObserver.onCompleted();
    } catch (StatusRuntimeException e) {
        responseObserver.onError(e); //(2)!
    } catch (Exception e) {
        responseObserver.onError(Status.INTERNAL
            .withDescription("Failed to get user")
            .withCause(e) //(3)!
            .asRuntimeException());
    }
}
  1. Строит ошибку NOT_FOUND с описанием
  2. Передает клиенту уже сопоставленную ошибку Status
  3. Сохраняет исходное исключение в качестве причины, чтобы телеметрия могла его записать
override fun getUser(request: Message.RequestEvent, responseObserver: StreamObserver<Message.ResponseEvent>) {
    try {
        val user = userService.getUser(request.name)
            ?: throw Status.NOT_FOUND
                .withDescription("User not found: ${request.name}")
                .asRuntimeException() //(1)!
        responseObserver.onNext(toResponse(user))
        responseObserver.onCompleted()
    } catch (e: StatusRuntimeException) {
        responseObserver.onError(e) //(2)!
    } catch (e: Exception) {
        responseObserver.onError(
            Status.INTERNAL
                .withDescription("Failed to get user")
                .withCause(e) //(3)!
                .asRuntimeException()
        )
    }
}
  1. Строит ошибку NOT_FOUND с описанием
  2. Передает клиенту уже сопоставленную ошибку Status
  3. Сохраняет исходное исключение в качестве причины, чтобы телеметрия могла его записать

Сигнатуры

Форма метода обработчика определяется контрактом proto и сгенерированным ...ImplBase:

Под Req и Resp подразумеваются сгенерированные типы сообщений запроса и ответа.

  • Унарный: void myMethod(Req request, StreamObserver<Resp> responseObserver)
  • Серверная потоковая передача: void myMethod(Req request, StreamObserver<Resp> responseObserver) (несколько onNext, один onCompleted)
  • Клиентская потоковая передача: StreamObserver<Req> myMethod(StreamObserver<Resp> responseObserver)
  • Двунаправленная потоковая передача: StreamObserver<Req> myMethod(StreamObserver<Resp> responseObserver)

Сгенерированный метод возвращает void (или StreamObserver для запроса), поэтому результаты доставляются через колбэки StreamObserver, а не через возвращаемое значение.

Под Req и Resp подразумеваются сгенерированные типы сообщений запроса и ответа.

  • Унарный: myMethod(request: Req, responseObserver: StreamObserver<Resp>)
  • Серверная потоковая передача: myMethod(request: Req, responseObserver: StreamObserver<Resp>) (несколько onNext, один onCompleted)
  • Клиентская потоковая передача: myMethod(responseObserver: StreamObserver<Resp>): StreamObserver<Req>
  • Двунаправленная потоковая передача: myMethod(responseObserver: StreamObserver<Resp>): StreamObserver<Req>

Сгенерированный метод возвращает Unit (или StreamObserver для запроса), поэтому результаты доставляются через колбэки StreamObserver, а не через возвращаемое значение.

Обработчики блокирующие: из метода обработчика можно напрямую обращаться к базе данных, HTTP-клиенту или gRPC-клиенту. Асинхронных, реактивных и suspend-сигнатур обработчиков нет — сервер выполняет обработчики на виртуальных потоках, смотрите Модель выполнения.

Модель выполнения

Каждому клиентскому соединению выделяется собственный однопоточный исполнитель на виртуальном потоке с именем grpc-<адрес клиента>. Исполнитель создается, когда транспорт готов, и останавливается, когда транспорт завершается.

  • Все колбэки перехватчиков и обработчиков для вызовов одного соединения выполняются на этом единственном виртуальном потоке — по одному за раз и в порядке поступления.
  • Блокироваться внутри обработчика безопасно: несущий поток освобождается. Но блокировка задерживает остальные вызовы того же соединения. Клиенту, которому нужны параллельные вызовы, следует открыть несколько соединений.
  • На время каждого колбэка Kora привязывает к этому потоку свой MDC и контекст OpenTelemetry, поэтому контекст логирования доступен внутри обработчика.

Перехватчики

io.grpc.ServerInterceptor обрабатывает вызов до того, как он будет передан в gRPC-сервис. Перехватчики подходят для сквозной логики: логирования, авторизации, трассировки, работы с Metadata и сопоставления ошибок.

В отличие от HTTP-сервера, модуль gRPC-сервера не имеет аннотаций @GrpcService или @InterceptWith: каждый ServerInterceptor, зарегистрированный как @Component, применяется глобально ко всем сервисам на сервере. Чтобы ограничить перехватчик одним сервисом или методом, анализируйте вызов во время выполнения — смотрите Ограничение области и авторизация.

Стандартные

При запуске сервера Kora добавляет один стандартный перехватчик:

  • TelemetryInterceptor — открывает GrpcServerObservation на каждый вызов, привязывает текущее наблюдение и контекст OpenTelemetry на время вызова и по его закрытию записывает логирование, метрики и трассировку в зависимости от подключенных модулей и настроек grpcServer.telemetry

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

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

Порядок выполнения

gRPC вызывает перехватчики в обратном порядке их регистрации, поэтому последний добавленный перехватчик выполняется первым (самый внешний). Поскольку Kora регистрирует пользовательские перехватчики первыми, а стандартный — последним, входящий вызов обрабатывается в таком порядке:

TelemetryInterceptor -> user interceptors -> handler

Следствия такого порядка:

  • TelemetryInterceptor оборачивает ваши перехватчики и обработчик, поэтому он наблюдает финальный Status — включая ошибки, выброшенные вашими перехватчиками или переданные через observer ответа.
  • Текущее наблюдение и контекст OpenTelemetry устанавливаются вокруг ваших перехватчиков и обработчика, поэтому они доступны внутри колбэков слушателя обработчика.
  • Когда пользовательских перехватчиков несколько, они выполняются в порядке, обратном порядку их регистрации в графе; не полагайтесь на конкретный порядок между ними для корректности.

Пользовательские

Чтобы добавить пользовательский перехватчик, создайте реализацию ServerInterceptor с аннотацией @Component:

@Component
public final class LoggingServerInterceptor implements ServerInterceptor {

    private static final Logger logger = LoggerFactory.getLogger(LoggingServerInterceptor.class);

    @Override
    public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call,
                                                                 Metadata headers,
                                                                 ServerCallHandler<ReqT, RespT> next) {
        logger.info("Incoming gRPC call: {}", call.getMethodDescriptor().getFullMethodName()); //(1)!

        return next.startCall(call, headers); //(2)!
    }
}
  1. getFullMethodName() возвращает service/method для перехватываемого вызова
  2. Передает вызов дальше; если вернуться, не вызвав startCall, вызов придется закрыть самостоятельно
@Component
class LoggingServerInterceptor : ServerInterceptor {

    private val logger = LoggerFactory.getLogger(LoggingServerInterceptor::class.java)

    override fun <ReqT : Any, RespT : Any> interceptCall(
        call: ServerCall<ReqT, RespT>,
        headers: Metadata,
        next: ServerCallHandler<ReqT, RespT>
    ): ServerCall.Listener<ReqT> {
        logger.info("Incoming gRPC call: {}", call.methodDescriptor.fullMethodName) //(1)!

        return next.startCall(call, headers) //(2)!
    }
}
  1. fullMethodName возвращает service/method для перехватываемого вызова
  2. Передает вызов дальше; если вернуться, не вызвав startCall, вызов придется закрыть самостоятельно

Ограничение области и авторизация

Поскольку перехватчик глобальный, ограничьте его конкретным сервисом или методом, анализируя call.getMethodDescriptor(): getServiceName() возвращает имя сервиса (сгенерированную константу ...Grpc.SERVICE_NAME), а getFullMethodName() возвращает service/method.

Заголовки запроса поступают в виде Metadata. Читайте заголовок с помощью Metadata.Key, а отклоняйте вызов, закрывая его с помощью Status и возвращая пустой слушатель, чтобы обработчик никогда не вызывался. Пример ниже применяет авторизацию по API-ключу только к одному сервису:

@ConfigSource("auth.apiKey")
public interface ApiKeyConfig {

    String value();
}
@Component
public final class ApiKeyServerInterceptor implements ServerInterceptor {

    private static final Metadata.Key<String> AUTHORIZATION =
        Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER); //(1)!

    private final ApiKeyConfig config;

    public ApiKeyServerInterceptor(ApiKeyConfig config) {
        this.config = config;
    }

    @Override
    public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call,
                                                                 Metadata headers,
                                                                 ServerCallHandler<ReqT, RespT> next) {
        if (!UserServiceGrpc.SERVICE_NAME.equals(call.getMethodDescriptor().getServiceName())) { //(2)!
            return next.startCall(call, headers);
        }

        var apiKey = headers.get(AUTHORIZATION); //(3)!
        if (!config.value().equals(apiKey)) {
            call.close(Status.UNAUTHENTICATED.withDescription("Invalid API key"), new Metadata()); //(4)!
            return new ServerCall.Listener<>() {}; //(5)!
        }

        return next.startCall(call, headers);
    }
}
  1. Metadata.Key для чтения заголовка authorization как ASCII-строки
  2. Применяет перехватчик только к UserService; другие сервисы проходят без изменений
  3. Читает значение заголовка из Metadata запроса
  4. Отклоняет вызов со статусом UNAUTHENTICATED
  5. Возвращает пустой слушатель, чтобы обработчик никогда не вызывался
@ConfigSource("auth.apiKey")
interface ApiKeyConfig {

    fun value(): String
}
@Component
class ApiKeyServerInterceptor(private val config: ApiKeyConfig) : ServerInterceptor {

    override fun <ReqT : Any, RespT : Any> interceptCall(
        call: ServerCall<ReqT, RespT>,
        headers: Metadata,
        next: ServerCallHandler<ReqT, RespT>
    ): ServerCall.Listener<ReqT> {
        if (UserServiceGrpc.SERVICE_NAME != call.methodDescriptor.serviceName) { //(2)!
            return next.startCall(call, headers)
        }

        val apiKey = headers.get(AUTHORIZATION) //(3)!
        if (config.value() != apiKey) {
            call.close(Status.UNAUTHENTICATED.withDescription("Invalid API key"), Metadata()) //(4)!
            return object : ServerCall.Listener<ReqT>() {} //(5)!
        }

        return next.startCall(call, headers)
    }

    companion object {
        private val AUTHORIZATION: Metadata.Key<String> =
            Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER) //(1)!
    }
}
  1. Metadata.Key для чтения заголовка authorization как ASCII-строки
  2. Применяет перехватчик только к UserService; другие сервисы проходят без изменений
  3. Читает значение заголовка из Metadata запроса
  4. Отклоняет вызов со статусом UNAUTHENTICATED
  5. Возвращает пустой слушатель, чтобы обработчик никогда не вызывался

Жизненный цикл и готовность

Сервером управляет компонент GrpcServer, который создается как компонент @Root и следует жизненному циклу приложения:

  • При запуске он собирает и стартует сервер на настроенном port. Если порт уже занят, запуск завершается ошибкой gRPC server failed to start on port '8090': port is already in use; stop the other process or configure a different port.
  • При выключении он выполняет штатное завершение: перестает принимать новые вызовы и ждет до shutdownWait завершения выполняющихся вызовов, затем принудительно завершает оставшиеся вызовы.

GrpcServer также реализует пробу готовности: сервер сообщает о неготовности во время запуска или выключения и о готовности только во время работы. В развертывании Kubernetes это позволяет пробе готовности отражать реальное состояние сервера и сливать трафик во время штатного завершения.

Рефлексия

Поддерживается gRPC Server Reflection, которая предоставляет информацию о доступных gRPC-сервисах на сервере. Рефлексия помогает клиентам и инструментам формировать запросы RPC во время выполнения без предварительно скомпилированной информации о сервисах. Например, ее использует gRPC CLI, который может исследовать описания proto сервера и отправлять тестовые вызовы RPC. gRPC Server Reflection поддерживается только для сервисов на основе proto.

Подробнее о gRPC Server Reflection можно узнать в руководстве grpc-java.

Зависимость

Необходимо дополнительно добавить зависимость gRPC Server Reflection.

Зависимость build.gradle:

implementation "io.grpc:grpc-services:1.83.1"

Зависимость build.gradle.kts:

implementation("io.grpc:grpc-services:1.83.1")

Конфигурация

Также необходимо включить сервис gRPC Server Reflection в конфигурации. Kora добавляет его на сервер только при наличии в приложении класса io.grpc.protobuf.services.ProtoReflectionServiceV1, поэтому одной конфигурации без зависимости недостаточно.

grpcServer {
    reflectionEnabled = false //(1)!
}
  1. Включает сервис gRPC Server Reflection (по умолчанию: false).
grpcServer:
  reflectionEnabled: false #(1)!
  1. Включает сервис gRPC Server Reflection (по умолчанию: false).

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

При включенной рефлексии инструменты вроде grpcurl могут обнаруживать сервисы и отправлять вызовы RPC без предварительно скомпилированного клиента. Для сервера, слушающего порт 8090:

grpcurl -plaintext localhost:8090 list #(1)!
grpcurl -plaintext localhost:8090 describe io.koraframework.generated.grpc.UserService #(2)!
grpcurl -plaintext -d '{"name": "Bob", "code": "123"}' \
    localhost:8090 io.koraframework.generated.grpc.UserService/createUser #(3)!
  1. Выводит список сервисов, предоставляемых сервером
  2. Описывает сервис и его методы
  3. Отправляет унарный RPC; -plaintext используется, потому что у сервера из примера нет TLS

Телеметрия

Наблюдаемость сервера обеспечивается TelemetryInterceptor через фасад GrpcServerTelemetry и настраивается в grpcServer.telemetry. Точки расширения находятся в io.koraframework.grpc.server.telemetry.

На каждый gRPC-вызов создается GrpcServerObservation: он собирает заголовки, отправленные и полученные сообщения, финальный Status и ошибку, а закрывается по завершении вызова. Фабрика GrpcServerTelemetryFactory по умолчанию зарегистрирована как @DefaultComponent, поэтому ее можно заменить целиком; либо отдельные части можно переопределить, зарегистрировав компонентом наследника DefaultGrpcServerLoggerFactory, DefaultGrpcServerMetricsFactory или DefaultGrpcServerBodyConverter. Если логирование, метрики и трассировка выключены все сразу, фабрика возвращает пустую реализацию телеметрии и вызовы не несут накладных расходов на наблюдаемость.

В логах, метриках и спанах сервер представляется именем kora-grpc.

Логирование

Логирование вызовов включается через grpcServer.telemetry.logging.enabled и пишется двумя логгерами:

  • io.koraframework.grpc.server.GrpcServer.requestGrpcCall received со структурированным полем grpcRequest
  • io.koraframework.grpc.server.GrpcServer.responseGrpcCall responded со структурированным полем grpcResponse

Структурированные поля содержат serverName, serverPort, serviceName и operation (service/method); в ответе дополнительно есть processingTime в миллисекундах, код Status в поле status и, для неудачного вызова, exceptionType. Вызов, завершившийся ошибкой, логируется на уровне WARN с приложенным исключением, успешный — на уровне INFO.

Уровень логгера добавляет детализацию сверх этого: DEBUG на логгере запроса добавляет Metadata запроса в поле headers, а TRACE добавляет тело сообщения, отрендеренное через DefaultGrpcServerBodyConverter.

logging.levels {
    "io.koraframework.grpc.server.GrpcServer.request" = "DEBUG"
    "io.koraframework.grpc.server.GrpcServer.response" = "TRACE"
}
logging:
  levels:
    "io.koraframework.grpc.server.GrpcServer.request": "DEBUG"
    "io.koraframework.grpc.server.GrpcServer.response": "TRACE"

Метрики

Метрикам нужен grpcServer.telemetry.metrics.enabled и MeterRegistry, который предоставляет модуль метрик. Модуль пишет единственный таймер rpc.server.duration с корзинами из grpcServer.telemetry.metrics.slo и тегами server.name, server.port, rpc.system (всегда grpc), rpc.service, rpc.method и rpc.grpc.status_code, плюс всё, что объявлено в grpcServer.telemetry.metrics.tags.

Метрики описаны в разделе Справочник метрик.

Трассировка

Трассировке нужен grpcServer.telemetry.tracing.enabled и Tracer, который предоставляет модуль трассировки. На каждый вызов создается спан вида SERVER с именем <service>/<method>; его родитель извлекается из Metadata запроса пропагатором W3C Trace Context, поэтому трасса, начатая вызывающей стороной, продолжается на сервере.

Спан несет атрибуты server.port, server.name, rpc.system, rpc.service, rpc.method и network.peer.address, плюс всё, что объявлено в grpcServer.telemetry.tracing.attributes; при закрытии добавляется rpc.grpc.status_code. Каждое отправленное и полученное сообщение добавляет событие rpc.message с атрибутом rpc.message.type. Статус Status, отличный от OK, или исключение переводят статус спана в ERROR.

Тестирование

gRPC-сервер занимает настоящий порт, поэтому @KoraAppTest поднимает его, а тест обращается к нему через обычный ManagedChannel:

@KoraAppTest(Application.class)
class UserServiceTests {

    @Test
    void createUser() {
        var channel = ManagedChannelBuilder.forAddress("localhost", 8090) //(1)!
            .usePlaintext()
            .build();

        try {
            var stub = UserServiceGrpc.newBlockingStub(channel); //(2)!
            var response = stub.createUser(Message.RequestEvent.newBuilder()
                .setName("Bob")
                .setCode("123")
                .build());

            assertFalse(response.getId().isEmpty());
        } finally {
            channel.shutdownNow();
        }
    }
}
  1. Порт, на котором поднят сервер, то есть grpcServer.port из тестовой конфигурации
  2. Блокирующая заглушка, сгенерированная из контракта proto
@KoraAppTest(Application::class)
class UserServiceTests {

    @Test
    fun createUser() {
        val channel = ManagedChannelBuilder.forAddress("localhost", 8090) //(1)!
            .usePlaintext()
            .build()

        try {
            val stub = UserServiceGrpc.newBlockingStub(channel) //(2)!
            val response = stub.createUser(
                Message.RequestEvent.newBuilder()
                    .setName("Bob")
                    .setCode("123")
                    .build()
            )

            assertFalse(response.id.isEmpty)
        } finally {
            channel.shutdownNow()
        }
    }
}
  1. Порт, на котором поднят сервер, то есть grpcServer.port из тестовой конфигурации
  2. Блокирующая заглушка, сгенерированная из контракта proto

Согласование версий: клиентской стороне теста нужен транспорт gRPC в тестовом classpath, и его версия должна совпадать с рантаймом gRPC, который приходит с io.koraframework:grpc-server, — 1.83.1. Закрепленная более старая версия компилируется без замечаний и падает только в рантайме с AbstractMethodError: ... does not define or inherit an implementation of the resolved method 'buildClientTransportServers(List, MetricRecorder)'.

Зависимость build.gradle:

testImplementation "io.koraframework:test-junit5"
testImplementation "io.grpc:grpc-netty:1.83.1"

Зависимость build.gradle.kts:

testImplementation("io.koraframework:test-junit5")
testImplementation("io.grpc:grpc-netty:1.83.1")