From c6f065299051f14c36fbb5db0a467104a9cf5d3f Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 11:01:58 +0300 Subject: [PATCH 1/7] Make SDK tracer span current around gRPC newCall --- .../tech/ydb/core/impl/BaseGrpcTransport.java | 45 +++++++++++-------- 1 file changed, 26 insertions(+), 19 deletions(-) diff --git a/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java b/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java index eaf2ff069..3ddb8e9c3 100644 --- a/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java +++ b/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java @@ -37,6 +37,7 @@ import tech.ydb.core.impl.call.UnaryCall; import tech.ydb.core.impl.pool.EndpointRecord; import tech.ydb.core.impl.pool.GrpcChannel; +import tech.ydb.core.tracing.Scope; import tech.ydb.core.tracing.Span; /** @@ -124,16 +125,19 @@ public CompletableFuture> unaryCall( return CompletableFuture.completedFuture(deadlineExpiredResult(method, settings)); } - ClientCall call = channel.getReadyChannel().newCall(method, options); - ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); + // Make SDK span current so gRPC ClientInterceptors nest under it at newCall. + try (Scope ignored = settings.getSpan().makeCurrent()) { + ClientCall call = channel.getReadyChannel().newCall(method, options); + ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); - if (logger.isTraceEnabled()) { - logger.trace("UnaryCall[{}] with method {} and endpoint {} created", - traceId, method.getFullMethodName(), endpoint.getHostAndPort()); + if (logger.isTraceEnabled()) { + logger.trace("UnaryCall[{}] with method {} and endpoint {} created", + traceId, method.getFullMethodName(), endpoint.getHostAndPort()); + } + Metadata metadata = makeMetadataFromSettings(settings, endpoint); + return new UnaryCall<>(traceId, endpoint.getHostAndPort(), call, handler, settings.getSpan()) + .startCall(request, metadata); } - Metadata metadata = makeMetadataFromSettings(settings, endpoint); - return new UnaryCall<>(traceId, endpoint.getHostAndPort(), call, handler, settings.getSpan()) - .startCall(request, metadata); } catch (UnexpectedResultException ex) { logger.warn("UnaryCall[{}] got unexpected status {}", traceId, ex.getStatus()); return CompletableFuture.completedFuture(Result.fail(ex)); @@ -163,19 +167,22 @@ public GrpcReadStream readStreamCall( return new EmptyStream<>(deadlineExpiredStatus(method, settings)); } - ClientCall call = channel.getReadyChannel().newCall(method, options); - ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); + // Make SDK span current so gRPC ClientInterceptors nest under it at newCall. + try (Scope ignored = settings.getSpan().makeCurrent()) { + ClientCall call = channel.getReadyChannel().newCall(method, options); + ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); - if (logger.isTraceEnabled()) { - logger.trace("ReadStreamCall[{}] with method {} and endpoint {} created", - traceId, method.getFullMethodName(), endpoint.getHostAndPort() - ); - } + if (logger.isTraceEnabled()) { + logger.trace("ReadStreamCall[{}] with method {} and endpoint {} created", + traceId, method.getFullMethodName(), endpoint.getHostAndPort() + ); + } - Metadata metadata = makeMetadataFromSettings(settings, endpoint); - GrpcFlowControl flowCtrl = settings.getFlowControl(); - return new ReadStreamCall<>(traceId, endpoint.getHostAndPort(), call, flowCtrl, request, metadata, handler, - settings.getSpan()); + Metadata metadata = makeMetadataFromSettings(settings, endpoint); + GrpcFlowControl flowCtrl = settings.getFlowControl(); + return new ReadStreamCall<>(traceId, endpoint.getHostAndPort(), call, flowCtrl, request, metadata, + handler, settings.getSpan()); + } } catch (UnexpectedResultException ex) { logger.warn("ReadStreamCall[{}] got unexpected status {}", traceId, ex.getStatus()); return new EmptyStream<>(ex.getStatus()); From c8692faec05552eb0440ce228a70bea5138dc41e Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 14:49:14 +0300 Subject: [PATCH 2/7] Make SDK tracer span current around gRPC newCall --- .../tech/ydb/core/impl/BaseGrpcTransport.java | 45 +++-- .../java/tech/ydb/query/impl/SessionImpl.java | 158 ++++++++++-------- .../java/tech/ydb/query/impl/SessionPool.java | 34 ++-- .../tech/ydb/query/impl/TableClientImpl.java | 90 +++++----- 4 files changed, 175 insertions(+), 152 deletions(-) diff --git a/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java b/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java index 3ddb8e9c3..eaf2ff069 100644 --- a/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java +++ b/core/src/main/java/tech/ydb/core/impl/BaseGrpcTransport.java @@ -37,7 +37,6 @@ import tech.ydb.core.impl.call.UnaryCall; import tech.ydb.core.impl.pool.EndpointRecord; import tech.ydb.core.impl.pool.GrpcChannel; -import tech.ydb.core.tracing.Scope; import tech.ydb.core.tracing.Span; /** @@ -125,19 +124,16 @@ public CompletableFuture> unaryCall( return CompletableFuture.completedFuture(deadlineExpiredResult(method, settings)); } - // Make SDK span current so gRPC ClientInterceptors nest under it at newCall. - try (Scope ignored = settings.getSpan().makeCurrent()) { - ClientCall call = channel.getReadyChannel().newCall(method, options); - ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); + ClientCall call = channel.getReadyChannel().newCall(method, options); + ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); - if (logger.isTraceEnabled()) { - logger.trace("UnaryCall[{}] with method {} and endpoint {} created", - traceId, method.getFullMethodName(), endpoint.getHostAndPort()); - } - Metadata metadata = makeMetadataFromSettings(settings, endpoint); - return new UnaryCall<>(traceId, endpoint.getHostAndPort(), call, handler, settings.getSpan()) - .startCall(request, metadata); + if (logger.isTraceEnabled()) { + logger.trace("UnaryCall[{}] with method {} and endpoint {} created", + traceId, method.getFullMethodName(), endpoint.getHostAndPort()); } + Metadata metadata = makeMetadataFromSettings(settings, endpoint); + return new UnaryCall<>(traceId, endpoint.getHostAndPort(), call, handler, settings.getSpan()) + .startCall(request, metadata); } catch (UnexpectedResultException ex) { logger.warn("UnaryCall[{}] got unexpected status {}", traceId, ex.getStatus()); return CompletableFuture.completedFuture(Result.fail(ex)); @@ -167,22 +163,19 @@ public GrpcReadStream readStreamCall( return new EmptyStream<>(deadlineExpiredStatus(method, settings)); } - // Make SDK span current so gRPC ClientInterceptors nest under it at newCall. - try (Scope ignored = settings.getSpan().makeCurrent()) { - ClientCall call = channel.getReadyChannel().newCall(method, options); - ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); - - if (logger.isTraceEnabled()) { - logger.trace("ReadStreamCall[{}] with method {} and endpoint {} created", - traceId, method.getFullMethodName(), endpoint.getHostAndPort() - ); - } + ClientCall call = channel.getReadyChannel().newCall(method, options); + ChannelStatusHandler handler = new ChannelStatusHandler(channel, settings); - Metadata metadata = makeMetadataFromSettings(settings, endpoint); - GrpcFlowControl flowCtrl = settings.getFlowControl(); - return new ReadStreamCall<>(traceId, endpoint.getHostAndPort(), call, flowCtrl, request, metadata, - handler, settings.getSpan()); + if (logger.isTraceEnabled()) { + logger.trace("ReadStreamCall[{}] with method {} and endpoint {} created", + traceId, method.getFullMethodName(), endpoint.getHostAndPort() + ); } + + Metadata metadata = makeMetadataFromSettings(settings, endpoint); + GrpcFlowControl flowCtrl = settings.getFlowControl(); + return new ReadStreamCall<>(traceId, endpoint.getHostAndPort(), call, flowCtrl, request, metadata, handler, + settings.getSpan()); } catch (UnexpectedResultException ex) { logger.warn("ReadStreamCall[{}] got unexpected status {}", traceId, ex.getStatus()); return new EmptyStream<>(ex.getStatus()); diff --git a/query/src/main/java/tech/ydb/query/impl/SessionImpl.java b/query/src/main/java/tech/ydb/query/impl/SessionImpl.java index 3334f3e05..1f3cbaa54 100644 --- a/query/src/main/java/tech/ydb/query/impl/SessionImpl.java +++ b/query/src/main/java/tech/ydb/query/impl/SessionImpl.java @@ -26,6 +26,7 @@ import tech.ydb.core.grpc.GrpcRequestSettings; import tech.ydb.core.operation.StatusExtractor; import tech.ydb.core.settings.BaseRequestSettings; +import tech.ydb.core.tracing.Scope; import tech.ydb.core.tracing.Span; import tech.ydb.core.utils.URITools; import tech.ydb.core.utils.UpdatableOptional; @@ -126,10 +127,14 @@ public CompletableFuture> beginTransaction(TxMode tx, B .setTxSettings(TxControl.txSettings(tx)) .build(); - return rpc.beginTransaction(request, makeOptions(settings).build()).thenApply(result -> { - updateSessionState(result.getStatus()); - return result.map(resp -> updateTransaction(new TransactionImpl(tx, resp.getTxMeta().getId()))); - }); + Span span = startSpan("ydb.BeginTransaction"); + try (Scope ignored = span.makeCurrent()) { + return Span.endOnResult(span, rpc.beginTransaction(request, makeOptions(settings, span).build())) + .thenApply(result -> { + updateSessionState(result.getStatus()); + return result.map(resp -> updateTransaction(new TransactionImpl(tx, resp.getTxMeta().getId()))); + }); + } } private QueryTransaction updateTransaction(TransactionImpl newTx) { @@ -331,14 +336,16 @@ GrpcReadStream createGrpcStream( public QueryStream createQuery(String query, TxMode tx, Params prms, ExecuteQuerySettings settings) { YdbQuery.TransactionControl tc = TxControl.txModeCtrl(tx, true); Span span = startSpan("ydb.ExecuteQuery"); - return new StreamImpl(createGrpcStream(query, tc, prms, settings, span), span) { - @Override - void handleTxMeta(String txID) { - if (txID != null && !txID.isEmpty()) { - logger.warn("{} got unexpected transaction id {}", SessionImpl.this, txID); + try (Scope ignored = span.makeCurrent()) { + return new StreamImpl(createGrpcStream(query, tc, prms, settings, span), span) { + @Override + void handleTxMeta(String txID) { + if (txID != null && !txID.isEmpty()) { + logger.warn("{} got unexpected transaction id {}", SessionImpl.this, txID); + } } - } - }; + }; + } } public CompletableFuture> delete(DeleteSessionSettings settings) { @@ -478,44 +485,46 @@ public QueryStream createQuery(String query, boolean commitAtEnd, Params prms, E : TxControl.txModeCtrl(txMode, commitAtEnd); Span span = startSpan("ydb.ExecuteQuery"); - return new StreamImpl(createGrpcStream(query, tc, prms, settings, span), span) { - @Override - void handleTxMeta(String txID) { - String newId = txID == null || txID.isEmpty() ? null : txID; - if (!txId.compareAndSet(currentId, newId)) { - logger.warn("{} lost transaction meta id {}", SessionImpl.this, newId); + try (Scope ignored = span.makeCurrent()) { + return new StreamImpl(createGrpcStream(query, tc, prms, settings, span), span) { + @Override + void handleTxMeta(String txID) { + String newId = txID == null || txID.isEmpty() ? null : txID; + if (!txId.compareAndSet(currentId, newId)) { + logger.warn("{} lost transaction meta id {}", SessionImpl.this, newId); + } } - } - @Override - void handleCompletion(Status status, Throwable th) { - if (th != null) { - currentStatusFuture.completeExceptionally( - new RuntimeException("Query on transaction failed with exception ", th)); - } - if (status.isSuccess()) { - if (commitAtEnd) { - currentStatusFuture.complete(Status.SUCCESS); + @Override + void handleCompletion(Status status, Throwable th) { + if (th != null) { + currentStatusFuture.completeExceptionally( + new RuntimeException("Query on transaction failed with exception ", th)); } - } else { - if (txId.compareAndSet(currentId, null)) { - logger.warn("{} transaction with id {} was failed", SessionImpl.this, currentId); + if (status.isSuccess()) { + if (commitAtEnd) { + currentStatusFuture.complete(Status.SUCCESS); + } + } else { + if (txId.compareAndSet(currentId, null)) { + logger.warn("{} transaction with id {} was failed", SessionImpl.this, currentId); + } + currentStatusFuture.complete(Status + .of(StatusCode.ABORTED) + .withIssues(Issue.of("Query on transaction failed with status " + + status, Issue.Severity.ERROR))); } - currentStatusFuture.complete(Status - .of(StatusCode.ABORTED) - .withIssues(Issue.of("Query on transaction failed with status " - + status, Issue.Severity.ERROR))); } - } - @Override - public void cancel() { - super.cancel(); - if (txId.compareAndSet(currentId, null)) { - logger.warn("{} transaction with id {} was cancelled", SessionImpl.this, currentId); + @Override + public void cancel() { + super.cancel(); + if (txId.compareAndSet(currentId, null)) { + logger.warn("{} transaction with id {} was cancelled", SessionImpl.this, currentId); + } } - } - }; + }; + } } @Override @@ -534,22 +543,25 @@ public CompletableFuture> commit(CommitTransactionSettings set .setTxId(transactionId) .build(); - return Span.endOnResult(span, rpc.commitTransaction(request, makeOptions(settings, span).build())) - .thenApply(res -> { - Status status = res.getStatus(); - currentStatusFuture.complete(status); - updateSessionState(status); - if (!txId.compareAndSet(transactionId, null)) { - logger.warn("{} lost commit response for transaction {}", SessionImpl.this, transactionId); - } - // TODO: CommitTransactionResponse must contain exec_stats - return res.map(resp -> new QueryInfo(null)); - }).whenComplete(((status, th) -> { - if (th != null) { - currentStatusFuture.completeExceptionally( - new RuntimeException("Transaction commit failed with exception", th)); - } - })); + try (Scope ignored = span.makeCurrent()) { + return Span.endOnResult(span, rpc.commitTransaction(request, makeOptions(settings, span).build())) + .thenApply(res -> { + Status status = res.getStatus(); + currentStatusFuture.complete(status); + updateSessionState(status); + if (!txId.compareAndSet(transactionId, null)) { + logger.warn("{} lost commit response for transaction {}", SessionImpl.this, + transactionId); + } + // TODO: CommitTransactionResponse must contain exec_stats + return res.map(resp -> new QueryInfo(null)); + }).whenComplete(((status, th) -> { + if (th != null) { + currentStatusFuture.completeExceptionally( + new RuntimeException("Transaction commit failed with exception", th)); + } + })); + } } @Override @@ -568,20 +580,22 @@ public CompletableFuture rollback(RollbackTransactionSettings settings) .setSessionId(sessionId) .setTxId(transactionId) .build(); - return Span.endOnResult(span, rpc.rollbackTransaction(request, makeOptions(settings, span).build())) - .thenApply(result -> { - updateSessionState(result.getStatus()); - if (!txId.compareAndSet(transactionId, null)) { - logger.warn("{} lost rollback response for transaction {}", SessionImpl.this, - transactionId); - } - return result.getStatus(); - }) - .whenComplete((status, th) -> { - currentStatusFuture.complete(Status - .of(StatusCode.ABORTED) - .withIssues(Issue.of("Transaction was rolled back", Issue.Severity.ERROR))); - }); + try (Scope ignored = span.makeCurrent()) { + return Span.endOnResult(span, rpc.rollbackTransaction(request, makeOptions(settings, span).build())) + .thenApply(result -> { + updateSessionState(result.getStatus()); + if (!txId.compareAndSet(transactionId, null)) { + logger.warn("{} lost rollback response for transaction {}", SessionImpl.this, + transactionId); + } + return result.getStatus(); + }) + .whenComplete((status, th) -> { + currentStatusFuture.complete(Status + .of(StatusCode.ABORTED) + .withIssues(Issue.of("Transaction was rolled back", Issue.Severity.ERROR))); + }); + } } } } diff --git a/query/src/main/java/tech/ydb/query/impl/SessionPool.java b/query/src/main/java/tech/ydb/query/impl/SessionPool.java index a7743c079..497fb4fa0 100644 --- a/query/src/main/java/tech/ydb/query/impl/SessionPool.java +++ b/query/src/main/java/tech/ydb/query/impl/SessionPool.java @@ -23,6 +23,7 @@ import tech.ydb.core.UnexpectedResultException; import tech.ydb.core.grpc.GrpcReadStream; import tech.ydb.core.metrics.Meter; +import tech.ydb.core.tracing.Scope; import tech.ydb.core.tracing.Span; import tech.ydb.core.utils.FutureTools; import tech.ydb.proto.query.YdbQuery; @@ -286,20 +287,25 @@ public CompletableFuture create() { long startNanos = System.nanoTime(); stats.requested.increment(); metrics.onSessionRequested(); - return Span.endOnResult(createSpan, SessionImpl.createSession(rpc, CREATE_SETTINGS, true, createSpan)) - .thenCompose(r -> { - metrics.onCreateTime(System.nanoTime() - startNanos); - - if (!r.isSuccess()) { - stats.failed.increment(); - metrics.onSessionFailed(r.getStatus()); - throw new UnexpectedResultException("create session problem", r.getStatus()); - } - metrics.onSessionCreated(); - PooledQuerySession session = new PooledQuerySession(rpc, r.getValue()); - return session.start(); - }) - .thenApply(Result::getValue); + try (Scope ignored = createSpan.makeCurrent()) { + return Span.endOnResult( + createSpan, + SessionImpl.createSession(rpc, CREATE_SETTINGS, true, createSpan) + ) + .thenCompose(r -> { + metrics.onCreateTime(System.nanoTime() - startNanos); + + if (!r.isSuccess()) { + stats.failed.increment(); + metrics.onSessionFailed(r.getStatus()); + throw new UnexpectedResultException("create session problem", r.getStatus()); + } + metrics.onSessionCreated(); + PooledQuerySession session = new PooledQuerySession(rpc, r.getValue()); + return session.start(); + }) + .thenApply(Result::getValue); + } } finally { ctx.detach(previous); } diff --git a/query/src/main/java/tech/ydb/query/impl/TableClientImpl.java b/query/src/main/java/tech/ydb/query/impl/TableClientImpl.java index 27ec97e6f..93401cb5b 100644 --- a/query/src/main/java/tech/ydb/query/impl/TableClientImpl.java +++ b/query/src/main/java/tech/ydb/query/impl/TableClientImpl.java @@ -15,6 +15,7 @@ import tech.ydb.core.StatusCode; import tech.ydb.core.UnexpectedResultException; import tech.ydb.core.grpc.GrpcTransport; +import tech.ydb.core.tracing.Scope; import tech.ydb.core.tracing.Span; import tech.ydb.core.tracing.Tracer; import tech.ydb.proto.ValueProtos; @@ -136,49 +137,54 @@ public CompletableFuture> executeDataQueryInternal( final List results = new ArrayList<>(); Span span = querySession.startSpan("ydb.ExecuteQuery"); - QueryStream stream = querySession.new StreamImpl(querySession.createGrpcStream(query, tc, prms, qs, span), - span) { - @Override - void handleTxMeta(String txID) { - txRef.set(txID); - } - }; - - CompletableFuture> future = stream.execute(new QueryStream.PartsHandler() { - @Override - public void onIssues(Issue[] issueArr) { - issues.addAll(Arrays.asList(issueArr)); - } - - @Override - public void onNextPart(QueryResultPart part) { - } // not used + try (Scope ignored = span.makeCurrent()) { + QueryStream stream = querySession.new StreamImpl( + querySession.createGrpcStream(query, tc, prms, qs, span), + span + ) { + @Override + void handleTxMeta(String txID) { + txRef.set(txID); + } + }; - @Override - public void onNextRawPart(long index, ValueProtos.ResultSet rs) { - int idx = (int) index; - while (results.size() <= idx) { - results.add(null); + CompletableFuture> future = stream.execute(new QueryStream.PartsHandler() { + @Override + public void onIssues(Issue[] issueArr) { + issues.addAll(Arrays.asList(issueArr)); } - if (results.get(idx) == null) { - results.set(idx, rs); - } else { - results.set(idx, results.get(idx).toBuilder().addAllRows(rs.getRowsList()).build()); + + @Override + public void onNextPart(QueryResultPart part) { + } // not used + + @Override + public void onNextRawPart(long index, ValueProtos.ResultSet rs) { + int idx = (int) index; + while (results.size() <= idx) { + results.add(null); + } + if (results.get(idx) == null) { + results.set(idx, rs); + } else { + results.set(idx, results.get(idx).toBuilder().addAllRows(rs.getRowsList()).build()); + } } - } - }); + }); - return future.thenApply(res -> { - if (!res.isSuccess()) { - return res.map(v -> null); - } - QueryStats stats = res.getValue().getStats(); - String txId = txRef.get(); - Status status = res.getStatus().withIssues(issues.toArray(new Issue[0])); - DataQueryResult value = new DataQueryResult(txId, results, stats != null ? stats.toProtobuf() : null); - return Result.success(value, status); - }); + return future.thenApply(res -> { + if (!res.isSuccess()) { + return res.map(v -> null); + } + QueryStats stats = res.getValue().getStats(); + String txId = txRef.get(); + Status status = res.getStatus().withIssues(issues.toArray(new Issue[0])); + DataQueryResult value = new DataQueryResult( + txId, results, stats != null ? stats.toProtobuf() : null); + return Result.success(value, status); + }); + } } @Override @@ -214,7 +220,9 @@ protected CompletableFuture commitTransactionInternal(String txId, Commi .withTraceId(settings.getTraceId()) .withRequestTimeout(settings.getTimeoutDuration()) .build(); - return Span.endOnStatus(span, querySession.commitById(txId, querySettings, span)); + try (Scope ignored = span.makeCurrent()) { + return Span.endOnStatus(span, querySession.commitById(txId, querySettings, span)); + } } @Override @@ -224,7 +232,9 @@ protected CompletableFuture rollbackTransactionInternal(String txId, Rol .withTraceId(settings.getTraceId()) .withRequestTimeout(settings.getTimeoutDuration()) .build(); - return Span.endOnStatus(span, querySession.rollbackById(txId, querySettings, span)); + try (Scope ignored = span.makeCurrent()) { + return Span.endOnStatus(span, querySession.rollbackById(txId, querySettings, span)); + } } private final class TracedTableTransaction implements TableTransaction { From 5395f3febad7fbbb9937aa19f74adcfbf50bba24 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 15:12:59 +0300 Subject: [PATCH 3/7] Test that GrpcTelemetry spans nest under SDK tracer spans --- ...nTelemetryQueryTracingIntegrationTest.java | 73 +++++++++++++++++++ 1 file changed, 73 insertions(+) diff --git a/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java b/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java index 1dc0bf4ba..9d2b68f3f 100644 --- a/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java +++ b/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java @@ -12,6 +12,7 @@ import io.opentelemetry.api.common.AttributeKey; import io.opentelemetry.api.trace.Tracer; import io.opentelemetry.context.Scope; +import io.opentelemetry.instrumentation.grpc.v1_6.GrpcTelemetry; import io.opentelemetry.sdk.OpenTelemetrySdk; import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; import io.opentelemetry.sdk.trace.SdkTracerProvider; @@ -33,6 +34,7 @@ import tech.ydb.core.UnexpectedResultException; import tech.ydb.core.grpc.GrpcTransport; import tech.ydb.core.tracing.OpenTelemetryTracer; +import tech.ydb.proto.query.v1.QueryServiceGrpc; import tech.ydb.query.QueryClient; import tech.ydb.query.QuerySession; import tech.ydb.query.QueryTransaction; @@ -74,9 +76,11 @@ public static void initTransport() { .setTracerProvider(tracerProvider) .build(); appTracer = openTelemetry.getTracer("test.app"); + GrpcTelemetry grpcTelemetry = GrpcTelemetry.create(openTelemetry); transport = GrpcTransport.forEndpoint(YDB.endpoint(), YDB.database()) .withAuthProvider(new TokenAuthProvider(YDB.authToken())) .withTracer(OpenTelemetryTracer.fromOpenTelemetry(openTelemetry)) + .addChannelInitializer(builder -> builder.intercept(grpcTelemetry.createClientInterceptor())) .build(); } @@ -140,6 +144,7 @@ public void queryClientSpansHaveBaseAttributes() { session.createQuery("SELECT 1", TxMode.NONE).execute().join().getStatus().expectSuccess(); assertSpanOK("ydb.ExecuteQuery", 1); + assertSpanOK("ydb.BeginTransaction", 0); assertSpanOK("ydb.Commit", 0); assertSpanOK("ydb.Rollback", 0); @@ -147,6 +152,7 @@ public void queryClientSpansHaveBaseAttributes() { txCommit.getStatus().expectSuccess(); assertSpanOK("ydb.ExecuteQuery", 1); + assertSpanOK("ydb.BeginTransaction", 1); assertSpanOK("ydb.Commit", 0); assertSpanOK("ydb.Rollback", 0); @@ -156,6 +162,7 @@ public void queryClientSpansHaveBaseAttributes() { commitResult.getStatus().expectSuccess(); assertSpanOK("ydb.ExecuteQuery", 2); + assertSpanOK("ydb.BeginTransaction", 1); assertSpanOK("ydb.Commit", 1); assertSpanOK("ydb.Rollback", 0); @@ -169,6 +176,7 @@ public void queryClientSpansHaveBaseAttributes() { assertSpanOK("ydb.CreateSession", 1); assertSpanOK("ydb.ExecuteQuery", 3); + assertSpanOK("ydb.BeginTransaction", 2); assertSpanOK("ydb.Commit", 1); assertSpanOK("ydb.Rollback", 1); } @@ -235,6 +243,71 @@ public void sdkSpanIsChildOfApplicationSpan() { } } + @Test + public void grpcInterceptorSpansNestUnderSdkSpans() { + try (QuerySession session = queryClient.createSession(Duration.ofSeconds(5)).join().getValue()) { + session.createQuery("SELECT 1", TxMode.NONE).execute().join().getStatus().expectSuccess(); + + Result txCommit = session.beginTransaction(TxMode.SERIALIZABLE_RW).join(); + txCommit.getStatus().expectSuccess(); + txCommit.getValue().commit().join().getStatus().expectSuccess(); + + Result txRollback = session.beginTransaction(TxMode.SERIALIZABLE_RW).join(); + txRollback.getStatus().expectSuccess(); + txRollback.getValue().rollback().join().expectSuccess(); + } + + List spans = exportedSpans(); + assertGrpcInterceptorChildOf("ydb.CreateSession", + QueryServiceGrpc.getCreateSessionMethod().getFullMethodName(), spans); + assertGrpcInterceptorChildOf("ydb.ExecuteQuery", + QueryServiceGrpc.getExecuteQueryMethod().getFullMethodName(), spans); + assertGrpcInterceptorChildOf("ydb.BeginTransaction", + QueryServiceGrpc.getBeginTransactionMethod().getFullMethodName(), spans); + assertGrpcInterceptorChildOf("ydb.Commit", + QueryServiceGrpc.getCommitTransactionMethod().getFullMethodName(), spans); + assertGrpcInterceptorChildOf("ydb.Rollback", + QueryServiceGrpc.getRollbackTransactionMethod().getFullMethodName(), spans); + } + + private static void assertGrpcInterceptorChildOf(String sdkSpanName, String grpcSpanName, List spans) { + int sdkCount = 0; + int nestedCount = 0; + for (SpanData sdkSpan : spans) { + if (!sdkSpanName.equals(sdkSpan.getName())) { + continue; + } + sdkCount++; + for (SpanData grpcSpan : spans) { + if (grpcSpanName.equals(grpcSpan.getName()) + && sdkSpan.getSpanId().equals(grpcSpan.getParentSpanId())) { + nestedCount++; + Assert.assertEquals(sdkSpan.getTraceId(), grpcSpan.getTraceId()); + } + } + } + Assert.assertTrue("Expected at least one " + sdkSpanName + ", spans=" + summarizeSpans(spans), sdkCount >= 1); + Assert.assertEquals("Each " + sdkSpanName + " must have nested " + grpcSpanName + + " interceptor span, spans=" + summarizeSpans(spans), sdkCount, nestedCount); + } + + private static String summarizeSpans(List spans) { + StringBuilder sb = new StringBuilder("["); + for (int i = 0; i < spans.size(); i++) { + SpanData span = spans.get(i); + if (i > 0) { + sb.append(", "); + } + sb.append(span.getName()) + .append("(id=") + .append(span.getSpanId()) + .append(", parent=") + .append(span.getParentSpanId()) + .append(')'); + } + return sb.append(']').toString(); + } + @Test public void retrySpanIsParentForRpcSpans() { AtomicInteger attempt = new AtomicInteger(); From 4b754faad85c544a25084106bdd5045ea8d8a2d0 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 15:15:07 +0300 Subject: [PATCH 4/7] Test that GrpcTelemetry spans nest under SDK tracer spans --- pom.xml | 2 ++ query/pom.xml | 6 ++++++ 2 files changed, 8 insertions(+) diff --git a/pom.xml b/pom.xml index 1aeabfa82..81a3d75eb 100644 --- a/pom.xml +++ b/pom.xml @@ -44,6 +44,8 @@ 17.0.0 1.59.0 + + 2.25.0-alpha diff --git a/query/pom.xml b/query/pom.xml index fd6263483..fc82caaf9 100644 --- a/query/pom.xml +++ b/query/pom.xml @@ -48,6 +48,12 @@ ${opentelemetry.version} test + + io.opentelemetry.instrumentation + opentelemetry-grpc-1.6 + ${opentelemetry.instrumentation.version} + test + junit junit From 5ddb0f16512cc5d6889d19482b0672c4a55c8009 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 15:38:56 +0300 Subject: [PATCH 5/7] Test that GrpcTelemetry spans nest under SDK tracer spans --- ...nTelemetryQueryTracingIntegrationTest.java | 23 +++---------------- 1 file changed, 3 insertions(+), 20 deletions(-) diff --git a/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java b/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java index 9d2b68f3f..f9d5add22 100644 --- a/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java +++ b/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java @@ -286,26 +286,9 @@ private static void assertGrpcInterceptorChildOf(String sdkSpanName, String grpc } } } - Assert.assertTrue("Expected at least one " + sdkSpanName + ", spans=" + summarizeSpans(spans), sdkCount >= 1); - Assert.assertEquals("Each " + sdkSpanName + " must have nested " + grpcSpanName - + " interceptor span, spans=" + summarizeSpans(spans), sdkCount, nestedCount); - } - - private static String summarizeSpans(List spans) { - StringBuilder sb = new StringBuilder("["); - for (int i = 0; i < spans.size(); i++) { - SpanData span = spans.get(i); - if (i > 0) { - sb.append(", "); - } - sb.append(span.getName()) - .append("(id=") - .append(span.getSpanId()) - .append(", parent=") - .append(span.getParentSpanId()) - .append(')'); - } - return sb.append(']').toString(); + Assert.assertTrue("Expected at least one " + sdkSpanName, sdkCount >= 1); + Assert.assertEquals("Each " + sdkSpanName + " must have nested " + grpcSpanName + " interceptor span", + sdkCount, nestedCount); } @Test From 8b9e7ca111df6a0bc789bd24ccb8f3ac9421dec9 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 15:53:06 +0300 Subject: [PATCH 6/7] rollback otel tests --- ...nTelemetryQueryTracingIntegrationTest.java | 56 ------------------- 1 file changed, 56 deletions(-) diff --git a/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java b/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java index f9d5add22..1dc0bf4ba 100644 --- a/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java +++ b/query/src/test/java/tech/ydb/query/opentelemetry/OpenTelemetryQueryTracingIntegrationTest.java @@ -12,7 +12,6 @@ import io.opentelemetry.api.common.AttributeKey; import io.opentelemetry.api.trace.Tracer; import io.opentelemetry.context.Scope; -import io.opentelemetry.instrumentation.grpc.v1_6.GrpcTelemetry; import io.opentelemetry.sdk.OpenTelemetrySdk; import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; import io.opentelemetry.sdk.trace.SdkTracerProvider; @@ -34,7 +33,6 @@ import tech.ydb.core.UnexpectedResultException; import tech.ydb.core.grpc.GrpcTransport; import tech.ydb.core.tracing.OpenTelemetryTracer; -import tech.ydb.proto.query.v1.QueryServiceGrpc; import tech.ydb.query.QueryClient; import tech.ydb.query.QuerySession; import tech.ydb.query.QueryTransaction; @@ -76,11 +74,9 @@ public static void initTransport() { .setTracerProvider(tracerProvider) .build(); appTracer = openTelemetry.getTracer("test.app"); - GrpcTelemetry grpcTelemetry = GrpcTelemetry.create(openTelemetry); transport = GrpcTransport.forEndpoint(YDB.endpoint(), YDB.database()) .withAuthProvider(new TokenAuthProvider(YDB.authToken())) .withTracer(OpenTelemetryTracer.fromOpenTelemetry(openTelemetry)) - .addChannelInitializer(builder -> builder.intercept(grpcTelemetry.createClientInterceptor())) .build(); } @@ -144,7 +140,6 @@ public void queryClientSpansHaveBaseAttributes() { session.createQuery("SELECT 1", TxMode.NONE).execute().join().getStatus().expectSuccess(); assertSpanOK("ydb.ExecuteQuery", 1); - assertSpanOK("ydb.BeginTransaction", 0); assertSpanOK("ydb.Commit", 0); assertSpanOK("ydb.Rollback", 0); @@ -152,7 +147,6 @@ public void queryClientSpansHaveBaseAttributes() { txCommit.getStatus().expectSuccess(); assertSpanOK("ydb.ExecuteQuery", 1); - assertSpanOK("ydb.BeginTransaction", 1); assertSpanOK("ydb.Commit", 0); assertSpanOK("ydb.Rollback", 0); @@ -162,7 +156,6 @@ public void queryClientSpansHaveBaseAttributes() { commitResult.getStatus().expectSuccess(); assertSpanOK("ydb.ExecuteQuery", 2); - assertSpanOK("ydb.BeginTransaction", 1); assertSpanOK("ydb.Commit", 1); assertSpanOK("ydb.Rollback", 0); @@ -176,7 +169,6 @@ public void queryClientSpansHaveBaseAttributes() { assertSpanOK("ydb.CreateSession", 1); assertSpanOK("ydb.ExecuteQuery", 3); - assertSpanOK("ydb.BeginTransaction", 2); assertSpanOK("ydb.Commit", 1); assertSpanOK("ydb.Rollback", 1); } @@ -243,54 +235,6 @@ public void sdkSpanIsChildOfApplicationSpan() { } } - @Test - public void grpcInterceptorSpansNestUnderSdkSpans() { - try (QuerySession session = queryClient.createSession(Duration.ofSeconds(5)).join().getValue()) { - session.createQuery("SELECT 1", TxMode.NONE).execute().join().getStatus().expectSuccess(); - - Result txCommit = session.beginTransaction(TxMode.SERIALIZABLE_RW).join(); - txCommit.getStatus().expectSuccess(); - txCommit.getValue().commit().join().getStatus().expectSuccess(); - - Result txRollback = session.beginTransaction(TxMode.SERIALIZABLE_RW).join(); - txRollback.getStatus().expectSuccess(); - txRollback.getValue().rollback().join().expectSuccess(); - } - - List spans = exportedSpans(); - assertGrpcInterceptorChildOf("ydb.CreateSession", - QueryServiceGrpc.getCreateSessionMethod().getFullMethodName(), spans); - assertGrpcInterceptorChildOf("ydb.ExecuteQuery", - QueryServiceGrpc.getExecuteQueryMethod().getFullMethodName(), spans); - assertGrpcInterceptorChildOf("ydb.BeginTransaction", - QueryServiceGrpc.getBeginTransactionMethod().getFullMethodName(), spans); - assertGrpcInterceptorChildOf("ydb.Commit", - QueryServiceGrpc.getCommitTransactionMethod().getFullMethodName(), spans); - assertGrpcInterceptorChildOf("ydb.Rollback", - QueryServiceGrpc.getRollbackTransactionMethod().getFullMethodName(), spans); - } - - private static void assertGrpcInterceptorChildOf(String sdkSpanName, String grpcSpanName, List spans) { - int sdkCount = 0; - int nestedCount = 0; - for (SpanData sdkSpan : spans) { - if (!sdkSpanName.equals(sdkSpan.getName())) { - continue; - } - sdkCount++; - for (SpanData grpcSpan : spans) { - if (grpcSpanName.equals(grpcSpan.getName()) - && sdkSpan.getSpanId().equals(grpcSpan.getParentSpanId())) { - nestedCount++; - Assert.assertEquals(sdkSpan.getTraceId(), grpcSpan.getTraceId()); - } - } - } - Assert.assertTrue("Expected at least one " + sdkSpanName, sdkCount >= 1); - Assert.assertEquals("Each " + sdkSpanName + " must have nested " + grpcSpanName + " interceptor span", - sdkCount, nestedCount); - } - @Test public void retrySpanIsParentForRpcSpans() { AtomicInteger attempt = new AtomicInteger(); From f7b543594f617baf2b0946518f0238e2edd306a8 Mon Sep 17 00:00:00 2001 From: KirillKurdyukov Date: Thu, 23 Jul 2026 16:42:40 +0300 Subject: [PATCH 7/7] rollback otel tests --- pom.xml | 2 -- query/pom.xml | 6 ------ 2 files changed, 8 deletions(-) diff --git a/pom.xml b/pom.xml index 81a3d75eb..1aeabfa82 100644 --- a/pom.xml +++ b/pom.xml @@ -44,8 +44,6 @@ 17.0.0 1.59.0 - - 2.25.0-alpha diff --git a/query/pom.xml b/query/pom.xml index fc82caaf9..fd6263483 100644 --- a/query/pom.xml +++ b/query/pom.xml @@ -48,12 +48,6 @@ ${opentelemetry.version} test - - io.opentelemetry.instrumentation - opentelemetry-grpc-1.6 - ${opentelemetry.instrumentation.version} - test - junit junit