Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions core/src/main/java/tech/ydb/core/grpc/Observability.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
package tech.ydb.core.grpc;

public final class Observability {
public static final String TRACING_CHAIN = "ydb-sdk-tracing";
public static final String METRICS_CHAIN = "ydb-sdk-metrics";

public static final String TRACING_CHAIN_VERSION = "0.1.0";
public static final String METRICS_CHAIN_VERSION = "0.1.0";

public static final String TRACING_CHAIN_TOKEN = TRACING_CHAIN + "/" + TRACING_CHAIN_VERSION;
public static final String METRICS_CHAIN_TOKEN = METRICS_CHAIN + "/" + METRICS_CHAIN_VERSION;

private Observability() {
}

public static String appendChains(String base, boolean tracing, boolean metrics) {
if (!tracing && !metrics) {
return base;
}
StringBuilder sb = new StringBuilder(base);
if (tracing) {
sb.append(';').append(TRACING_CHAIN_TOKEN);
}
if (metrics) {
sb.append(';').append(METRICS_CHAIN_TOKEN);
}
return sb.toString();
}
}
1 change: 0 additions & 1 deletion core/src/main/java/tech/ydb/core/grpc/YdbHeaders.java
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@ private YdbHeaders() { }
public static ClientInterceptor createMetadataInterceptor(GrpcTransportBuilder builder) {
Metadata extraHeaders = new Metadata();
extraHeaders.put(YdbHeaders.DATABASE, builder.getDatabase());
extraHeaders.put(YdbHeaders.BUILD_INFO, builder.getBuildInfo());
String appName = builder.getApplicationName();
if (appName != null) {
extraHeaders.put(YdbHeaders.APPLICATION_NAME, appName);
Expand Down
8 changes: 8 additions & 0 deletions core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,10 @@ protected BaseGrpcTransport(EndpointRecord serverEndpoint) {

protected abstract GrpcChannel getChannel(GrpcRequestSettings settings);

protected String getBuildInfo() {
return null;
}

protected void pessimizeEndpoint(EndpointRecord endpoint, String reason) {
// nothing to pessimize
}
Expand Down Expand Up @@ -244,6 +248,10 @@ private static Status deadlineExpiredStatus(MethodDescriptor<?, ?> method, GrpcR

private Metadata makeMetadataFromSettings(GrpcRequestSettings settings, EndpointRecord endpoint) {
Metadata metadata = new Metadata();
String buildInfo = getBuildInfo();
if (buildInfo != null) {
metadata.put(YdbHeaders.BUILD_INFO, buildInfo);
}
String token = getAuthCallOptions().getToken();
if (token != null) {
metadata.put(YdbHeaders.AUTH_TICKET, token);
Expand Down
22 changes: 22 additions & 0 deletions core/src/main/java/tech/ydb/core/impl/BuildInfoChainSupport.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
package tech.ydb.core.impl;

import tech.ydb.core.grpc.GrpcTransport;

/**
* Internal bridge for enabling SDK observability adoption chains on the transport after it has been built.
* <p>
* Unlike tracing (configured on the builder before {@code build()}), metrics adoption becomes known only
* once a higher-level client is wired with a real {@code Meter}, which happens after the transport exists.
* This class avoids adding a public transport hook by unwrapping the default SDK implementation and calling
* a package-private extension point on it; non-SDK transports silently no-op.
*/
public final class BuildInfoChainSupport {
private BuildInfoChainSupport() {
}

public static void enableMetricsChain(GrpcTransport transport) {
if (transport instanceof YdbTransportImpl) {
((YdbTransportImpl) transport).enableMetricsChain();
}
}
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package tech.ydb.core.impl;

import java.util.concurrent.ScheduledExecutorService;
import java.util.function.Supplier;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand All @@ -22,18 +23,21 @@ public class FixedCallOptionsTransport extends BaseGrpcTransport {
private final AuthCallOptions callOptions;
private final String database;
private final GrpcChannel channel;
private final Supplier<String> buildInfoSupplier;

public FixedCallOptionsTransport(
ScheduledExecutorService scheduler,
AuthCallOptions callOptions,
String database,
EndpointRecord endpoint,
ManagedChannelFactory channelFactory) {
ManagedChannelFactory channelFactory,
Supplier<String> buildInfoSupplier) {
super(endpoint);
this.scheduler = scheduler;
this.callOptions = callOptions;
this.database = database;
this.channel = new GrpcChannel(endpoint, channelFactory);
this.buildInfoSupplier = buildInfoSupplier;
}

@Override
Expand All @@ -46,6 +50,11 @@ public String getDatabase() {
return database;
}

@Override
protected String getBuildInfo() {
return buildInfoSupplier.get();
}

@Override
protected void shutdown() {
channel.shutdown();
Expand Down
30 changes: 29 additions & 1 deletion core/src/main/java/tech/ydb/core/impl/YdbTransportImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,14 @@
import tech.ydb.core.grpc.GrpcRequestSettings;
import tech.ydb.core.grpc.GrpcTransport;
import tech.ydb.core.grpc.GrpcTransportBuilder;
import tech.ydb.core.grpc.Observability;
import tech.ydb.core.impl.auth.AuthCallOptions;
import tech.ydb.core.impl.pool.EndpointPool;
import tech.ydb.core.impl.pool.EndpointRecord;
import tech.ydb.core.impl.pool.GrpcChannel;
import tech.ydb.core.impl.pool.GrpcChannelPool;
import tech.ydb.core.impl.pool.ManagedChannelFactory;
import tech.ydb.core.tracing.NoopTracer;
import tech.ydb.core.tracing.Tracer;

/**
Expand All @@ -38,13 +40,19 @@ public class YdbTransportImpl extends BaseGrpcTransport {
private final YdbDiscovery discovery;
private final Tracer tracer;

private final String baseBuildInfo;
private final boolean tracingChainEnabled;
private volatile boolean metricsChainEnabled;

public YdbTransportImpl(GrpcTransportBuilder builder) {
super(builder);
BalancingSettings balancingSettings = getBalancingSettings(builder);
Duration discoveryTimeout = Duration.ofMillis(builder.getDiscoveryTimeoutMillis());

this.database = Strings.nullToEmpty(builder.getDatabase());
this.tracer = builder.getTracer();
this.baseBuildInfo = builder.getBuildInfo();
this.tracingChainEnabled = tracer != NoopTracer.getInstance();

logger.info("Create YDB transport with endpoint {} and {}", serverEndpoint, balancingSettings);

Expand Down Expand Up @@ -124,6 +132,19 @@ public Tracer getTracer() {
return tracer;
}

void enableMetricsChain() {
this.metricsChainEnabled = true;
}

@Override
protected String getBuildInfo() {
return baseBuildInfo;
}

private String discoveryBuildInfo() {
return Observability.appendChains(baseBuildInfo, tracingChainEnabled, metricsChainEnabled);
}

@Override
public AuthCallOptions getAuthCallOptions() {
return callOptions;
Expand Down Expand Up @@ -167,7 +188,14 @@ public CompletableFuture<Boolean> handleEndpoints(List<EndpointRecord> endpoints

@Override
public GrpcTransport createDiscoveryTransport() {
return new FixedCallOptionsTransport(scheduler, callOptions, database, serverEndpoint, channelFactory);
return new FixedCallOptionsTransport(
scheduler,
callOptions,
database,
serverEndpoint,
channelFactory,
YdbTransportImpl.this::discoveryBuildInfo
);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,8 @@ public AuthCallOptions(

AuthRpcProvider<? super GrpcAuthRpc> authProvider = builder.getAuthProvider();
if (authProvider != null) {
GrpcAuthRpc rpc = new GrpcAuthRpc(endpoints, scheduler, builder.getDatabase(), channelFactory);
GrpcAuthRpc rpc = new GrpcAuthRpc(
endpoints, scheduler, builder.getDatabase(), channelFactory, builder.getBuildInfo());
authIdentity = builder.getAuthProvider().createAuthIdentity(rpc);
} else {
authIdentity = null;
Expand Down
11 changes: 8 additions & 3 deletions core/src/main/java/tech/ydb/core/impl/auth/GrpcAuthRpc.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,20 +19,23 @@ public class GrpcAuthRpc {
private final ScheduledExecutorService scheduler;
private final String database;
private final ManagedChannelFactory channelFactory;
private final String buildInfo;
private final AtomicInteger endpointIdx = new AtomicInteger();

public GrpcAuthRpc(
List<EndpointRecord> endpoints,
ScheduledExecutorService scheduler,
String database,
ManagedChannelFactory channelFactory) {
ManagedChannelFactory channelFactory,
String buildInfo) {
if (endpoints == null || endpoints.isEmpty()) {
throw new IllegalStateException("Empty endpoints list for auth rpc");
}
this.endpoints = endpoints;
this.scheduler = scheduler;
this.database = database;
this.channelFactory = channelFactory;
this.buildInfo = buildInfo;
}

public ExecutorService getExecutor() {
Expand All @@ -55,13 +58,15 @@ public void changeEndpoint() {
}

public GrpcTransport createTransport() {
// For auth provider we use transport without auth (with default CallOptions)
// For auth provider we use transport without auth (with default CallOptions).
// Auth requests report only the base build-info chain, matching regular (non-discovery) requests.
return new FixedCallOptionsTransport(
scheduler,
new AuthCallOptions(),
database,
endpoints.get(endpointIdx.get()),
channelFactory
channelFactory,
() -> buildInfo
);
}

Expand Down
41 changes: 41 additions & 0 deletions core/src/test/java/tech/ydb/core/grpc/ObservabilityTest.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
package tech.ydb.core.grpc;

import org.junit.Assert;
import org.junit.Test;

public class ObservabilityTest {

private static final String BASE = "ydb-java-sdk/1.2.3";

@Test
public void disabledObservabilityLeavesBaseUntouched() {
Assert.assertEquals(BASE, Observability.appendChains(BASE, false, false));
}

@Test
public void tracingOnlyAppendsTracingChain() {
Assert.assertEquals(
BASE + ";" + Observability.TRACING_CHAIN_TOKEN,
Observability.appendChains(BASE, true, false));
}

@Test
public void metricsOnlyAppendsMetricsChain() {
Assert.assertEquals(
BASE + ";" + Observability.METRICS_CHAIN_TOKEN,
Observability.appendChains(BASE, false, true));
}

@Test
public void bothChainsKeepTracingBeforeMetrics() {
Assert.assertEquals(
BASE + ";" + Observability.TRACING_CHAIN_TOKEN + ";" + Observability.METRICS_CHAIN_TOKEN,
Observability.appendChains(BASE, true, true));
}

@Test
public void chainTokensFollowNameSlashVersionFormat() {
Assert.assertEquals("ydb-sdk-tracing/0.1.0", Observability.TRACING_CHAIN_TOKEN);
Assert.assertEquals("ydb-sdk-metrics/0.1.0", Observability.METRICS_CHAIN_TOKEN);
}
}
10 changes: 10 additions & 0 deletions core/src/test/java/tech/ydb/core/impl/MockedCall.java
Original file line number Diff line number Diff line change
Expand Up @@ -15,15 +15,25 @@

public abstract class MockedCall<ResT, RespT> extends ClientCall<ResT, RespT> {
private final Executor executor;
private volatile Metadata lastHeaders;

protected MockedCall(Executor executor) {
this.executor = executor;
}

protected abstract void complete(Listener<RespT> listener);

/**
* Returns the request headers captured on the most recent {@link #start} call, or {@code null} if the
* call has not been started yet. Handy for asserting per-request metadata such as x-ydb-sdk-build-info.
*/
public Metadata getLastHeaders() {
return lastHeaders;
}

@Override
public void start(Listener<RespT> listener, Metadata headers) {
this.lastHeaders = headers;
executor.execute(() -> complete(listener));
}

Expand Down
4 changes: 3 additions & 1 deletion core/src/test/java/tech/ydb/core/impl/YdbDiscoveryTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -255,7 +255,9 @@ public Instant instant() {
@Override
public GrpcTransport createDiscoveryTransport() {
EndpointRecord discovery = new EndpointRecord("unknown", 1234);
return new FixedCallOptionsTransport(scheduler, new AuthCallOptions(), "/test", discovery, channelFactory);
return new FixedCallOptionsTransport(
scheduler, new AuthCallOptions(), "/test", discovery, channelFactory, () -> null
);
}

@Override
Expand Down
Loading
Loading