gRPC сервер
Модуль запускает gRPC-сервер на основе grpc-java и подключает к нему обработчики из графа приложения.
Обработчик — это BindableService, обычно класс, который наследует сгенерированный ...ImplBase и реализует унарные или потоковые методы RPC.
Kora создает NettyServerBuilder, добавляет сервисы сервера, пользовательские и стандартные ServerInterceptor, управляет жизненным циклом сервера и участвует в проверках готовности приложения.
Если параметров конфигурации недостаточно, итоговый NettyServerBuilder можно дополнительно настроить в коде через GrpcServerBuilderConfigurer.
Если нужен пошаговый разбор перед справочным описанием, смотрите gRPC-сервер и продвинутый gRPC-сервер.
Подключение¶
Зависимость build.gradle:
implementation "ru.tinkoff.kora:grpc-server"
implementation "io.grpc:grpc-protobuf:1.74.0"
implementation "javax.annotation:javax.annotation-api:1.3.2"
Модуль:
Зависимость build.gradle.kts:
implementation("ru.tinkoff.kora:grpc-server")
implementation("io.grpc:grpc-protobuf:1.74.0")
implementation("javax.annotation:javax.annotation-api:1.3.2")
Модуль:
Плагин¶
Код для gRPC-сервера генерируется с помощью gradle-плагина protobuf.
Плагин в build.gradle:
plugins {
id "com.google.protobuf" version "0.9.4"
}
protobuf {
protoc { artifact = "com.google.protobuf:protoc:3.25.3" }
plugins {
grpc { artifact = "io.grpc:protoc-gen-grpc-java:1.74.0" }
}
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.9.4")
}
protobuf {
protoc { artifact = "com.google.protobuf:protoc:3.25.3" }
plugins {
id("grpc") { artifact = "io.grpc:protoc-gen-grpc-java:1.74.0" }
}
generateProtoTasks {
ofSourceSet("main").forEach { it.plugins { id("grpc") { } } }
}
}
kotlin {
sourceSets.main {
kotlin.srcDir("build/generated/source/proto/main/grpc")
kotlin.srcDir("build/generated/source/proto/main/java")
}
}
Конфигурация¶
Обычно нужно задать только port; все остальные параметры имеют значения по умолчанию.
Минимальная конфигурация, которая привязывает порт и включает логирование:
Основные параметры конфигурации:
- Порт
gRPC-сервера(по умолчанию:8090). - Максимальный размер входящего сообщения (по умолчанию:
4MiB). - Включает сервис
gRPC Server Reflection(по умолчанию:false). - Включает виртуальные потоки для обработки вызовов, требует
Java 21+(по умолчанию:false).
- Порт
gRPC-сервера(по умолчанию:8090). - Максимальный размер входящего сообщения (по умолчанию:
4MiB). - Включает сервис
gRPC Server Reflection(по умолчанию:false).
Полная конфигурация
Пример полной конфигурации, описанной в классе GrpcServerConfig:
grpcServer {
port = 8090 //(1)!
maxMessageSize = "4MiB" //(2)!
reflectionEnabled = false //(3)!
shutdownWait = "30s" //(4)!
maxConnectionAge = "0s" //(5)!
maxConnectionAgeGrace = "0s" //(6)!
keepAliveTime = "0s" //(7)!
keepAliveTimeout = "0s" //(8)!
telemetry {
logging {
enabled = false //(9)!
}
metrics {
enabled = true //(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"
}
}
}
}
- Порт
gRPC-сервера(по умолчанию:8090). - Максимальный размер входящего сообщения (по умолчанию:
4MiB). Может быть указан в виде числа байт или как4MiB,4MB,1000Kbи подобных значений. - Включает сервис
gRPC Server Reflection(по умолчанию:false). - Время ожидания обработки перед выключением сервера при штатном завершении (по умолчанию:
30s). - Задает пользовательское максимальное время жизни соединения, после которого соединение штатно завершается (по умолчанию: не задано, опционально). К значению добавляется случайное отклонение +/-10%.
- Задает дополнительное время для штатного завершения соединения после достижения максимального времени жизни соединения (по умолчанию: не задано, опционально). Вызовы
RPC, которые не успевают завершиться, отменяются, чтобы соединение могло завершиться. - Задает интервал между кадрами
PING(по умолчанию: не задано, опционально). - Тайм-аут подтверждения кадра
PING(по умолчанию: не задано, опционально). Если подтверждение не получено за это время, соединение закрывается. - Включает логирование модуля (по умолчанию:
false). - Включает метрики модуля (по умолчанию:
true). - Настройка SLO для метрики DistributionSummary (по умолчанию:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Теги метрик (по умолчанию:
{}). - Включает трассировку модуля (по умолчанию:
true). - Атрибуты трассировки (по умолчанию:
{}).
grpcServer:
port: 8090 #(1)!
maxMessageSize: "4MiB" #(2)!
reflectionEnabled: false #(3)!
shutdownWait: "30s" #(4)!
maxConnectionAge: "0s" #(5)!
maxConnectionAgeGrace: "0s" #(6)!
keepAliveTime: "0s" #(7)!
keepAliveTimeout: "0s" #(8)!
telemetry:
logging:
enabled: false #(9)!
metrics:
enabled: true #(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
- Порт
gRPC-сервера(по умолчанию:8090). - Максимальный размер входящего сообщения (по умолчанию:
4MiB). Может быть указан в виде числа байт или как4MiB,4MB,1000Kbи подобных значений. - Включает сервис
gRPC Server Reflection(по умолчанию:false). - Время ожидания обработки перед выключением сервера при штатном завершении (по умолчанию:
30s). - Задает пользовательское максимальное время жизни соединения, после которого соединение штатно завершается (по умолчанию: не задано, опционально). К значению добавляется случайное отклонение +/-10%.
- Задает дополнительное время для штатного завершения соединения после достижения максимального времени жизни соединения (по умолчанию: не задано, опционально). Вызовы
RPC, которые не успевают завершиться, отменяются, чтобы соединение могло завершиться. - Задает интервал между кадрами
PING(по умолчанию: не задано, опционально). - Тайм-аут подтверждения кадра
PING(по умолчанию: не задано, опционально). Если подтверждение не получено за это время, соединение закрывается. - Включает логирование модуля (по умолчанию:
false). - Включает метрики модуля (по умолчанию:
true). - Настройка SLO для метрики DistributionSummary (по умолчанию:
ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Теги метрик (по умолчанию:
{}). - Включает трассировку модуля (по умолчанию:
true). - Атрибуты трассировки (по умолчанию:
{}).
Обработчики¶
Обработчик — это класс, который наследует сгенерированный ...ImplBase и регистрируется в графе приложения с помощью аннотации @Component.
Класс ...ImplBase создается из контракта proto с помощью gradle-плагина protobuf; вы переопределяете его методы RPC, чтобы реализовать поведение сервера.
Обычные компоненты Kora, такие как сервисы и репозитории, можно внедрить в обработчик через его конструктор.
Рассмотрим контракт proto с единственным унарным методом:
syntax = "proto3";
package ru.tinkoff.kora.generated.grpc;
service UserService {
rpc createUser(RequestEvent) returns (ResponseEvent) {} //(1)!
}
message RequestEvent {
string name = 1;
string code = 2;
}
message ResponseEvent {
bytes id = 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)!
}
}
- Сгенерированный метод получает сообщение запроса и
StreamObserverдля отправки ответа - Отправляет клиенту одно сообщение ответа
- Сигнализирует о завершении вызова; для унарного метода вызывается ровно один раз, после единственного
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)!
}
}
- Сгенерированный метод получает сообщение запроса и
StreamObserverдля отправки ответа - Отправляет клиенту одно сообщение ответа
- Сигнализирует о завершении вызова; для унарного метода вызывается ровно один раз, после единственного
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)!
}
- Отправляет одно из нескольких сообщений ответа
- Завершает поток ответа после последнего сообщения
override fun getAllUsers(request: Message.RequestEvent, responseObserver: StreamObserver<Message.ResponseEvent>) {
userService.findAll().forEach { responseObserver.onNext(toResponse(it)) } //(1)!
responseObserver.onCompleted() //(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();
}
};
}
- Собирает каждое входящее сообщение запроса
- Пробрасывает ошибку потока со стороны клиента
- Формирует единственный агрегированный ответ после того, как клиент завершил отправку
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()
}
}
}
- Собирает каждое входящее сообщение запроса
- Пробрасывает ошибку потока со стороны клиента
- Формирует единственный агрегированный ответ после того, как клиент завершил отправку
Двунаправленная потоковая передача¶
Для двунаправленного потокового 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)!
}
};
}
- Отвечает на каждое входящее сообщение по мере его поступления
- Завершает поток ответа, когда клиент прекращает отправку
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)!
}
}
}
- Отвечает на каждое входящее сообщение по мере его поступления
- Завершает поток ответа, когда клиент прекращает отправку
Обработка ошибок¶
Описание: gRPC представляет ошибки вызова с помощью кода io.grpc.Status
и необязательного описания, а не с помощью кодов ответа HTTP.
Чтобы завершить вызов с ошибкой, завершите observer ответа вызовом responseObserver.onError(status.asRuntimeException())
или выбросьте StatusRuntimeException из обработчика.
Автоматически зарегистрированный TelemetryInterceptor наблюдает финальный Status при закрытии вызова
(в close, onHalfClose, onCancel и onComplete) и соответствующим образом записывает логирование, метрики и трассировку.
Причины: выбирайте код 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());
}
}
- Строит ошибку
NOT_FOUNDс описанием - Передает клиенту уже сопоставленную ошибку
Status - Сохраняет исходное исключение в качестве причины, чтобы телеметрия могла его записать
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()
)
}
}
- Строит ошибку
NOT_FOUNDс описанием - Передает клиенту уже сопоставленную ошибку
Status - Сохраняет исходное исключение в качестве причины, чтобы телеметрия могла его записать
Сигнатуры¶
Форма метода обработчика определяется контрактом 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>
Когда вы генерируете корутинные заглушки с помощью плагина grpc-kotlin (io.grpc:protoc-gen-grpc-kotlin)
и наследуете сгенерированный ...CoroutineImplBase, методы обработчика могут быть suspend-функциями (а потоковые методы могут использовать Flow).
Kora автоматически регистрирует CoroutineContextInjectInterceptor, который внедряет Context Kora в CoroutineContext обработчика;
он активируется только при наличии kotlinx-coroutines в classpath.
Перехватчики¶
io.grpc.ServerInterceptor обрабатывает вызов до того, как он будет передан в gRPC-сервис.
Перехватчики подходят для сквозной логики: логирования, авторизации, трассировки, работы с Metadata и сопоставления ошибок.
В отличие от HTTP-сервера, модуль gRPC-сервера не имеет аннотаций @GrpcService или @InterceptWith:
каждый ServerInterceptor, зарегистрированный как @Component, применяется глобально ко всем сервисам на сервере.
Чтобы ограничить перехватчик одним сервисом или методом, анализируйте вызов во время выполнения — смотрите Ограничение области и авторизация.
Стандартные¶
При запуске сервера Kora добавляет стандартные перехватчики:
TelemetryInterceptor— включает телеметрию сервера (логирование, метрики, трассировку) в зависимости от подключенных модулей и настроекgrpcServer.telemetry, а также сопоставляет финальныйStatus/исключение при закрытии вызоваContextServerInterceptor— пробрасываетContextKora в обработку вызова, чтобы он был доступен внутри обработчикаCoroutineContextInjectInterceptor— добавляет поддержкуCoroutineContextдля корутинных обработчиков наKotlin(активен только при наличииkotlinx-coroutinesв classpath)
Пользовательские бины ServerInterceptor из графа приложения добавляются в NettyServerBuilder перед стандартными перехватчиками.
Для полной настройки NettyServerBuilder используйте GrpcServerBuilderConfigurer.
Порядок выполнения¶
gRPC вызывает перехватчики в обратном порядке их регистрации, поэтому последний добавленный перехватчик выполняется первым (самый внешний). Поскольку Kora регистрирует пользовательские перехватчики первыми, а стандартные — последними, входящий вызов обрабатывается в таком порядке:
CoroutineContextInjectInterceptor -> ContextServerInterceptor -> TelemetryInterceptor -> user interceptors -> handler
Следствия такого порядка:
ContextKora иCoroutineContextKotlin устанавливаются вокруг ваших перехватчиков и обработчика, поэтому они доступны внутри колбэков слушателя обработчика.TelemetryInterceptorоборачивает ваши перехватчики и обработчик, поэтому он наблюдает финальныйStatus(включая ошибки, выброшенные или переданные через observer ответа).- Когда пользовательских перехватчиков несколько, они выполняются в порядке, обратном порядку их регистрации в графе; не полагайтесь на конкретный порядок между ними для корректности.
Пользовательские¶
Чтобы добавить пользовательский перехватчик, создайте реализацию ServerInterceptor с аннотацией @Component:
@Component
public class GrpcExceptionHandlerServerInterceptor implements ServerInterceptor {
@Override
public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> serverCall,
Metadata metadata,
ServerCallHandler<ReqT, RespT> serverCallHandler) {
// do something
return serverCallHandler.startCall(serverCall, metadata);
}
}
@Component
class GrpcExceptionHandlerServerInterceptor : ServerInterceptor {
override fun <ReqT, RespT> interceptCall(
serverCall: ServerCall<ReqT, RespT>,
metadata: Metadata,
serverCallHandler: ServerCallHandler<ReqT, RespT>
): ServerCall.Listener<ReqT> {
// do something
return serverCallHandler.startCall(serverCall, metadata)
}
}
Ограничение области и авторизация¶
Поскольку перехватчик глобальный, ограничьте его конкретным сервисом или методом, анализируя call.getMethodDescriptor():
getServiceName() возвращает имя сервиса (сгенерированную константу ...Grpc.SERVICE_NAME), а getFullMethodName() возвращает service/method.
Заголовки запроса поступают в виде Metadata.
Читайте заголовок с помощью Metadata.Key, а отклоняйте вызов, закрывая его с помощью Status и возвращая пустой слушатель, чтобы обработчик никогда не вызывался.
Пример ниже применяет авторизацию по API-ключу только к одному сервису:
@Component
public final class ApiKeyServerInterceptor implements ServerInterceptor {
private static final Metadata.Key<String> AUTHORIZATION =
Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER); //(1)!
@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 (apiKey == null || !apiKey.equals("secret")) {
call.close(Status.UNAUTHENTICATED.withDescription("Invalid API key"), new Metadata()); //(4)!
return new ServerCall.Listener<>() {}; //(5)!
}
return next.startCall(call, headers);
}
}
Metadata.Keyдля чтения заголовкаauthorizationкак ASCII-строки- Применяет перехватчик только к
UserService; другие сервисы проходят без изменений - Читает значение заголовка из
Metadataзапроса - Отклоняет вызов со статусом
UNAUTHENTICATED - Возвращает пустой слушатель, чтобы обработчик никогда не вызывался
@Component
class ApiKeyServerInterceptor : 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 (apiKey != "secret") {
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)!
}
}
Metadata.Keyдля чтения заголовкаauthorizationкак ASCII-строки- Применяет перехватчик только к
UserService; другие сервисы проходят без изменений - Читает значение заголовка из
Metadataзапроса - Отклоняет вызов со статусом
UNAUTHENTICATED - Возвращает пустой слушатель, чтобы обработчик никогда не вызывался
Жизненный цикл и готовность¶
Сервером управляет компонент GrpcNettyServer, который создается как компонент @Root
и следует жизненному циклу приложения:
- При запуске он создает и стартует сервер
Nettyна настроенномport. Если порт уже занят, запуск завершается с понятной ошибкой. - При выключении он выполняет штатное завершение: перестает принимать новые вызовы и ждет до
shutdownWaitзавершения выполняющихся вызовов, затем принудительно завершает оставшиеся вызовы.
GrpcNettyServer также реализует пробу готовности: сервер сообщает о неготовности во время запуска или выключения
и о готовности только во время работы. В развертывании Kubernetes это позволяет пробе готовности отражать реальное состояние сервера и сливать трафик во время штатного завершения.
Рефлексия¶
Поддерживается gRPC Server Reflection, которая предоставляет информацию о доступных gRPC-сервисах на сервере.
Рефлексия помогает клиентам и инструментам формировать запросы RPC во время выполнения без предварительно скомпилированной информации о сервисах.
Например, ее использует gRPC CLI, который может исследовать описания proto сервера и отправлять тестовые вызовы RPC.
gRPC Server Reflection поддерживается только для сервисов на основе proto.
Подробнее о gRPC Server Reflection можно узнать в руководстве grpc-java.
Зависимость¶
Необходимо дополнительно добавить зависимость gRPC Server Reflection.
Зависимость build.gradle:
Зависимость build.gradle.kts:
Конфигурация¶
Также необходимо включить сервис gRPC Server Reflection в конфигурации.
Kora добавляет его на сервер только при наличии в приложении класса io.grpc.protobuf.services.ProtoReflectionService, поэтому одной конфигурации без зависимости недостаточно.
Использование¶
При включенной рефлексии инструменты вроде grpcurl могут обнаруживать сервисы и отправлять вызовы RPC без предварительно скомпилированного клиента.
Для сервера, слушающего порт 8090:
grpcurl -plaintext localhost:8090 list #(1)!
grpcurl -plaintext localhost:8090 describe ru.tinkoff.kora.generated.grpc.UserService #(2)!
grpcurl -plaintext -d '{"name": "Bob", "code": "123"}' \
localhost:8090 ru.tinkoff.kora.generated.grpc.UserService/createUser #(3)!
- Выводит список сервисов, предоставляемых сервером
- Описывает сервис и его методы
- Отправляет унарный
RPC;-plaintextиспользуется, потому что у сервера из примера нетTLS
Телеметрия¶
gRPC Server использует контракт телеметрии для логирования, метрик и трассировки вызовов.
Конфигурация телеметрии (секция telemetry { logging / metrics / tracing }) описана в разделе Конфигурация.
Точки расширения находятся в ru.tinkoff.kora.grpc.server.common.telemetry.
Наблюдаемость сервера обеспечивается TelemetryInterceptor через фасад GrpcServerTelemetry и настраивается в grpcServer.telemetry.
Для каждого gRPC-вызова создаётся GrpcServerTelemetry.GrpcServerTelemetryContext, который закрывается по завершении вызова.
Вызов описывается через параметры обработчика телеметрии, включая сервис, метод, статус ответа и длительность.
Фабрика по умолчанию DefaultGrpcServerTelemetryFactory объединяет три фабрики:
- GrpcServerLoggerFactory строит GrpcServerLogger для логирования начала/конца вызова;
- GrpcServerMetricsFactory строит GrpcServerMetrics для записи метрик вызовов;
- GrpcServerTracerFactory строит GrpcServerTracer для распределённой трассировки.
Метрики и трассировка описаны в разделе Справочник метрик.