diff --git a/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3MetricsExemplar.java b/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3MetricsExemplar.java new file mode 100644 index 000000000000..0ff3763e7fcc --- /dev/null +++ b/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3MetricsExemplar.java @@ -0,0 +1,289 @@ +/* + * Copyright 2026 Google LLC + * + * Redistribution and use in source and binary forms, with or without + * modification, are permitted provided that the following conditions are + * met: + * + * * Redistributions of source code must retain the above copyright + * notice, this list of conditions and the following disclaimer. + * * Redistributions in binary form must reproduce the above + * copyright notice, this list of conditions and the following disclaimer + * in the documentation and/or other materials provided with the + * distribution. + * * Neither the name of Google LLC nor the names of its + * contributors may be used to endorse or promote products derived from + * this software without specific prior written permission. + * + * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS + * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT + * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR + * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT + * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, + * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT + * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, + * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY + * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT + * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE + * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. + */ + +package com.google.showcase.v1beta1.it; + +import static com.google.common.truth.Truth.assertThat; + +import com.google.api.client.http.javanet.NetHttpTransport; +import com.google.api.gax.core.NoCredentialsProvider; +import com.google.api.gax.tracing.ApiTracerFactory; +import com.google.api.gax.tracing.CompositeTracerFactory; +import com.google.api.gax.tracing.OpenTelemetryMetricsFactory; +import com.google.api.gax.tracing.OpenTelemetryTracingFactory; +import com.google.showcase.v1beta1.EchoClient; +import com.google.showcase.v1beta1.EchoRequest; +import com.google.showcase.v1beta1.EchoSettings; +import com.google.showcase.v1beta1.stub.EchoStub; +import com.google.showcase.v1beta1.stub.EchoStubSettings; +import io.grpc.ManagedChannelBuilder; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.metrics.ExemplarFilter; +import io.opentelemetry.sdk.metrics.SdkMeterProvider; +import io.opentelemetry.sdk.metrics.data.ExemplarData; +import io.opentelemetry.sdk.metrics.data.HistogramPointData; +import io.opentelemetry.sdk.metrics.data.MetricData; +import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.data.SpanData; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; +import java.io.IOException; +import java.time.Duration; +import java.util.Arrays; +import java.util.Collection; +import java.util.List; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** Integration tests for Feature 3 (F3.1–F3.2): T3 Tracing and M3 Metrics Exemplar Correlation. */ +class ITOtelT3MetricsExemplar { + private static final String SHOWCASE_SERVER_ADDRESS = "localhost"; + private static final long SHOWCASE_SERVER_PORT = 7469; + private static final String SHOWCASE_GRPC_ENDPOINT = + String.format("%s:%s", SHOWCASE_SERVER_ADDRESS, SHOWCASE_SERVER_PORT); + private static final String SHOWCASE_HTTPJSON_ENDPOINT = + String.format("http://%s:%s", SHOWCASE_SERVER_ADDRESS, SHOWCASE_SERVER_PORT); + private static final String SHOWCASE_SERVICE_NAME = "showcase"; + + private InMemorySpanExporter spanExporter; + private InMemoryMetricReader metricReader; + private OpenTelemetrySdk openTelemetrySdk; + + @BeforeEach + void setUp() { + spanExporter = InMemorySpanExporter.create(); + metricReader = InMemoryMetricReader.create(); + + SdkTracerProvider tracerProvider = + SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) + .build(); + + SdkMeterProvider meterProvider = + SdkMeterProvider.builder() + .registerMetricReader(metricReader) + .setExemplarFilter(ExemplarFilter.traceBased()) + .build(); + + openTelemetrySdk = + OpenTelemetrySdk.builder() + .setTracerProvider(tracerProvider) + .setMeterProvider(meterProvider) + .build(); + } + + @AfterEach + void tearDown() { + if (openTelemetrySdk != null) { + openTelemetrySdk.close(); + } + } + + // F3.1: HTTP M3 metric records T3 span as exemplar + @Test + void testHttpJson_m3ExemplarMatchesT3Span() throws Exception { + // Verifies that for HTTP/JSON, client request duration metrics attach exemplars + // pointing directly to the overall T3 operation span (traceId and spanId match). + ApiTracerFactory compositeTracerFactory = createCompositeTracerFactory(); + EchoSettings settings = createEchoSettings(true); + EchoStub stub = createStubWithServiceName(settings, compositeTracerFactory); + + try (EchoClient client = EchoClient.create(stub)) { + client.echo(EchoRequest.newBuilder().setContent("exemplar-test-http").build()); + + List spans = waitAndCollectSpans(2); + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + Awaitility.await() + .atMost(Duration.ofSeconds(5)) + .untilAsserted( + () -> { + Collection metrics = metricReader.collectAllMetrics(); + MetricData durationMetric = + metrics.stream() + .filter(m -> m.getName().equals("gcp.client.request.duration")) + .findFirst() + .orElseThrow( + () -> new AssertionError("Duration metric not found in: " + metrics)); + + Collection points = + durationMetric.getHistogramData().getPoints(); + assertThat(points).hasSize(1); + HistogramPointData point = points.iterator().next(); + List exemplars = point.getExemplars(); + assertThat(exemplars).isNotEmpty(); + + ExemplarData exemplar = exemplars.get(0); + assertThat(exemplar.getSpanContext().getTraceId()).isEqualTo(t3Span.getTraceId()); + assertThat(exemplar.getSpanContext().getSpanId()).isEqualTo(t3Span.getSpanId()); + }); + } + } + + // F3.2: gRPC M3 metric records T3 span as exemplar + @Test + void testGrpc_m3ExemplarMatchesT3Span() throws Exception { + // Verifies that for gRPC, client request duration metrics attach exemplars + // pointing directly to the overall T3 operation span (traceId and spanId match). + ApiTracerFactory compositeTracerFactory = createCompositeTracerFactory(); + EchoSettings settings = createEchoSettings(false); + EchoStub stub = createStubWithServiceName(settings, compositeTracerFactory); + + try (EchoClient client = EchoClient.create(stub)) { + client.echo(EchoRequest.newBuilder().setContent("exemplar-test-grpc").build()); + + List spans = waitAndCollectSpans(2); + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + Awaitility.await() + .atMost(Duration.ofSeconds(5)) + .untilAsserted( + () -> { + Collection metrics = metricReader.collectAllMetrics(); + MetricData durationMetric = + metrics.stream() + .filter(m -> m.getName().equals("gcp.client.request.duration")) + .findFirst() + .orElseThrow( + () -> new AssertionError("Duration metric not found in: " + metrics)); + + Collection points = + durationMetric.getHistogramData().getPoints(); + assertThat(points).hasSize(1); + HistogramPointData point = points.iterator().next(); + List exemplars = point.getExemplars(); + assertThat(exemplars).isNotEmpty(); + + ExemplarData exemplar = exemplars.get(0); + assertThat(exemplar.getSpanContext().getTraceId()).isEqualTo(t3Span.getTraceId()); + assertThat(exemplar.getSpanContext().getSpanId()).isEqualTo(t3Span.getSpanId()); + }); + } + } + + /** + * Creates a composite tracer factory combining both OpenTelemetry tracing and metrics factories. + * + * @return the configured {@link CompositeTracerFactory} + */ + private CompositeTracerFactory createCompositeTracerFactory() { + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + OpenTelemetryMetricsFactory metricsFactory = new OpenTelemetryMetricsFactory(openTelemetrySdk); + return new CompositeTracerFactory(Arrays.asList(tracingFactory, metricsFactory)); + } + + /** + * Waits until the in-memory span exporter records at least {@code minSpans} completed spans. + * + * @param minSpans the minimum number of spans expected + * @return the list of completed {@link SpanData} items + */ + private List waitAndCollectSpans(int minSpans) { + Awaitility.await() + .atMost(Duration.ofSeconds(5)) + .until(() -> spanExporter.getFinishedSpanItems().size() >= minSpans); + return spanExporter.getFinishedSpanItems(); + } + + /** + * Constructs {@link EchoSettings} configured for the local Showcase test server. + * + * @param isHttpJson {@code true} for HTTP/JSON transport; {@code false} for gRPC transport + * @return the configured {@link EchoSettings} + * @throws Exception if transport provider initialization fails + */ + private EchoSettings createEchoSettings(boolean isHttpJson) throws Exception { + if (isHttpJson) { + return EchoSettings.newHttpJsonBuilder() + .setCredentialsProvider(NoCredentialsProvider.create()) + .setTransportChannelProvider( + EchoSettings.defaultHttpJsonTransportProviderBuilder() + .setHttpTransport(new NetHttpTransport.Builder().build()) + .build()) + .setEndpoint(SHOWCASE_HTTPJSON_ENDPOINT) + .build(); + } else { + return EchoSettings.newBuilder() + .setCredentialsProvider(NoCredentialsProvider.create()) + .setTransportChannelProvider( + EchoSettings.defaultGrpcTransportProviderBuilder() + .setChannelConfigurator(ManagedChannelBuilder::usePlaintext) + .build()) + .setEndpoint(SHOWCASE_GRPC_ENDPOINT) + .build(); + } + } + + /** + * Instantiates an {@link EchoStub} with custom service name and tracer factory. + * + * @param settings the client settings to base the stub on + * @param tracerFactory the tracer factory to register with the stub + * @return the initialized {@link EchoStub} + * @throws IOException if stub creation fails + */ + private EchoStub createStubWithServiceName(EchoSettings settings, ApiTracerFactory tracerFactory) + throws IOException { + EchoStubSettings.Builder builder = + (EchoStubSettings.Builder) settings.getStubSettings().toBuilder(); + builder.setTracerFactory(tracerFactory); + return new ExtendedEchoStubSettings(builder).createStub(); + } + + /** Extended {@link EchoStubSettings} that overrides {@link #getServiceName()} for testing. */ + private static class ExtendedEchoStubSettings extends EchoStubSettings { + /** + * Constructs settings wrapping the specified builder. + * + * @param builder the settings builder + * @throws IOException if base settings construction fails + */ + protected ExtendedEchoStubSettings(EchoStubSettings.Builder builder) throws IOException { + super(builder); + } + + @Override + public String getServiceName() { + return SHOWCASE_SERVICE_NAME; + } + } +} diff --git a/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3T4Hierarchy.java b/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3T4Hierarchy.java new file mode 100644 index 000000000000..6d2773eb308f --- /dev/null +++ b/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3T4Hierarchy.java @@ -0,0 +1,576 @@ +/* + * Copyright 2026 Google LLC + * + * Redistribution and use in source and binary forms, with or without + * modification, are permitted provided that the following conditions are + * met: + * + * * Redistributions of source code must retain the above copyright + * notice, this list of conditions and the following disclaimer. + * * Redistributions in binary form must reproduce the above + * copyright notice, this list of conditions and the following disclaimer + * in the documentation and/or other materials provided with the + * distribution. + * * Neither the name of Google LLC nor the names of its + * contributors may be used to endorse or promote products derived from + * this software without specific prior written permission. + * + * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS + * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT + * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR + * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT + * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, + * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT + * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, + * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY + * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT + * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE + * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. + */ + +package com.google.showcase.v1beta1.it; + +import static com.google.common.truth.Truth.assertThat; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import com.google.api.client.http.javanet.NetHttpTransport; +import com.google.api.gax.core.NoCredentialsProvider; +import com.google.api.gax.retrying.RetrySettings; +import com.google.api.gax.rpc.ApiException; +import com.google.api.gax.rpc.StatusCode; +import com.google.api.gax.tracing.ObservabilityAttributes; +import com.google.api.gax.tracing.OpenTelemetryTracingFactory; +import com.google.common.collect.ImmutableSet; +import com.google.rpc.Code; +import com.google.rpc.Status; +import com.google.showcase.v1beta1.AttemptSequenceRequest; +import com.google.showcase.v1beta1.CreateSequenceRequest; +import com.google.showcase.v1beta1.Sequence; +import com.google.showcase.v1beta1.SequenceServiceClient; +import com.google.showcase.v1beta1.SequenceServiceSettings; +import com.google.showcase.v1beta1.it.util.TestClientInitializer; +import com.google.showcase.v1beta1.stub.SequenceServiceStubSettings; +import io.grpc.ManagedChannelBuilder; +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.data.SpanData; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; +import java.io.IOException; +import java.time.Duration; +import java.util.Comparator; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** Integration tests for Feature 2 (F2.1–F2.4): T3/T4 Span Hierarchy and Retry Aggregation. */ +class ITOtelT3T4Hierarchy { + private static final String SHOWCASE_SERVER_ADDRESS = "localhost"; + private static final long SHOWCASE_SERVER_PORT = 7469; + private static final String SHOWCASE_GRPC_ENDPOINT = + String.format("%s:%s", SHOWCASE_SERVER_ADDRESS, SHOWCASE_SERVER_PORT); + private static final String SHOWCASE_HTTPJSON_ENDPOINT = + String.format("http://%s:%s", SHOWCASE_SERVER_ADDRESS, SHOWCASE_SERVER_PORT); + private static final String SHOWCASE_SERVICE_NAME = "showcase"; + + private InMemorySpanExporter spanExporter; + private OpenTelemetrySdk openTelemetrySdk; + private SequenceServiceClient setupGrpcClient; + private SequenceServiceClient setupHttpJsonClient; + + @BeforeEach + void setUp() throws Exception { + spanExporter = InMemorySpanExporter.create(); + SdkTracerProvider tracerProvider = + SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) + .build(); + openTelemetrySdk = OpenTelemetrySdk.builder().setTracerProvider(tracerProvider).build(); + + setupGrpcClient = TestClientInitializer.createGrpcSequenceClient(); + setupHttpJsonClient = TestClientInitializer.createHttpJsonSequenceClient(); + } + + @AfterEach + void tearDown() { + if (setupGrpcClient != null) { + setupGrpcClient.close(); + } + if (setupHttpJsonClient != null) { + setupHttpJsonClient.close(); + } + if (openTelemetrySdk != null) { + openTelemetrySdk.close(); + } + } + + // F2.1: HTTP T3/T4 retry succeeds (1 T3 span, 2 T4 child spans) + @Test + void testHttpJson_retrySucceeds() throws Exception { + // Verifies HTTP/JSON retry behavior: transient failure on attempt 0 is retried and succeeds on + // attempt 1. + // Asserts: + // 1. Exactly one overall INTERNAL operation span (T3). + // 2. Both attempt spans (T4) have parent_span_id == T3.span_id. + // 3. T3 operation span aggregates the successful status (OK) and 200 HTTP code. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + + // Sequence: attempt 1 -> UNAVAILABLE, attempt 2 -> OK + Sequence sequence = + Sequence.newBuilder() + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.OK.getNumber()).build()) + .build()) + .build(); + Sequence createdSequence = + setupHttpJsonClient.createSequence( + CreateSequenceRequest.newBuilder().setSequence(sequence).build()); + + // Flush spans created by setup client call + spanExporter.reset(); + + RetrySettings retrySettings = + RetrySettings.newBuilder() + .setInitialRetryDelayDuration(Duration.ofMillis(50L)) + .setRetryDelayMultiplier(1.5) + .setMaxRetryDelayDuration(Duration.ofMillis(200L)) + .setInitialRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setMaxRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setTotalTimeoutDuration(Duration.ofMillis(3000L)) + .setMaxAttempts(3) + .build(); + + try (SequenceServiceClient client = + createSequenceClient( + true, tracingFactory, retrySettings, ImmutableSet.of(StatusCode.Code.UNAVAILABLE))) { + client.attemptSequence( + AttemptSequenceRequest.newBuilder().setName(createdSequence.getName()).build()); + + List spans = waitAndCollectSpans(3); + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + List t4Spans = + spans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .sorted(Comparator.comparingLong(SpanData::getStartEpochNanos)) + .collect(Collectors.toList()); + + assertThat(t4Spans).hasSize(2); + + // Verify T3/T4 Parentage: both T4 child spans have parent_span_id == T3.span_id + for (SpanData t4 : t4Spans) { + assertThat(t4.getParentSpanId()).isEqualTo(t3Span.getSpanId()); + } + + // Verify T3 retry aggregation on success + assertThat(t3Span.getStatus().getStatusCode()) + .isEqualTo(io.opentelemetry.api.trace.StatusCode.UNSET); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNull(); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("OK"); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.longKey(ObservabilityAttributes.HTTP_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo(200L); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.HTTP_URL_TEMPLATE_ATTRIBUTE))) + .isNotEmpty(); + } + } + + // F2.2: HTTP T3/T4 retries exhausted + @Test + void testHttpJson_retriesExhausted() throws Exception { + // Verifies HTTP/JSON behavior when retries are exhausted after repeated failures. + // Asserts: + // 1. Exactly one overall INTERNAL operation span (T3). + // 2. All attempt spans (T4) are linked to the T3 span as parent. + // 3. T3 operation span aggregates ERROR status, 503 HTTP status, and UNAVAILABLE RPC status. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + + // Sequence: 3 UNAVAILABLE responses + Sequence sequence = + Sequence.newBuilder() + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .build(); + Sequence createdSequence = + setupHttpJsonClient.createSequence( + CreateSequenceRequest.newBuilder().setSequence(sequence).build()); + + spanExporter.reset(); + + // maxAttempts = 2 will exhaust retries after 2 attempts + RetrySettings retrySettings = + RetrySettings.newBuilder() + .setInitialRetryDelayDuration(Duration.ofMillis(50L)) + .setRetryDelayMultiplier(1.5) + .setMaxRetryDelayDuration(Duration.ofMillis(200L)) + .setInitialRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setMaxRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setTotalTimeoutDuration(Duration.ofMillis(3000L)) + .setMaxAttempts(2) + .build(); + + try (SequenceServiceClient client = + createSequenceClient( + true, tracingFactory, retrySettings, ImmutableSet.of(StatusCode.Code.UNAVAILABLE))) { + assertThrows( + ApiException.class, + () -> + client.attemptSequence( + AttemptSequenceRequest.newBuilder().setName(createdSequence.getName()).build())); + + List spans = waitAndCollectSpans(3); + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + List t4Spans = + spans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .sorted(Comparator.comparingLong(SpanData::getStartEpochNanos)) + .collect(Collectors.toList()); + + assertThat(t4Spans).hasSize(2); + + // Verify T3/T4 Parentage: all T4 child spans have parent_span_id == T3.span_id + for (SpanData t4 : t4Spans) { + assertThat(t4.getParentSpanId()).isEqualTo(t3Span.getSpanId()); + } + + // Verify T3 attributes on retries exhausted + assertThat(t3Span.getStatus().getStatusCode()) + .isEqualTo(io.opentelemetry.api.trace.StatusCode.ERROR); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.longKey(ObservabilityAttributes.HTTP_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo(503L); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("UNAVAILABLE"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isNotEmpty(); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNotNull(); + } + } + + // F2.3: gRPC T3/T4 retry succeeds + @Test + void testGrpc_retrySucceeds() throws Exception { + // Verifies gRPC retry behavior: transient failure on attempt 0 is retried and succeeds on + // attempt 1. + // Asserts: + // 1. Exactly one overall INTERNAL operation span (T3). + // 2. Both attempt spans (T4) have parent_span_id == T3.span_id. + // 3. T3 operation span aggregates the successful status (OK). + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + + // Sequence: attempt 1 -> UNAVAILABLE, attempt 2 -> OK + Sequence sequence = + Sequence.newBuilder() + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.OK.getNumber()).build()) + .build()) + .build(); + Sequence createdSequence = + setupGrpcClient.createSequence( + CreateSequenceRequest.newBuilder().setSequence(sequence).build()); + + spanExporter.reset(); + + RetrySettings retrySettings = + RetrySettings.newBuilder() + .setInitialRetryDelayDuration(Duration.ofMillis(50L)) + .setRetryDelayMultiplier(1.5) + .setMaxRetryDelayDuration(Duration.ofMillis(200L)) + .setInitialRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setMaxRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setTotalTimeoutDuration(Duration.ofMillis(3000L)) + .setMaxAttempts(3) + .build(); + + try (SequenceServiceClient client = + createSequenceClient( + false, tracingFactory, retrySettings, ImmutableSet.of(StatusCode.Code.UNAVAILABLE))) { + client.attemptSequence( + AttemptSequenceRequest.newBuilder().setName(createdSequence.getName()).build()); + + List spans = waitAndCollectSpans(3); + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + List t4Spans = + spans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .sorted(Comparator.comparingLong(SpanData::getStartEpochNanos)) + .collect(Collectors.toList()); + + assertThat(t4Spans).hasSize(2); + + // Verify T3/T4 Parentage: both T4 child spans have parent_span_id == T3.span_id + for (SpanData t4 : t4Spans) { + assertThat(t4.getParentSpanId()).isEqualTo(t3Span.getSpanId()); + } + + // Verify T3 retry aggregation on success + assertThat(t3Span.getStatus().getStatusCode()) + .isEqualTo(io.opentelemetry.api.trace.StatusCode.UNSET); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNull(); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("OK"); + } + } + + // F2.4: gRPC T3/T4 retries exhausted + @Test + void testGrpc_retriesExhausted() throws Exception { + // Verifies gRPC behavior when retries are exhausted after repeated failures. + // Asserts: + // 1. Exactly one overall INTERNAL operation span (T3). + // 2. All attempt spans (T4) are linked to the T3 span as parent. + // 3. T3 operation span aggregates ERROR status and UNAVAILABLE RPC status. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + + // Sequence: 3 UNAVAILABLE responses + Sequence sequence = + Sequence.newBuilder() + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .addResponses( + Sequence.Response.newBuilder() + .setStatus(Status.newBuilder().setCode(Code.UNAVAILABLE.getNumber()).build()) + .build()) + .build(); + Sequence createdSequence = + setupGrpcClient.createSequence( + CreateSequenceRequest.newBuilder().setSequence(sequence).build()); + + spanExporter.reset(); + + // maxAttempts = 2 will exhaust retries after 2 attempts + RetrySettings retrySettings = + RetrySettings.newBuilder() + .setInitialRetryDelayDuration(Duration.ofMillis(50L)) + .setRetryDelayMultiplier(1.5) + .setMaxRetryDelayDuration(Duration.ofMillis(200L)) + .setInitialRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setMaxRpcTimeoutDuration(Duration.ofMillis(1000L)) + .setTotalTimeoutDuration(Duration.ofMillis(3000L)) + .setMaxAttempts(2) + .build(); + + try (SequenceServiceClient client = + createSequenceClient( + false, tracingFactory, retrySettings, ImmutableSet.of(StatusCode.Code.UNAVAILABLE))) { + assertThrows( + ApiException.class, + () -> + client.attemptSequence( + AttemptSequenceRequest.newBuilder().setName(createdSequence.getName()).build())); + + List spans = waitAndCollectSpans(3); + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + List t4Spans = + spans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .sorted(Comparator.comparingLong(SpanData::getStartEpochNanos)) + .collect(Collectors.toList()); + + assertThat(t4Spans).hasSize(2); + + // Verify T3/T4 Parentage: all T4 child spans have parent_span_id == T3.span_id + for (SpanData t4 : t4Spans) { + assertThat(t4.getParentSpanId()).isEqualTo(t3Span.getSpanId()); + } + + // Verify T3 attributes on retries exhausted + assertThat(t3Span.getStatus().getStatusCode()) + .isEqualTo(io.opentelemetry.api.trace.StatusCode.ERROR); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("UNAVAILABLE"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isNotEmpty(); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNotNull(); + } + } + + /** + * Waits until the in-memory span exporter records at least {@code minSpans} completed spans. + * + * @param minSpans the minimum number of spans expected + * @return the list of completed {@link SpanData} items + */ + private List waitAndCollectSpans(int minSpans) { + Awaitility.await() + .atMost(Duration.ofSeconds(5)) + .until(() -> spanExporter.getFinishedSpanItems().size() >= minSpans); + return spanExporter.getFinishedSpanItems(); + } + + /** + * Constructs a {@link SequenceServiceClient} configured with custom retry settings and tracing. + * + * @param isHttpJson {@code true} for HTTP/JSON transport; {@code false} for gRPC transport + * @param tracingFactory the tracer factory to register with the client + * @param retrySettings the custom retry settings to apply to attemptSequence + * @param retryableCodes the set of status codes considered retryable + * @return the configured {@link SequenceServiceClient} + * @throws Exception if client initialization fails + */ + private SequenceServiceClient createSequenceClient( + boolean isHttpJson, + OpenTelemetryTracingFactory tracingFactory, + RetrySettings retrySettings, + Set retryableCodes) + throws Exception { + SequenceServiceSettings.Builder settingsBuilder = + isHttpJson + ? SequenceServiceSettings.newHttpJsonBuilder() + : SequenceServiceSettings.newBuilder(); + + settingsBuilder + .attemptSequenceSettings() + .setRetrySettings(retrySettings) + .setRetryableCodes(retryableCodes); + + settingsBuilder.setCredentialsProvider(NoCredentialsProvider.create()); + + if (isHttpJson) { + settingsBuilder + .setTransportChannelProvider( + SequenceServiceSettings.defaultHttpJsonTransportProviderBuilder() + .setHttpTransport(new NetHttpTransport.Builder().build()) + .build()) + .setEndpoint(SHOWCASE_HTTPJSON_ENDPOINT); + } else { + settingsBuilder + .setTransportChannelProvider( + SequenceServiceSettings.defaultGrpcTransportProviderBuilder() + .setChannelConfigurator(ManagedChannelBuilder::usePlaintext) + .build()) + .setEndpoint(SHOWCASE_GRPC_ENDPOINT); + } + + SequenceServiceStubSettings.Builder stubSettingsBuilder = + settingsBuilder.getStubSettingsBuilder(); + stubSettingsBuilder.setTracerFactory(tracingFactory); + return SequenceServiceClient.create( + new ExtendedSequenceServiceStubSettings(stubSettingsBuilder).createStub()); + } + + /** + * Extended {@link SequenceServiceStubSettings} that overrides {@link #getServiceName()} for + * testing. + */ + private static class ExtendedSequenceServiceStubSettings extends SequenceServiceStubSettings { + /** + * Constructs settings wrapping the specified builder. + * + * @param builder the settings builder + * @throws IOException if base settings construction fails + */ + protected ExtendedSequenceServiceStubSettings(SequenceServiceStubSettings.Builder builder) + throws IOException { + super(builder); + } + + @Override + public String getServiceName() { + return SHOWCASE_SERVICE_NAME; + } + } +} diff --git a/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3Tracing.java b/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3Tracing.java new file mode 100644 index 000000000000..d751a6f26720 --- /dev/null +++ b/java-showcase/gapic-showcase/src/test/java/com/google/showcase/v1beta1/it/ITOtelT3Tracing.java @@ -0,0 +1,534 @@ +/* + * Copyright 2026 Google LLC + * + * Redistribution and use in source and binary forms, with or without + * modification, are permitted provided that the following conditions are + * met: + * + * * Redistributions of source code must retain the above copyright + * notice, this list of conditions and the following disclaimer. + * * Redistributions in binary form must reproduce the above + * copyright notice, this list of conditions and the following disclaimer + * in the documentation and/or other materials provided with the + * distribution. + * * Neither the name of Google LLC nor the names of its + * contributors may be used to endorse or promote products derived from + * this software without specific prior written permission. + * + * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS + * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT + * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR + * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT + * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, + * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT + * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, + * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY + * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT + * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE + * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. + */ + +package com.google.showcase.v1beta1.it; + +import static com.google.common.truth.Truth.assertThat; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import com.google.api.client.http.javanet.NetHttpTransport; +import com.google.api.gax.core.NoCredentialsProvider; +import com.google.api.gax.rpc.ApiException; +import com.google.api.gax.tracing.ObservabilityAttributes; +import com.google.api.gax.tracing.OpenTelemetryTracingFactory; +import com.google.rpc.Code; +import com.google.rpc.Status; +import com.google.showcase.v1beta1.BlockRequest; +import com.google.showcase.v1beta1.BlockResponse; +import com.google.showcase.v1beta1.EchoClient; +import com.google.showcase.v1beta1.EchoRequest; +import com.google.showcase.v1beta1.EchoSettings; +import com.google.showcase.v1beta1.stub.EchoStub; +import com.google.showcase.v1beta1.stub.EchoStubSettings; +import io.grpc.ManagedChannelBuilder; +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.api.trace.StatusCode; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.data.SpanData; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; +import java.io.IOException; +import java.time.Duration; +import java.util.List; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** Integration tests for Feature 1 (F1.1–F1.8): Client-level T3 tracing across HTTP and gRPC. */ +class ITOtelT3Tracing { + private static final String SHOWCASE_SERVER_ADDRESS = "localhost"; + private static final long SHOWCASE_SERVER_PORT = 7469; + private static final String SHOWCASE_GRPC_ENDPOINT = + String.format("%s:%s", SHOWCASE_SERVER_ADDRESS, SHOWCASE_SERVER_PORT); + private static final String SHOWCASE_HTTPJSON_ENDPOINT = + String.format("http://%s:%s", SHOWCASE_SERVER_ADDRESS, SHOWCASE_SERVER_PORT); + private static final String SHOWCASE_SERVICE_NAME = "showcase"; + + private InMemorySpanExporter spanExporter; + private OpenTelemetrySdk openTelemetrySdk; + + @BeforeEach + void setUp() { + spanExporter = InMemorySpanExporter.create(); + SdkTracerProvider tracerProvider = + SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) + .build(); + openTelemetrySdk = OpenTelemetrySdk.builder().setTracerProvider(tracerProvider).build(); + } + + @AfterEach + void tearDown() { + if (openTelemetrySdk != null) { + openTelemetrySdk.close(); + } + } + + // F1.1: HTTP no traces emitted unless enabled. + @Test + void testTracingDisabled_httpjson() throws Exception { + // Verifies that when OpenTelemetry tracing is not configured on the client settings, + // no spans are recorded for HTTP/JSON calls. + EchoSettings settings = createEchoSettings(true); + try (EchoClient client = EchoClient.create(settings)) { + client.echo(EchoRequest.newBuilder().setContent("test-f1-1").build()); + List spans = spanExporter.getFinishedSpanItems(); + assertThat(spans).isEmpty(); + } + } + + // F1.2: HTTP T3 success case name and attributes conform to requirements. + @Test + void testT3Success_httpjson() throws Exception { + // Verifies that a successful HTTP/JSON call produces a T3 INTERNAL span with proper + // semantic convention attributes (http rpc.system, server.address, server.port, 200 + // http.response.status_code, + // url.template) and UNSET status. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + EchoSettings settings = createEchoSettings(true); + EchoStub stub = createStubWithServiceName(settings, tracingFactory); + + try (EchoClient client = EchoClient.create(stub)) { + client.echo(EchoRequest.newBuilder().setContent("test-f1-2").build()); + + List spans = waitAndCollectSpans(2); + assertThat(spans).isNotEmpty(); + + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + assertThat(t3Span.getKind()).isEqualTo(SpanKind.INTERNAL); + assertThat(t3Span.getName()).isNotEmpty(); + assertThat(t3Span.getStatus().getStatusCode()).isEqualTo(StatusCode.UNSET); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.RPC_SYSTEM_NAME_ATTRIBUTE))) + .isEqualTo("http"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.SERVER_ADDRESS_ATTRIBUTE))) + .isEqualTo(SHOWCASE_SERVER_ADDRESS); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.longKey(ObservabilityAttributes.SERVER_PORT_ATTRIBUTE))) + .isEqualTo(SHOWCASE_SERVER_PORT); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey(ObservabilityAttributes.GCP_CLIENT_SERVICE_ATTRIBUTE))) + .isEqualTo(SHOWCASE_SERVICE_NAME); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.longKey(ObservabilityAttributes.HTTP_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo(200L); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.HTTP_URL_TEMPLATE_ATTRIBUTE))) + .isEqualTo("v1beta1/echo:echo"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNull(); + } + } + + // F1.3: HTTP T3 server failures case name and attributes conform to requirements. + @Test + void testT3ServerFailure_httpjson() throws Exception { + // Verifies that a server-side error on HTTP/JSON produces a T3 INTERNAL span with ERROR status, + // 400 http.response.status_code, and error.type. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + EchoSettings settings = createEchoSettings(true); + EchoStub stub = createStubWithServiceName(settings, tracingFactory); + + try (EchoClient client = EchoClient.create(stub)) { + EchoRequest request = + EchoRequest.newBuilder() + .setError( + Status.newBuilder() + .setCode(Code.INVALID_ARGUMENT.getNumber()) + .setMessage("Server-side failure message") + .build()) + .build(); + + assertThrows(ApiException.class, () -> client.echo(request)); + + List spans = waitAndCollectSpans(2); + assertThat(spans).isNotEmpty(); + + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + assertThat(t3Span.getStatus().getStatusCode()).isEqualTo(StatusCode.ERROR); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.longKey(ObservabilityAttributes.HTTP_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo(400L); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isNotEmpty(); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNotNull(); + } + } + + // F1.4: HTTP T3 client failures case name and attributes conform to requirements. + @Test + void testT3ClientFailure_httpjson() throws Exception { + // Verifies that a client-side timeout on HTTP/JSON produces a T3 INTERNAL span with ERROR + // status, + // 504 http.response.status_code, and error.type. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + EchoSettings settings = createEchoSettings(true); + // Configure 1000ms timeout for blockCallable + EchoStubSettings.Builder builder = + (EchoStubSettings.Builder) settings.getStubSettings().toBuilder(); + builder.setTracerFactory(tracingFactory); + builder.blockSettings().setSimpleTimeoutNoRetriesDuration(Duration.ofMillis(1000L)); + EchoStub stub = new ExtendedEchoStubSettings(builder).createStub(); + + try (EchoClient client = EchoClient.create(stub)) { + BlockRequest request = + BlockRequest.newBuilder() + .setSuccess(BlockResponse.newBuilder().setContent("content").build()) + .setResponseDelay(com.google.protobuf.Duration.newBuilder().setSeconds(5).build()) + .build(); + + assertThrows(ApiException.class, () -> client.block(request)); + + List spans = waitAndCollectSpans(2); + assertThat(spans).isNotEmpty(); + + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + assertThat(t3Span.getStatus().getStatusCode()).isEqualTo(StatusCode.ERROR); + // In GAX, client timeout ApiException maps to HTTP 504 + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.longKey(ObservabilityAttributes.HTTP_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo(504L); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isNotEmpty(); + String errorType = + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE)); + assertThat(errorType).isNotNull(); + } + } + + // F1.5: gRPC no traces emitted unless enabled. + @Test + void testTracingDisabled_grpc() throws Exception { + // Verifies that when OpenTelemetry tracing is not configured on the client settings, + // no spans are recorded for gRPC calls. + EchoSettings settings = createEchoSettings(false); + try (EchoClient client = EchoClient.create(settings)) { + client.echo(EchoRequest.newBuilder().setContent("test-f1-5").build()); + List spans = spanExporter.getFinishedSpanItems(); + assertThat(spans).isEmpty(); + } + } + + // F1.6: gRPC T3 success case name and attributes conform to requirements. + @Test + void testT3Success_grpc() throws Exception { + // Verifies that a successful gRPC call produces a T3 INTERNAL span with proper + // semantic convention attributes (grpc rpc.system, server.address, server.port, OK + // rpc.response.status_code) + // and UNSET status. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + EchoSettings settings = createEchoSettings(false); + EchoStub stub = createStubWithServiceName(settings, tracingFactory); + + try (EchoClient client = EchoClient.create(stub)) { + client.echo(EchoRequest.newBuilder().setContent("test-f1-6").build()); + + List spans = waitAndCollectSpans(2); + assertThat(spans).isNotEmpty(); + + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + assertThat(t3Span.getKind()).isEqualTo(SpanKind.INTERNAL); + assertThat(t3Span.getName()).isEqualTo("google.showcase.v1beta1.Echo/Echo"); + assertThat(t3Span.getStatus().getStatusCode()).isEqualTo(StatusCode.UNSET); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.RPC_SYSTEM_NAME_ATTRIBUTE))) + .isEqualTo("grpc"); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("OK"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.SERVER_ADDRESS_ATTRIBUTE))) + .isEqualTo(SHOWCASE_SERVER_ADDRESS); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.longKey(ObservabilityAttributes.SERVER_PORT_ATTRIBUTE))) + .isEqualTo(SHOWCASE_SERVER_PORT); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey(ObservabilityAttributes.GCP_CLIENT_SERVICE_ATTRIBUTE))) + .isEqualTo(SHOWCASE_SERVICE_NAME); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNull(); + } + } + + // F1.7: gRPC T3 server failures case name and attributes conform to requirements. + @Test + void testT3ServerFailure_grpc() throws Exception { + // Verifies that a server-side error on gRPC produces a T3 INTERNAL span with ERROR status, + // INVALID_ARGUMENT rpc.response.status_code, and error.type. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + EchoSettings settings = createEchoSettings(false); + EchoStub stub = createStubWithServiceName(settings, tracingFactory); + + try (EchoClient client = EchoClient.create(stub)) { + EchoRequest request = + EchoRequest.newBuilder() + .setError( + Status.newBuilder() + .setCode(Code.INVALID_ARGUMENT.getNumber()) + .setMessage("Server-side failure message") + .build()) + .build(); + + assertThrows(ApiException.class, () -> client.echo(request)); + + List spans = waitAndCollectSpans(2); + assertThat(spans).isNotEmpty(); + + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + assertThat(t3Span.getStatus().getStatusCode()).isEqualTo(StatusCode.ERROR); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("INVALID_ARGUMENT"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isNotEmpty(); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNotNull(); + } + } + + // F1.8: gRPC T3 client failures case name and attributes conform to requirements. + @Test + void testT3ClientFailure_grpc() throws Exception { + // Verifies that a client-side timeout on gRPC produces a T3 INTERNAL span with ERROR status, + // DEADLINE_EXCEEDED rpc.response.status_code, and error.type. + OpenTelemetryTracingFactory tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk); + EchoSettings settings = createEchoSettings(false); + // Configure 1000ms timeout for blockCallable + EchoStubSettings.Builder builder = + (EchoStubSettings.Builder) settings.getStubSettings().toBuilder(); + builder.setTracerFactory(tracingFactory); + builder.blockSettings().setSimpleTimeoutNoRetriesDuration(Duration.ofMillis(1000L)); + EchoStub stub = new ExtendedEchoStubSettings(builder).createStub(); + + try (EchoClient client = EchoClient.create(stub)) { + BlockRequest request = + BlockRequest.newBuilder() + .setSuccess(BlockResponse.newBuilder().setContent("content").build()) + .setResponseDelay(com.google.protobuf.Duration.newBuilder().setSeconds(5).build()) + .build(); + + assertThrows(ApiException.class, () -> client.block(request)); + + List spans = waitAndCollectSpans(2); + assertThat(spans).isNotEmpty(); + + SpanData t3Span = + spans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("T3 INTERNAL span not found in: " + spans)); + + assertThat(t3Span.getStatus().getStatusCode()).isEqualTo(StatusCode.ERROR); + assertThat( + t3Span + .getAttributes() + .get( + AttributeKey.stringKey( + ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("DEADLINE_EXCEEDED"); + assertThat( + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isNotEmpty(); + String errorType = + t3Span + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE)); + assertThat(errorType).isNotNull(); + } + } + + /** + * Waits until the in-memory span exporter records at least {@code minSpans} completed spans. + * + * @param minSpans the minimum number of spans expected + * @return the list of completed {@link SpanData} items + */ + private List waitAndCollectSpans(int minSpans) { + Awaitility.await() + .atMost(Duration.ofSeconds(5)) + .until(() -> spanExporter.getFinishedSpanItems().size() >= minSpans); + return spanExporter.getFinishedSpanItems(); + } + + /** + * Constructs {@link EchoSettings} configured for the local Showcase test server. + * + * @param isHttpJson {@code true} for HTTP/JSON transport; {@code false} for gRPC transport + * @return the configured {@link EchoSettings} + * @throws Exception if transport provider initialization fails + */ + private EchoSettings createEchoSettings(boolean isHttpJson) throws Exception { + if (isHttpJson) { + return EchoSettings.newHttpJsonBuilder() + .setCredentialsProvider(NoCredentialsProvider.create()) + .setTransportChannelProvider( + EchoSettings.defaultHttpJsonTransportProviderBuilder() + .setHttpTransport(new NetHttpTransport.Builder().build()) + .build()) + .setEndpoint(SHOWCASE_HTTPJSON_ENDPOINT) + .build(); + } else { + return EchoSettings.newBuilder() + .setCredentialsProvider(NoCredentialsProvider.create()) + .setTransportChannelProvider( + EchoSettings.defaultGrpcTransportProviderBuilder() + .setChannelConfigurator(ManagedChannelBuilder::usePlaintext) + .build()) + .setEndpoint(SHOWCASE_GRPC_ENDPOINT) + .build(); + } + } + + /** + * Instantiates an {@link EchoStub} with custom service name and tracer factory. + * + * @param settings the client settings to base the stub on + * @param tracingFactory the tracer factory to register with the stub + * @return the initialized {@link EchoStub} + * @throws IOException if stub creation fails + */ + private EchoStub createStubWithServiceName( + EchoSettings settings, OpenTelemetryTracingFactory tracingFactory) throws IOException { + EchoStubSettings.Builder builder = + (EchoStubSettings.Builder) settings.getStubSettings().toBuilder(); + builder.setTracerFactory(tracingFactory); + return new ExtendedEchoStubSettings(builder).createStub(); + } + + /** Extended {@link EchoStubSettings} that overrides {@link #getServiceName()} for testing. */ + private static class ExtendedEchoStubSettings extends EchoStubSettings { + /** + * Constructs settings wrapping the specified builder. + * + * @param builder the settings builder + * @throws IOException if base settings construction fails + */ + protected ExtendedEchoStubSettings(EchoStubSettings.Builder builder) throws IOException { + super(builder); + } + + @Override + public String getServiceName() { + return SHOWCASE_SERVICE_NAME; + } + } +} diff --git a/sdk-platform-java/gax-java/dependencies.properties b/sdk-platform-java/gax-java/dependencies.properties index 6682d09eff89..135749d324d9 100644 --- a/sdk-platform-java/gax-java/dependencies.properties +++ b/sdk-platform-java/gax-java/dependencies.properties @@ -97,6 +97,7 @@ maven.io_opentelemetry_opentelemetry_sdk_testing=io.opentelemetry:opentelemetry- maven.io_opentelemetry_opentelemetry_sdk=io.opentelemetry:opentelemetry-sdk:1.57.0 maven.io_opentelemetry_opentelemetry_sdk_common=io.opentelemetry:opentelemetry-sdk-common:1.57.0 maven.io_opentelemetry_opentelemetry_sdk_metrics=io.opentelemetry:opentelemetry-sdk-metrics:1.57.0 +maven.io_opentelemetry_opentelemetry_sdk_trace=io.opentelemetry:opentelemetry-sdk-trace:1.57.0 maven.com_google_guava_guava_testlib=com.google.guava:guava-testlib:32.1.3-jre maven.org_awaitility_awaitility=org.awaitility:awaitility:4.3.0 diff --git a/sdk-platform-java/gax-java/gax/BUILD.bazel b/sdk-platform-java/gax-java/gax/BUILD.bazel index c98939ba7e18..a77329ede135 100644 --- a/sdk-platform-java/gax-java/gax/BUILD.bazel +++ b/sdk-platform-java/gax-java/gax/BUILD.bazel @@ -49,6 +49,7 @@ _TEST_COMPILE_DEPS = [ "@io_opentelemetry_opentelemetry_sdk//jar", "@io_opentelemetry_opentelemetry_sdk_metrics//jar", "@io_opentelemetry_opentelemetry_sdk_common//jar", + "@io_opentelemetry_opentelemetry_sdk_trace//jar", "@com_google_guava_guava_testlib//jar", "@org_awaitility_awaitility//jar", ] diff --git a/sdk-platform-java/gax-java/gax/pom.xml b/sdk-platform-java/gax-java/gax/pom.xml index 02c64480defb..8db559aabf12 100644 --- a/sdk-platform-java/gax-java/gax/pom.xml +++ b/sdk-platform-java/gax-java/gax/pom.xml @@ -104,6 +104,11 @@ opentelemetry-sdk-common test + + io.opentelemetry + opentelemetry-sdk-trace + test + org.junit.jupiter junit-jupiter-api diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/CompositeTracer.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/CompositeTracer.java index 118b2d89a334..e3c8461ae859 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/CompositeTracer.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/CompositeTracer.java @@ -31,7 +31,6 @@ import com.google.api.core.InternalApi; import com.google.common.collect.ImmutableList; -import java.util.ArrayList; import java.util.List; import java.util.Map; import org.jspecify.annotations.NullMarked; @@ -45,6 +44,8 @@ @NullMarked @InternalApi class CompositeTracer extends BaseApiTracer { + private static final Scope NO_OP_SCOPE = () -> {}; + private final List children; public CompositeTracer(List children) { @@ -53,60 +54,146 @@ public CompositeTracer(List children) { @Override public Scope inScope() { - final List childScopes = new ArrayList<>(children.size()); + if (children.isEmpty()) { + return NO_OP_SCOPE; + } + if (children.size() == 1) { + Scope scope = children.get(0).inScope(); + return scope != null ? scope : NO_OP_SCOPE; + } + + Scope[] childScopes = new Scope[children.size()]; + int scopeCount = 0; try { for (ApiTracer child : children) { - childScopes.add(child.inScope()); + Scope scope = child.inScope(); + if (scope != null) { + childScopes[scopeCount++] = scope; + } } - } catch (RuntimeException e) { - for (int i = childScopes.size() - 1; i >= 0; i--) { + if (scopeCount == 0) { + return NO_OP_SCOPE; + } + if (scopeCount == 1) { + return childScopes[0]; + } + return new CompositeScope(childScopes, scopeCount); + } catch (Throwable t) { + for (int i = scopeCount - 1; i >= 0; i--) { try { - childScopes.get(i).close(); - } catch (RuntimeException suppressed) { - e.addSuppressed(suppressed); + childScopes[i].close(); + } catch (Throwable suppressed) { + if (t != suppressed) { + t.addSuppressed(suppressed); + } } } - throw e; + throw throwException(t); } + } - return () -> { - RuntimeException exception = null; - for (int i = childScopes.size() - 1; i >= 0; i--) { - try { - childScopes.get(i).close(); - } catch (RuntimeException e) { - if (exception == null) { - exception = e; - } else { - exception.addSuppressed(e); + /** + * An aggregate {@link Scope} that encapsulates and closes multiple child scopes in reverse order. + */ + private static class CompositeScope implements Scope { + private final Scope[] scopes; + private final int count; + + /** + * Constructs a {@code CompositeScope} managing the given array of active child scopes. + * + * @param scopes the array containing active child scopes + * @param count the number of valid scopes in the array + */ + CompositeScope(Scope[] scopes, int count) { + this.scopes = scopes; + this.count = count; + } + + /** Closes all managed child scopes in reverse order, suppressing secondary exceptions. */ + @Override + public void close() { + Throwable firstException = null; + for (int i = count - 1; i >= 0; i--) { + Scope scope = scopes[i]; + if (scope != null) { + scopes[i] = null; + try { + scope.close(); + } catch (Throwable t) { + if (firstException == null) { + firstException = t; + } else if (firstException != t) { + firstException.addSuppressed(t); + } } } } - if (exception != null) { - throw exception; + if (firstException != null) { + throw throwException(firstException); } - }; + } + } + + /** + * Rethrows or wraps a {@link Throwable} without losing runtime exception or error fidelity. + * + * @param t the throwable to rethrow or wrap + * @return a {@link RuntimeException} wrapping {@code t} if {@code t} is a checked exception + */ + private static RuntimeException throwException(Throwable t) { + if (t instanceof RuntimeException) { + return (RuntimeException) t; + } + if (t instanceof Error) { + throw (Error) t; + } + return new RuntimeException(t); + } + + /** + * Enters the tracer's ambient scope safely, returning a no-op {@link Scope} if entering fails. + * This avoids allocating runnables, throwing exceptions, or returning null during lifecycle + * notifications. + * + *

Entering the ambient scope ensures active trace span context is present on the thread so + * that OpenTelemetry metric measurements recorded in lifecycle callbacks attach exemplars + * pointing to the active span. + */ + private Scope enterScope() { + try { + return inScope(); + } catch (RuntimeException e) { + // Ignore to prevent disrupting the lifecycle notification + return NO_OP_SCOPE; + } } @Override public void operationSucceeded() { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).operationSucceeded(); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).operationSucceeded(); + } } } @Override public void operationCancelled() { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).operationCancelled(); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).operationCancelled(); + } } } @Override public void operationFailed(Throwable error) { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).operationFailed(error); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).operationFailed(error); + } } } @@ -134,43 +221,55 @@ public void attemptStarted(Object request, int attemptNumber) { @Override public void attemptSucceeded() { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).attemptSucceeded(); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).attemptSucceeded(); + } } } @Override public void attemptCancelled() { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).attemptCancelled(); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).attemptCancelled(); + } } } @Override public void attemptFailed(Throwable error, org.threeten.bp.Duration delay) { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).attemptFailed(error, delay); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).attemptFailed(error, delay); + } } } @Override public void attemptFailedDuration(Throwable error, java.time.Duration delay) { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).attemptFailedDuration(error, delay); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).attemptFailedDuration(error, delay); + } } } @Override public void attemptFailedRetriesExhausted(Throwable error) { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).attemptFailedRetriesExhausted(error); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).attemptFailedRetriesExhausted(error); + } } } @Override public void attemptPermanentFailure(Throwable error) { - for (int i = children.size() - 1; i >= 0; i--) { - children.get(i).attemptPermanentFailure(error); + try (Scope s = enterScope()) { + for (int i = children.size() - 1; i >= 0; i--) { + children.get(i).attemptPermanentFailure(error); + } } } diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/OpenTelemetryTracingTracer.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/OpenTelemetryTracingTracer.java index cb41da9ccbeb..a4a87463c26d 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/OpenTelemetryTracingTracer.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/tracing/OpenTelemetryTracingTracer.java @@ -35,10 +35,14 @@ import io.opentelemetry.api.trace.Span; import io.opentelemetry.api.trace.SpanBuilder; import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.api.trace.StatusCode; import io.opentelemetry.api.trace.Tracer; +import io.opentelemetry.api.trace.propagation.W3CTraceContextPropagator; +import io.opentelemetry.context.Context; import java.util.HashMap; import java.util.Map; import java.util.concurrent.CancellationException; +import java.util.concurrent.locks.ReentrantLock; import org.jspecify.annotations.NullMarked; import org.jspecify.annotations.Nullable; @@ -51,15 +55,30 @@ class OpenTelemetryTracingTracer implements ApiTracer { private final Tracer tracer; private final Map attemptAttributes; private final String attemptSpanName; + private final String operationSpanName; private final ApiTracerContext apiTracerContext; - private @Nullable Span attemptSpan; + // Captures the active trace context from the calling thread at RPC initiation. + // This allows the operation span and attempt spans to link back to the caller's trace. + private final Context parentContext; + // Trace context containing the operationSpan, serving as the parent for attempt spans. + private final Context operationContext; + // Lock coordinates attempt transitions and operation completion across threads. + private final ReentrantLock lock = new ReentrantLock(); + private boolean operationCompleted; + // operationSpan and attemptSpan are volatile to ensure fresh reads for thread-safe snapshotting. + private volatile @Nullable Span operationSpan; + private volatile @Nullable Span attemptSpan; @Override - public void injectTraceContext(java.util.Map carrier) { - if (attemptSpan != null) { - io.opentelemetry.context.Context context = - io.opentelemetry.context.Context.current().with(attemptSpan); - io.opentelemetry.api.trace.propagation.W3CTraceContextPropagator.getInstance() + public void injectTraceContext(Map carrier) { + // Prefer the active attempt span (T4) so outgoing RPC wire context reflects the specific + // attempt; + // fall back to the overall operation span (T3) if no attempt is currently in-flight. + Span currentAttempt = attemptSpan; + Span spanToInject = currentAttempt != null ? currentAttempt : operationSpan; + if (spanToInject != null) { + Context context = Context.current().with(spanToInject); + W3CTraceContextPropagator.getInstance() .inject( context, carrier, @@ -71,6 +90,20 @@ public void injectTraceContext(java.util.Map carrier) { } } + @Override + @SuppressWarnings("MustBeClosedChecker") + public Scope inScope() { + // Attach the active attempt span to the current execution thread context; + // fall back to the overall operation span when between attempts. + Span currentAttempt = attemptSpan; + Span currentSpan = currentAttempt != null ? currentAttempt : operationSpan; + if (currentSpan == null) { + return () -> {}; + } + io.opentelemetry.context.Scope otelScope = currentSpan.makeCurrent(); + return otelScope::close; + } + /** * Creates a new instance of {@code OpenTelemetryTracingTracer}. * @@ -78,11 +111,7 @@ public void injectTraceContext(java.util.Map carrier) { * @param apiTracerContext the {@link ApiTracerContext} to use for recording spans */ OpenTelemetryTracingTracer(Tracer tracer, ApiTracerContext apiTracerContext) { - this.tracer = tracer; - this.apiTracerContext = apiTracerContext; - this.attemptSpanName = resolveAttemptSpanName(apiTracerContext); - this.attemptAttributes = new HashMap<>(); - buildAttributes(); + this(tracer, apiTracerContext, resolveAttemptSpanName(apiTracerContext)); } /** @@ -96,13 +125,72 @@ public void injectTraceContext(java.util.Map carrier) { @InternalApi OpenTelemetryTracingTracer( Tracer tracer, ApiTracerContext apiTracerContext, String attemptSpanName) { + this(tracer, apiTracerContext, attemptSpanName, resolveOperationSpanName(attemptSpanName)); + } + + /** + * Creates a new instance of {@code OpenTelemetryTracingTracer} with explicitly provided attempt + * and operation span names. + * + * @param tracer the {@link Tracer} to use for recording spans + * @param apiTracerContext the {@link ApiTracerContext} to use for recording spans + * @param attemptSpanName the name of the individual attempt spans + * @param operationSpanName the name of the overall client request operation span + */ + @InternalApi + OpenTelemetryTracingTracer( + Tracer tracer, + ApiTracerContext apiTracerContext, + String attemptSpanName, + String operationSpanName) { this.tracer = tracer; - this.attemptSpanName = attemptSpanName; this.apiTracerContext = apiTracerContext; + this.operationSpanName = operationSpanName; + this.attemptSpanName = attemptSpanName; this.attemptAttributes = new HashMap<>(); + this.parentContext = Context.current(); buildAttributes(); + this.operationSpan = startOperationSpan(); + this.operationContext = parentContext.with(this.operationSpan); } + /** + * Starts and initializes the operation-level client request span (T3). + * + * @return the newly started {@link Span} for the overall operation + */ + private Span startOperationSpan() { + SpanBuilder operationSpanBuilder = tracer.spanBuilder(operationSpanName); + operationSpanBuilder.setSpanKind(SpanKind.INTERNAL); + operationSpanBuilder.setParent(parentContext); + operationSpanBuilder.setAllAttributes( + ObservabilityUtils.toOtelAttributes(this.attemptAttributes)); + return operationSpanBuilder.startSpan(); + } + + /** + * Derives the operation-level span name from the attempt span name. + * + * @param attemptSpanName the attempt span name + * @return the operation span name + */ + private static String resolveOperationSpanName(String attemptSpanName) { + if (!Strings.isNullOrEmpty(attemptSpanName)) { + if (attemptSpanName.endsWith("/attempt")) { + String name = attemptSpanName.substring(0, attemptSpanName.length() - "/attempt".length()); + return name.isEmpty() ? "operation" : name; + } + return "attempt".equals(attemptSpanName) ? "operation" : attemptSpanName; + } + return "operation"; + } + + /** + * Resolves the canonical attempt-level span name based on transport and context. + * + * @param apiTracerContext the tracer context containing transport and method metadata + * @return the attempt span name + */ private static String resolveAttemptSpanName(ApiTracerContext apiTracerContext) { if (apiTracerContext.transport() == ApiTracerContext.Transport.GRPC) { // gRPC Uses the full method name as span name. @@ -118,34 +206,122 @@ private static String resolveAttemptSpanName(ApiTracerContext apiTracerContext) } } + /** Copies attempt-level attributes from the tracer context into the local attribute cache. */ private void buildAttributes() { this.attemptAttributes.putAll(this.apiTracerContext.getAttemptAttributes()); } @Override public void attemptStarted(Object request, int attemptNumber) { - Map currentAttemptAttributes = new HashMap<>(this.attemptAttributes); - - if (attemptNumber > 0) { - ApiTracerContext.Transport transport = apiTracerContext.transport(); - if (transport == ApiTracerContext.Transport.GRPC) { - currentAttemptAttributes.put( - ObservabilityAttributes.GRPC_RESEND_COUNT_ATTRIBUTE, (long) attemptNumber); - } else if (transport == ApiTracerContext.Transport.HTTP) { - currentAttemptAttributes.put( - ObservabilityAttributes.HTTP_RESEND_COUNT_ATTRIBUTE, (long) attemptNumber); + Span oldSpan = null; + lock.lock(); + try { + // Prevent creating new attempt spans if the overall operation has already concluded. + if (operationCompleted || operationSpan == null) { + return; + } + // If a previous attempt was not explicitly closed before a retry started, + // capture it so it can be ended cleanly outside the lock without blocking. + if (attemptSpan != null) { + oldSpan = attemptSpan; + attemptSpan = null; } + Map currentAttemptAttributes = new HashMap<>(this.attemptAttributes); + + if (attemptNumber > 0) { + ApiTracerContext.Transport transport = apiTracerContext.transport(); + if (transport == ApiTracerContext.Transport.GRPC) { + currentAttemptAttributes.put( + ObservabilityAttributes.GRPC_RESEND_COUNT_ATTRIBUTE, (long) attemptNumber); + } else if (transport == ApiTracerContext.Transport.HTTP) { + currentAttemptAttributes.put( + ObservabilityAttributes.HTTP_RESEND_COUNT_ATTRIBUTE, (long) attemptNumber); + } + } + + SpanBuilder spanBuilder = tracer.spanBuilder(attemptSpanName); + + // Attempt spans are of the CLIENT kind + spanBuilder.setSpanKind(SpanKind.CLIENT); + + // Link attempt span to operation context (parent T3 span) + spanBuilder.setParent(operationContext); + + // Pass the combined attributes to the new SpanBuilder method + spanBuilder.setAllAttributes(ObservabilityUtils.toOtelAttributes(currentAttemptAttributes)); + + this.attemptSpan = spanBuilder.startSpan(); + } finally { + lock.unlock(); + } + // End lingering previous attempt outside the lock to avoid holding the lock during callbacks. + if (oldSpan != null) { + endSpan(oldSpan, null); } + } - SpanBuilder spanBuilder = tracer.spanBuilder(attemptSpanName); + /** + * Signals that the overall logical operation succeeded. + * + *

Closes any remaining in-flight attempt span and ends the operation span. + */ + @Override + public void operationSucceeded() { + recordErrorAndEndOperation(null); + } - // Attempt spans are of the CLIENT kind - spanBuilder.setSpanKind(SpanKind.CLIENT); + /** + * Signals that the overall logical operation was cancelled. + * + *

Closes any remaining in-flight attempt span with a {@link CancellationException} and ends + * the operation span with an ERROR status. + */ + @Override + public void operationCancelled() { + recordErrorAndEndOperation(new CancellationException()); + } - // Pass the combined attributes to the new SpanBuilder method - spanBuilder.setAllAttributes(ObservabilityUtils.toOtelAttributes(currentAttemptAttributes)); + /** + * Signals that the overall logical operation failed permanently. + * + *

Closes any remaining in-flight attempt span and ends the operation span with the provided + * error details and an ERROR status. + * + * @param error the cause of the operation failure + */ + @Override + public void operationFailed(Throwable error) { + recordErrorAndEndOperation(error); + } - this.attemptSpan = spanBuilder.startSpan(); + /** + * Records error details and ends both the active attempt span and the operation span in a + * thread-safe manner. + * + * @param error the exception associated with the operation failure, or {@code null} if successful + */ + private void recordErrorAndEndOperation(@Nullable Throwable error) { + Span localOperationSpan; + Span localAttemptSpan; + lock.lock(); + try { + operationCompleted = true; + localOperationSpan = operationSpan; + if (localOperationSpan == null) { + return; + } + operationSpan = null; + localAttemptSpan = attemptSpan; + attemptSpan = null; + } finally { + lock.unlock(); + } + + if (localAttemptSpan != null) { + endSpan(localAttemptSpan, error); + } + + endSpan(localOperationSpan, error); } @Override @@ -154,13 +330,16 @@ public void attemptSucceeded() { } @Override - public void responseHeadersReceived(java.util.Map headers) { - if (attemptSpan == null) { + public void responseHeadersReceived(Map headers) { + // Snapshot to a local variable to prevent race conditions if another thread + // clears attemptSpan concurrently. + Span currentSpan = attemptSpan; + if (currentSpan == null) { return; } long contentLength = extractContentLength(headers); if (contentLength >= 0) { - attemptSpan.setAttribute(ObservabilityAttributes.HTTP_RESPONSE_BODY_SIZE, contentLength); + currentSpan.setAttribute(ObservabilityAttributes.HTTP_RESPONSE_BODY_SIZE, contentLength); } } @@ -215,42 +394,65 @@ public void attemptPermanentFailure(Throwable error) { recordErrorAndEndAttempt(error); } + /** + * Records error details and ends the current attempt span in a thread-safe manner. + * + * @param error the exception associated with the attempt failure, or {@code null} if successful + */ private void recordErrorAndEndAttempt(@Nullable Throwable error) { - if (attemptSpan == null) { - return; + Span localAttemptSpan; + lock.lock(); + try { + localAttemptSpan = attemptSpan; + if (localAttemptSpan == null) { + return; + } + attemptSpan = null; + } finally { + lock.unlock(); } + + endSpan(localAttemptSpan, error); + } + + /** + * Attaches response status attributes and error messages to the span and ends it. + * + *

This method runs outside of synchronization locks to avoid blocking threads during + * OpenTelemetry span completion callbacks. + * + * @param span the span to finish + * @param error the exception that caused the span to end, or {@code null} if successful + */ + private void endSpan(Span span, @Nullable Throwable error) { Map responseAttributes = ObservabilityUtils.getResponseAttributes(error, this.apiTracerContext.transport()); if (!responseAttributes.isEmpty()) { - attemptSpan.setAllAttributes(ObservabilityUtils.toOtelAttributes(responseAttributes)); + span.setAllAttributes(ObservabilityUtils.toOtelAttributes(responseAttributes)); } - if (error != null && !Strings.isNullOrEmpty(error.getMessage())) { - attemptSpan.setAttribute( - ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE, error.getMessage()); - } - - endAttempt(); - } - - private void endAttempt() { - if (attemptSpan == null) { - return; + if (error != null) { + span.setStatus(StatusCode.ERROR); + if (!Strings.isNullOrEmpty(error.getMessage())) { + span.setAttribute(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE, error.getMessage()); + } } - attemptSpan.end(); - attemptSpan = null; + span.end(); } @Override public void requestUrlResolved(String url) { - if (attemptSpan == null) { + // Snapshot to a local variable to prevent race conditions if another thread + // clears attemptSpan concurrently. + Span currentSpan = attemptSpan; + if (currentSpan == null) { return; } String sanitizedUrlString = ObservabilityUtils.sanitizeUrlFull(url); if (sanitizedUrlString.isEmpty()) { return; } - attemptSpan.setAttribute(ObservabilityAttributes.HTTP_URL_FULL_ATTRIBUTE, sanitizedUrlString); + currentSpan.setAttribute(ObservabilityAttributes.HTTP_URL_FULL_ATTRIBUTE, sanitizedUrlString); } } diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/CompositeTracerTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/CompositeTracerTest.java index 02e5012dc32c..37881959c3e4 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/CompositeTracerTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/CompositeTracerTest.java @@ -35,6 +35,7 @@ import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -125,6 +126,26 @@ void testInScope_childScopeCloseThrows() { inOrder.verify(scope1).close(); } + @Test + void testInScope_compositeScopeClose_isIdempotent() { + ApiTracer.Scope scope1 = + mock(ApiTracer.Scope.class, Mockito.withSettings().withoutAnnotations()); + ApiTracer.Scope scope2 = + mock(ApiTracer.Scope.class, Mockito.withSettings().withoutAnnotations()); + + when(child1.inScope()).thenReturn(scope1); + when(child2.inScope()).thenReturn(scope2); + + ApiTracer.Scope compositeScope = compositeTracer.inScope(); + + compositeScope.close(); + // Subsequent close should be idempotent and not invoke underlying scopes again + compositeScope.close(); + + verify(scope2, times(1)).close(); + verify(scope1, times(1)).close(); + } + @Test void testOperationSucceeded() { compositeTracer.operationSucceeded(); diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerFactoryTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerFactoryTest.java index 3c78cef6dbd7..4a58a3b07bfd 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerFactoryTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerFactoryTest.java @@ -75,6 +75,7 @@ void setUp() { when(openTelemetry.getTracer(anyString())).thenReturn(tracer); when(tracer.spanBuilder(anyString())).thenReturn(spanBuilder); when(spanBuilder.setSpanKind(any())).thenReturn(spanBuilder); + when(spanBuilder.setParent(any())).thenReturn(spanBuilder); when(spanBuilder.setAllAttributes(any(Attributes.class))).thenReturn(spanBuilder); when(spanBuilder.startSpan()).thenReturn(span); @@ -228,7 +229,7 @@ void testNewTracer_withContext_grpc_usesFullMethodName() { tracerInstance.attemptStarted(null, 1); - verify(tracer).spanBuilder("google.cloud.v1.Service/Method"); + verify(tracer, atLeastOnce()).spanBuilder("google.cloud.v1.Service/Method"); } @ParameterizedTest @@ -255,7 +256,7 @@ void testNewTracer_withContext_http_usesHttpMethodAndPathTemplate( tracerInstance.attemptStarted(null, 1); - verify(tracer).spanBuilder(expectedSpanName); + verify(tracer, atLeastOnce()).spanBuilder(expectedSpanName); } @Test @@ -273,7 +274,7 @@ void testNewTracer_withContext_http_noHttpMethodOrPathTemplate_usesFullMethodNam tracerInstance.attemptStarted(null, 1); - verify(tracer).spanBuilder("google.cloud.v1.Service.Method"); + verify(tracer, atLeastOnce()).spanBuilder("google.cloud.v1.Service.Method"); } @Test @@ -309,7 +310,7 @@ void testNewTracer_mergesFactoryContext() { tracerInstance.attemptStarted(null, 1); ArgumentCaptor attributesCaptor = ArgumentCaptor.forClass(Attributes.class); - verify(spanBuilder).setAllAttributes(attributesCaptor.capture()); + verify(spanBuilder, atLeastOnce()).setAllAttributes(attributesCaptor.capture()); Attributes attributes = attributesCaptor.getValue(); assertThat(attributes.asMap()) diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerIntegrationTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerIntegrationTest.java new file mode 100644 index 000000000000..b479a6a5eccb --- /dev/null +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerIntegrationTest.java @@ -0,0 +1,368 @@ +/* + * Copyright 2026 Google LLC + * + * Redistribution and use in source and binary forms, with or without + * modification, are permitted provided that the following conditions are + * met: + * + * * Redistributions of source code must retain the above copyright + * notice, this list of conditions and the following disclaimer. + * * Redistributions in binary form must reproduce the above + * copyright notice, this list of conditions and the following disclaimer + * in the documentation and/or other materials provided with the + * distribution. + * * Neither the name of Google LLC nor the names of its + * contributors may be used to endorse or promote products derived from + * this software without specific prior written permission. + * + * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS + * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT + * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR + * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT + * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, + * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT + * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, + * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY + * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT + * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE + * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. + */ +package com.google.api.gax.tracing; + +import static com.google.common.truth.Truth.assertThat; + +import com.google.api.gax.rpc.LibraryMetadata; +import com.google.api.gax.tracing.ApiTracerContext.Transport; +import com.google.api.gax.tracing.ApiTracerFactory.OperationType; +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.api.trace.SpanKind; +import io.opentelemetry.api.trace.Tracer; +import io.opentelemetry.context.Scope; +import io.opentelemetry.sdk.OpenTelemetrySdk; +import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter; +import io.opentelemetry.sdk.trace.SdkTracerProvider; +import io.opentelemetry.sdk.trace.data.SpanData; +import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; +import java.time.Duration; +import java.util.List; +import java.util.stream.Collectors; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class OpenTelemetryTracingTracerIntegrationTest { + + private static final String FULL_METHOD_NAME = "google.fake.v1.FakeService/FakeMethod"; + private static final LibraryMetadata LIBRARY_METADATA = + LibraryMetadata.newBuilder() + .setRepository("googleapis/google-cloud-java") + .setArtifactName("google-cloud-fake") + .setVersion("1.0.0") + .build(); + private static final ApiTracerContext TRACER_CONTEXT = + ApiTracerContext.newBuilder() + .setFullMethodName(FULL_METHOD_NAME) + .setTransport(Transport.GRPC) + .setLibraryMetadata(LIBRARY_METADATA) + .setOperationType(OperationType.Unary) + .build(); + + private InMemorySpanExporter spanExporter; + private SdkTracerProvider tracerProvider; + private OpenTelemetrySdk openTelemetrySdk; + private Tracer tracer; + private ApiTracerFactory tracingFactory; + + @BeforeEach + void setUp() { + spanExporter = InMemorySpanExporter.create(); + tracerProvider = + SdkTracerProvider.builder() + .addSpanProcessor(SimpleSpanProcessor.create(spanExporter)) + .build(); + openTelemetrySdk = OpenTelemetrySdk.builder().setTracerProvider(tracerProvider).build(); + tracer = openTelemetrySdk.getTracer("test-tracer"); + tracingFactory = new OpenTelemetryTracingFactory(openTelemetrySdk).withContext(TRACER_CONTEXT); + } + + @AfterEach + void tearDown() { + tracerProvider.close(); + } + + @Test + void testAttemptSpan_linkedToParentContextFromCallingThread() { + // Verifies that when an application trace parent exists, the operation span (T3) + // links to the application parent, and the attempt span (T4) links to the operation span. + Span parentSpan = tracer.spanBuilder("application-parent-operation").startSpan(); + ApiTracer apiTracer; + try (Scope scope = parentSpan.makeCurrent()) { + apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + } + + apiTracer.attemptStarted(new Object(), 0); + apiTracer.attemptSucceeded(); + apiTracer.operationSucceeded(); + parentSpan.end(); + + List finishedSpans = spanExporter.getFinishedSpanItems(); + assertThat(finishedSpans).hasSize(3); // root, operation, attempt + + SpanData attemptSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .findFirst() + .orElseThrow(() -> new AssertionError("Attempt span not found")); + SpanData operationSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("Operation span not found")); + SpanData rootSpan = + finishedSpans.stream() + .filter(s -> s.getName().equals("application-parent-operation")) + .findFirst() + .orElseThrow(() -> new AssertionError("Parent span not found")); + + assertThat(attemptSpan.getKind()).isEqualTo(SpanKind.CLIENT); + assertThat(attemptSpan.getParentSpanId()).isEqualTo(operationSpan.getSpanContext().getSpanId()); + assertThat(operationSpan.getParentSpanId()).isEqualTo(rootSpan.getSpanContext().getSpanId()); + assertThat(attemptSpan.getSpanContext().getTraceId()) + .isEqualTo(rootSpan.getSpanContext().getTraceId()); + assertThat(operationSpan.getSpanContext().getTraceId()) + .isEqualTo(rootSpan.getSpanContext().getTraceId()); + } + + @Test + void testAttemptSpan_withoutParentContext_hasNoParent() { + // Verifies that without an external parent context, the operation span acts as root + // and the attempt span is a child of the operation span. + ApiTracer apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + + apiTracer.attemptStarted(new Object(), 0); + apiTracer.attemptSucceeded(); + apiTracer.operationSucceeded(); + + List finishedSpans = spanExporter.getFinishedSpanItems(); + assertThat(finishedSpans).hasSize(2); // operation, attempt + + SpanData attemptSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .findFirst() + .orElseThrow(() -> new AssertionError("Attempt span not found")); + SpanData operationSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("Operation span not found")); + + assertThat(operationSpan.getParentSpanContext().isValid()).isFalse(); + assertThat(attemptSpan.getParentSpanId()).isEqualTo(operationSpan.getSpanContext().getSpanId()); + } + + @Test + void testSequentialAttempts_closesPreviousAttemptSpanAndLinksAllToParent() { + Span parentSpan = tracer.spanBuilder("application-parent-operation").startSpan(); + ApiTracer apiTracer; + try (Scope scope = parentSpan.makeCurrent()) { + apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + } + + // Start attempt 0 (e.g. transient failure without explicit endAttempt before retry) + apiTracer.attemptStarted(new Object(), 0); + + // Start attempt 1 - should automatically end attempt 0 + apiTracer.attemptStarted(new Object(), 1); + assertThat(spanExporter.getFinishedSpanItems()).hasSize(1); + SpanData attempt0Span = spanExporter.getFinishedSpanItems().get(0); + + // Complete attempt 1 and operation + apiTracer.attemptSucceeded(); + apiTracer.operationSucceeded(); + parentSpan.end(); + + List finishedSpans = spanExporter.getFinishedSpanItems(); + assertThat(finishedSpans).hasSize(4); // attempt 0, attempt 1, operation, parent + + SpanData operationSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("Operation span not found")); + + assertThat(attempt0Span.getParentSpanId()) + .isEqualTo(operationSpan.getSpanContext().getSpanId()); + + SpanData attempt1Span = + finishedSpans.stream() + .filter( + s -> + s.getKind() == SpanKind.CLIENT + && !s.getSpanContext() + .getSpanId() + .equals(attempt0Span.getSpanContext().getSpanId())) + .findFirst() + .orElseThrow(() -> new AssertionError("Attempt 1 span not found")); + + assertThat(attempt1Span.getParentSpanId()) + .isEqualTo(operationSpan.getSpanContext().getSpanId()); + assertThat(operationSpan.getParentSpanId()).isEqualTo(parentSpan.getSpanContext().getSpanId()); + assertThat(attempt1Span.getSpanContext().getTraceId()) + .isEqualTo(parentSpan.getSpanContext().getTraceId()); + } + + @Test + void testOperationFailed_endsActiveAttemptSpanWithErrorAttributes() { + // Verifies that when the operation fails permanently while an attempt is in-flight, + // both the attempt span and the operation span end with the error status and message. + Span parentSpan = tracer.spanBuilder("application-parent-operation").startSpan(); + ApiTracer apiTracer; + try (Scope scope = parentSpan.makeCurrent()) { + apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + } + + apiTracer.attemptStarted(new Object(), 0); + // Operation fails while attempt was still in flight + apiTracer.operationFailed(new RuntimeException("network timeout")); + parentSpan.end(); + + List finishedSpans = spanExporter.getFinishedSpanItems(); + assertThat(finishedSpans).hasSize(3); // parent, operation, attempt + + SpanData attemptSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .findFirst() + .orElseThrow(() -> new AssertionError("Attempt span not found")); + SpanData operationSpan = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .findFirst() + .orElseThrow(() -> new AssertionError("Operation span not found")); + + assertThat(attemptSpan.getParentSpanId()).isEqualTo(operationSpan.getSpanContext().getSpanId()); + assertThat(operationSpan.getParentSpanId()).isEqualTo(parentSpan.getSpanContext().getSpanId()); + assertThat( + attemptSpan + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isEqualTo("network timeout"); + } + + @Test + void testAttemptStarted_afterOperationCompleted_doesNotEmitNewSpan() { + // Verifies that after operation completion, subsequent attemptStarted calls + // do not create orphan attempt spans. + ApiTracer apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + + apiTracer.operationSucceeded(); + + // Any attempts started after operation completed must be ignored + apiTracer.attemptStarted(new Object(), 0); + + // Only the operation span was emitted and ended + List finishedSpans = spanExporter.getFinishedSpanItems(); + assertThat(finishedSpans).hasSize(1); + assertThat(finishedSpans.get(0).getKind()).isEqualTo(SpanKind.INTERNAL); + } + + @Test + void testRetrySucceeds_operationAggregatesSuccessAttributes() { + // Verifies that when a transient failure is retried and succeeds: + // 1. Exactly one overall INTERNAL operation span (T3) is created. + // 2. Both attempt spans (T4) have the operation span as their parent. + // 3. The operation span aggregates the successful status (OK) from the final attempt. + ApiTracer apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + + // Attempt 0 fails with transient error + apiTracer.attemptStarted(new Object(), 0); + apiTracer.attemptFailedDuration(new RuntimeException("transient 503"), Duration.ofMillis(10)); + + // Attempt 1 succeeds + apiTracer.attemptStarted(new Object(), 1); + apiTracer.attemptSucceeded(); + apiTracer.operationSucceeded(); + + SpanData operationSpan = + verifySpanHierarchyAndGetOperationSpan(spanExporter.getFinishedSpanItems()); + + assertThat(operationSpan.getStatus().getStatusCode()) + .isEqualTo(io.opentelemetry.api.trace.StatusCode.UNSET); + assertThat( + operationSpan + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE))) + .isEqualTo("OK"); + assertThat( + operationSpan + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNull(); + } + + @Test + void testRetriesExhausted_operationAggregatesFailureAttributes() { + // Verifies that when retries are exhausted: + // 1. Exactly one overall INTERNAL operation span (T3) is created. + // 2. Both attempt spans (T4) have the operation span as their parent. + // 3. The operation span aggregates the final ERROR status and error attributes. + ApiTracer apiTracer = tracingFactory.newTracer(BaseApiTracer.getInstance(), TRACER_CONTEXT); + + // Attempt 0 fails with transient error + apiTracer.attemptStarted(new Object(), 0); + apiTracer.attemptFailedDuration(new RuntimeException("transient 503"), Duration.ofMillis(10)); + + // Attempt 1 fails and exhausts retries + apiTracer.attemptStarted(new Object(), 1); + RuntimeException finalError = new RuntimeException("unavailable: retries exhausted"); + apiTracer.attemptFailedRetriesExhausted(finalError); + apiTracer.operationFailed(finalError); + + SpanData operationSpan = + verifySpanHierarchyAndGetOperationSpan(spanExporter.getFinishedSpanItems()); + + assertThat(operationSpan.getStatus().getStatusCode()) + .isEqualTo(io.opentelemetry.api.trace.StatusCode.ERROR); + assertThat( + operationSpan + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE))) + .isEqualTo("unavailable: retries exhausted"); + assertThat( + operationSpan + .getAttributes() + .get(AttributeKey.stringKey(ObservabilityAttributes.ERROR_TYPE_ATTRIBUTE))) + .isNotNull(); + } + + /** + * Verifies that exactly one internal operation span and two client attempt spans exist in the + * finished spans, and that each attempt span has the operation span as its parent. + * + * @param finishedSpans the list of recorded finished spans + * @return the verified operation {@link SpanData} + */ + private SpanData verifySpanHierarchyAndGetOperationSpan(List finishedSpans) { + assertThat(finishedSpans).hasSize(3); + + List internalSpans = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.INTERNAL) + .collect(Collectors.toList()); + assertThat(internalSpans).hasSize(1); + SpanData operationSpan = internalSpans.get(0); + + List attemptSpans = + finishedSpans.stream() + .filter(s -> s.getKind() == SpanKind.CLIENT) + .collect(Collectors.toList()); + assertThat(attemptSpans).hasSize(2); + + for (SpanData attempt : attemptSpans) { + assertThat(attempt.getParentSpanId()).isEqualTo(operationSpan.getSpanContext().getSpanId()); + } + return operationSpan; + } +} diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerTest.java index 33fa2efcc0da..6d1659ee0c71 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/tracing/OpenTelemetryTracingTracerTest.java @@ -30,9 +30,12 @@ package com.google.api.gax.tracing; import static com.google.common.truth.Truth.assertThat; +import static io.opentelemetry.api.trace.StatusCode.ERROR; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -51,6 +54,7 @@ import io.opentelemetry.api.trace.Tracer; import java.net.ConnectException; import java.net.SocketTimeoutException; +import java.util.HashMap; import java.util.Map; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -64,15 +68,32 @@ class OpenTelemetryTracingTracerTest { @Mock private Tracer tracer; @Mock private SpanBuilder spanBuilder; @Mock private Span span; + @Mock private SpanBuilder operationSpanBuilder; + @Mock private Span operationSpan; private OpenTelemetryTracingTracer openTelemetryTracingTracer; private static final String ATTEMPT_SPAN_NAME = "Service/Method/attempt"; @BeforeEach void setUp() { - when(tracer.spanBuilder(anyString())).thenReturn(spanBuilder); - when(spanBuilder.setSpanKind(any(SpanKind.class))).thenReturn(spanBuilder); - when(spanBuilder.setAllAttributes(any(Attributes.class))).thenReturn(spanBuilder); - when(spanBuilder.startSpan()).thenReturn(span); + lenient().when(tracer.spanBuilder(anyString())).thenReturn(spanBuilder); + lenient().when(spanBuilder.setSpanKind(any(SpanKind.class))).thenReturn(spanBuilder); + lenient().when(spanBuilder.setParent(any())).thenReturn(spanBuilder); + lenient().when(spanBuilder.setAllAttributes(any(Attributes.class))).thenReturn(spanBuilder); + lenient().when(spanBuilder.startSpan()).thenReturn(span); + + lenient() + .when(operationSpanBuilder.setSpanKind(any(SpanKind.class))) + .thenReturn(operationSpanBuilder); + lenient().when(operationSpanBuilder.setParent(any())).thenReturn(operationSpanBuilder); + lenient() + .when(operationSpanBuilder.setAllAttributes(any(Attributes.class))) + .thenReturn(operationSpanBuilder); + lenient().when(operationSpanBuilder.startSpan()).thenReturn(operationSpan); + lenient() + .when(operationSpan.storeInContext(any(io.opentelemetry.context.Context.class))) + .thenAnswer(invocation -> invocation.getArgument(0)); + lenient().when(tracer.spanBuilder("Service/Method")).thenReturn(operationSpanBuilder); + openTelemetryTracingTracer = new OpenTelemetryTracingTracer(tracer, ApiTracerContext.empty(), ATTEMPT_SPAN_NAME); } @@ -680,4 +701,128 @@ void testInjectTraceContext_addsHeaders() { assertThat(carrier.get("traceparent")).contains("00000000000000000000000000000001"); assertThat(carrier.get("traceparent")).contains("0000000000000002"); } + + @Test + void testAttemptStarted_setsParentToParentContext() { + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + verify(spanBuilder).setParent(any(io.opentelemetry.context.Context.class)); + } + + @Test + void testOperationSucceeded_endsActiveAttemptSpan() { + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + openTelemetryTracingTracer.operationSucceeded(); + + verify(span).end(); + verify(operationSpan).end(); + } + + @Test + void testOperationFailed_endsActiveAttemptSpanWithErrorAttributes() { + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + openTelemetryTracingTracer.operationFailed(new RuntimeException("operation failed")); + + verify(span).setAttribute(ObservabilityAttributes.STATUS_MESSAGE_ATTRIBUTE, "operation failed"); + verify(span).end(); + verify(operationSpan).setStatus(ERROR); + verify(operationSpan).end(); + } + + @Test + void testOperationCancelled_endsActiveAttemptSpanWithCancellation() { + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + openTelemetryTracingTracer.operationCancelled(); + + ArgumentCaptor attrsCaptor = ArgumentCaptor.forClass(Attributes.class); + verify(span).setAllAttributes(attrsCaptor.capture()); + verify(span).end(); + verify(operationSpan).setStatus(ERROR); + verify(operationSpan).end(); + + assertThat(attrsCaptor.getValue().asMap()) + .containsEntry( + AttributeKey.stringKey(ObservabilityAttributes.RPC_RESPONSE_STATUS_ATTRIBUTE), + "CANCELLED"); + } + + @Test + void testInScope_withAttemptSpan() { + // Verifies that inScope() activates the current attempt span if an attempt is currently active. + io.opentelemetry.context.Scope mockScope = mock(io.opentelemetry.context.Scope.class); + when(span.makeCurrent()).thenReturn(mockScope); + + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + try (ApiTracer.Scope scope = openTelemetryTracingTracer.inScope()) { + verify(span).makeCurrent(); + } + verify(mockScope).close(); + } + + @Test + void testInScope_withOperationSpanFallback() { + // Verifies that inScope() falls back to activating the operation span when no attempt span is + // active. + io.opentelemetry.context.Scope mockScope = mock(io.opentelemetry.context.Scope.class); + when(operationSpan.makeCurrent()).thenReturn(mockScope); + + try (ApiTracer.Scope scope = openTelemetryTracingTracer.inScope()) { + verify(operationSpan).makeCurrent(); + } + verify(mockScope).close(); + } + + @Test + void testInjectTraceContext_withOperationSpanFallback() { + // Verifies that injectTraceContext() injects the operation span context into the carrier + // when between attempts so that context propagation doesn't drop trace state. + io.opentelemetry.api.trace.SpanContext mockSpanContext = + io.opentelemetry.api.trace.SpanContext.create( + "00000000000000000000000000000003", + "0000000000000004", + io.opentelemetry.api.trace.TraceFlags.getSampled(), + io.opentelemetry.api.trace.TraceState.getDefault()); + Span realSpan = Span.wrap(mockSpanContext); + when(operationSpanBuilder.startSpan()).thenReturn(realSpan); + + openTelemetryTracingTracer = + new OpenTelemetryTracingTracer(tracer, ApiTracerContext.empty(), ATTEMPT_SPAN_NAME); + + Map carrier = new HashMap<>(); + openTelemetryTracingTracer.injectTraceContext(carrier); + + assertThat(carrier).containsKey("traceparent"); + assertThat(carrier.get("traceparent")).contains("00000000000000000000000000000003"); + assertThat(carrier.get("traceparent")).contains("0000000000000004"); + } + + @Test + void testAttemptStarted_whenPreviousAttemptActive_closesOldSpan() { + // Verifies that starting a new retry attempt cleanly closes any lingering previous attempt + // span. + Span span1 = mock(Span.class); + Span span2 = mock(Span.class); + + when(spanBuilder.startSpan()).thenReturn(span1, span2); + + openTelemetryTracingTracer.attemptStarted(new Object(), 0); + + // Start a second attempt before the first attempt was ended + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + verify(span1).end(); + verify(span2, never()).end(); + + // Now complete the second attempt + openTelemetryTracingTracer.attemptSucceeded(); + verify(span2).end(); + } + + @Test + void testAttemptStarted_afterOperationCompleted_doesNotStartNewSpan() { + // Verifies that after operation completion, late callbacks cannot spawn new attempt spans. + openTelemetryTracingTracer.operationSucceeded(); + + // Attempting to start a new attempt after operation completion should be a no-op + openTelemetryTracingTracer.attemptStarted(new Object(), 1); + verify(spanBuilder, never()).startSpan(); + } }