diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java index a4aee6142f40..1e7c0e3fb45f 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java @@ -200,6 +200,7 @@ public class BigQueryConnection extends BigQueryNoOpsConnection { List queryProperties; Map authProperties; Map overrideProperties; + Map proxyProperties; Credentials credentials; boolean useStatelessQueryMode; int numBufferedRows; @@ -299,7 +300,7 @@ public class BigQueryConnection extends BigQueryNoOpsConnection { String.valueOf(ds.getRequestGoogleDriveScope()), BigQueryJdbcUrlUtility.REQUEST_GOOGLE_DRIVE_SCOPE_PROPERTY_NAME); - Map proxyProperties = + this.proxyProperties = BigQueryJdbcProxyUtility.parseProxyProperties(ds, this.connectionClassName); this.sslTrustStorePath = ds.getSSLTrustStorePath(); @@ -1204,14 +1205,20 @@ private OpenTelemetry getOpenTelemetryInstance() { this.customOpenTelemetry, this.gcpTelemetryCredentials, effectiveProjectId, - this.credentials); + this.credentials, + this.proxyProperties); boolean hasExternalOtel = this.customOpenTelemetry != null || this.useGlobalOpenTelemetry; Logging localLoggingClient = null; if (this.enableGcpLogExporter && !hasExternalOtel) { localLoggingClient = BigQueryJdbcOpenTelemetry.createLoggingClient( - true, null, this.gcpTelemetryCredentials, effectiveProjectId, this.credentials); + true, + null, + this.gcpTelemetryCredentials, + effectiveProjectId, + this.credentials, + this.headerProvider); } if (this.enableGcpLogExporter || hasExternalOtel) { diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java index 91404005c330..3ecca298480d 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java @@ -16,6 +16,7 @@ package com.google.cloud.bigquery.jdbc; +import com.google.api.gax.rpc.HeaderProvider; import com.google.auth.Credentials; import com.google.auth.oauth2.GoogleCredentials; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; @@ -32,9 +33,16 @@ import io.opentelemetry.context.Context; import io.opentelemetry.context.Scope; import io.opentelemetry.exporter.otlp.http.trace.OtlpHttpSpanExporter; +import io.opentelemetry.exporter.otlp.http.trace.OtlpHttpSpanExporterBuilder; import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter; import io.opentelemetry.sdk.OpenTelemetrySdk; import io.opentelemetry.sdk.autoconfigure.AutoConfiguredOpenTelemetrySdk; +import io.opentelemetry.sdk.common.export.ProxyOptions; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.net.Proxy; +import java.net.ProxySelector; +import java.net.SocketAddress; import java.net.URI; import java.nio.charset.StandardCharsets; import java.sql.SQLException; @@ -85,11 +93,25 @@ private static final class SdkCacheKey { private final String projectId; private final String credentialsHashOrPath; private final boolean enableTrace; - - SdkCacheKey(String projectId, String credentialsHashOrPath, boolean enableTrace) { + private final String proxyHost; + private final String proxyPort; + + SdkCacheKey( + String projectId, + String credentialsHashOrPath, + boolean enableTrace, + Map proxyProperties) { this.projectId = projectId; this.credentialsHashOrPath = credentialsHashOrPath; this.enableTrace = enableTrace; + this.proxyHost = + proxyProperties != null + ? proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME) + : null; + this.proxyPort = + proxyProperties != null + ? proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME) + : null; } @Override @@ -99,12 +121,14 @@ public boolean equals(Object o) { SdkCacheKey that = (SdkCacheKey) o; return enableTrace == that.enableTrace && Objects.equals(projectId, that.projectId) - && Objects.equals(credentialsHashOrPath, that.credentialsHashOrPath); + && Objects.equals(credentialsHashOrPath, that.credentialsHashOrPath) + && Objects.equals(proxyHost, that.proxyHost) + && Objects.equals(proxyPort, that.proxyPort); } @Override public int hashCode() { - return Objects.hash(projectId, credentialsHashOrPath, enableTrace); + return Objects.hash(projectId, credentialsHashOrPath, enableTrace, proxyHost, proxyPort); } } @@ -181,7 +205,8 @@ public static Logging createLoggingClient( OpenTelemetry customOpenTelemetry, String effectiveCredentials, String effectiveProjectId, - Credentials fallbackCredentials) { + Credentials fallbackCredentials, + HeaderProvider headerProvider) { if (!enableGcpLogExporter || customOpenTelemetry != null) { return null; @@ -200,6 +225,9 @@ public static Logging createLoggingClient( if (credentials != null) { loggingOptionsBuilder.setCredentials(credentials); } + if (headerProvider != null) { + loggingOptionsBuilder.setHeaderProvider(headerProvider); + } return loggingOptionsBuilder.build().getService(); } catch (Exception e) { throw new BigQueryJdbcRuntimeException("Failed to initialize Logging client", e); @@ -317,7 +345,8 @@ public static OpenTelemetry getOpenTelemetry( OpenTelemetry customOpenTelemetry, String gcpTelemetryCredentials, String gcpTelemetryProjectId, - Credentials fallbackCredentials) { + Credentials fallbackCredentials, + Map proxyProperties) { if (customOpenTelemetry != null) { return customOpenTelemetry; @@ -335,7 +364,8 @@ public static OpenTelemetry getOpenTelemetry( new SdkCacheKey( gcpTelemetryProjectId, getCredentialsIdentifier(gcpTelemetryCredentials), - enableGcpTraceExporter); + enableGcpTraceExporter, + proxyProperties); CachedSdk fastCheck = sdkCache.get(key); if (fastCheck != null) { CachedSdk result = @@ -415,11 +445,17 @@ public static OpenTelemetry getOpenTelemetry( final Credentials finalCredentials = credentials; if (spanExporter instanceof OtlpHttpSpanExporter) { - return ((OtlpHttpSpanExporter) spanExporter) - .toBuilder() - .setHeaders( - () -> getAuthHeaders(finalCredentials, gcpTelemetryProjectId)) - .build(); + OtlpHttpSpanExporterBuilder builder = + ((OtlpHttpSpanExporter) spanExporter).toBuilder(); + builder.setHeaders( + () -> getAuthHeaders(finalCredentials, gcpTelemetryProjectId)); + + ProxyOptions proxyOptions = createProxyOptions(proxyProperties); + if (proxyOptions != null) { + builder.setProxy(proxyOptions); + } + + return builder.build(); } if (spanExporter instanceof OtlpGrpcSpanExporter) { return ((OtlpGrpcSpanExporter) spanExporter) @@ -506,4 +542,37 @@ public static T withTracing( span.end(); } } + + private static ProxyOptions createProxyOptions(Map proxyProperties) { + if (proxyProperties == null) { + return null; + } + + final String host = proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME); + final String portStr = proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME); + if (host == null || host.isEmpty() || portStr == null || portStr.isEmpty()) { + return null; + } + + int port; + try { + port = Integer.parseInt(portStr); + } catch (NumberFormatException e) { + throw new BigQueryJdbcRuntimeException("Invalid proxy port number: " + portStr, e); + } + + ProxySelector proxySelector = + new ProxySelector() { + @Override + public List select(URI uri) { + return Collections.singletonList( + new Proxy(Proxy.Type.HTTP, InetSocketAddress.createUnresolved(host, port))); + } + + @Override + public void connectFailed(URI uri, SocketAddress sa, IOException ioe) {} + }; + + return ProxyOptions.create(proxySelector); + } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java index 17a714b22bc7..5798438bcfde 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java @@ -620,7 +620,7 @@ public void testOpenTelemetryPrecedenceHierarchy( .when( () -> BigQueryJdbcOpenTelemetry.createLoggingClient( - anyBoolean(), any(), any(), any(), any())) + anyBoolean(), any(), any(), any(), any(), any())) .thenReturn(mockLogging); // Stub getOpenTelemetry to return the expected mock based on inputs @@ -634,6 +634,7 @@ public void testOpenTelemetryPrecedenceHierarchy( hasCustom ? eq(mockCustomOtel) : isNull(), any(), any(), + any(), any())) .thenAnswer( invocation -> { diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java index 93d2550a3dbb..e1577e6010fd 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java @@ -52,7 +52,7 @@ public void testGetOpenTelemetry_withCustomSdk_returnsCustom() { OpenTelemetry result = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, false, false, mockCustomOtel, null, null, null); + false, false, false, mockCustomOtel, null, null, null, null); assertThat(result).isSameInstanceAs(mockCustomOtel); } @@ -64,7 +64,7 @@ public void testGetOpenTelemetry_withCustomSdkAndFlags_returnsCustom() { // Custom SDK always takes precedence over individual flags OpenTelemetry result = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, true, mockCustomOtel, null, null, null); + false, true, true, mockCustomOtel, null, null, null, null); assertThat(result).isSameInstanceAs(mockCustomOtel); } @@ -72,7 +72,8 @@ public void testGetOpenTelemetry_withCustomSdkAndFlags_returnsCustom() { @Test public void testGetOpenTelemetry_noFlags_returnsNoop() { OpenTelemetry result = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(false, false, false, null, null, null, null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + false, false, false, null, null, null, null, null); assertThat(result).isSameInstanceAs(OpenTelemetry.noop()); } @@ -88,10 +89,10 @@ public void testGetTracer_respectsScopeName() { public void testGetOpenTelemetry_cachesSdkInstances() { OpenTelemetry result1 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); assertThat(result1).isSameInstanceAs(result2); } @@ -100,10 +101,10 @@ public void testGetOpenTelemetry_cachesSdkInstances() { public void testGetOpenTelemetry_createsNewInstanceForDifferentKey() { OpenTelemetry result1 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project2", null); + false, true, false, null, null, "project2", null, null); assertThat(result1).isNotSameInstanceAs(result2); } @@ -111,10 +112,11 @@ public void testGetOpenTelemetry_createsNewInstanceForDifferentKey() { @Test public void testGetOpenTelemetry_createsNewInstanceForDifferentTraceFlag() { OpenTelemetry result1 = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(false, true, true, null, null, "project1", null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + false, true, true, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, false, true, null, null, "project1", null); + false, false, true, null, null, "project1", null, null); assertThat(result1).isNotSameInstanceAs(result2); } @@ -122,10 +124,11 @@ public void testGetOpenTelemetry_createsNewInstanceForDifferentTraceFlag() { @Test public void testGetOpenTelemetry_ignoresEnableLogFlagInCacheKey() { OpenTelemetry result1 = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(false, true, true, null, null, "project1", null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + false, true, true, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); assertThat(result1).isSameInstanceAs(result2); } @@ -133,7 +136,8 @@ public void testGetOpenTelemetry_ignoresEnableLogFlagInCacheKey() { @Test public void testGetOpenTelemetry_withUseGlobalOTel_returnsGlobal() { OpenTelemetry result = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(true, false, false, null, null, null, null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + true, false, false, null, null, null, null, null); assertThat(result).isSameInstanceAs(GlobalOpenTelemetry.get()); } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java index 444b92fcefc0..814c8329bb5f 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java @@ -23,7 +23,9 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import com.google.api.gax.core.FixedCredentialsProvider; +import com.google.api.gax.grpc.InstantiatingGrpcChannelProvider; import com.google.api.gax.paging.Page; +import com.google.api.gax.rpc.TransportChannelProvider; import com.google.auth.oauth2.GoogleCredentials; import com.google.cloud.ServiceOptions; import com.google.cloud.bigquery.jdbc.BigQueryConnection; @@ -36,7 +38,18 @@ import com.google.devtools.cloudtrace.v1.Trace; import com.google.devtools.cloudtrace.v1.TraceSpan; import com.google.gson.JsonObject; +import io.grpc.HttpConnectProxiedSocketAddress; +import io.grpc.ProxiedSocketAddress; +import io.grpc.ProxyDetector; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.api.trace.SpanContext; +import io.opentelemetry.api.trace.TraceFlags; +import io.opentelemetry.api.trace.TraceState; +import io.opentelemetry.context.Scope; +import io.opentelemetry.sdk.trace.IdGenerator; import java.io.File; +import java.net.InetSocketAddress; +import java.net.SocketAddress; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.sql.Connection; @@ -125,7 +138,6 @@ public void testExecute_withOpenTelemetryGcpExporter() throws Exception { @Test public void testExecute_withErrorCorrelation() throws Exception { - // Step 1: Connect with GCP Exporters enabled via DataSource DataSource ds = DataSource.fromUrl(CONNECTION_URL); ds.setEnableGcpTraceExporter(true); @@ -227,6 +239,39 @@ public void testExecute_withGrpcProtocol() throws Exception { } } + @Test + public void testExecute_withHttpProtocol_andDirectTraceVerification() throws Exception { + JsonObject authJson = getAuthJson(); + System.setProperty("otel.exporter.otlp.protocol", "http/protobuf"); + + try { + DataSource ds = DataSource.fromUrl(CONNECTION_URL); + ds.setEnableGcpTraceExporter(true); + ds.setGcpTelemetryProjectId(PROJECT_ID); + ds.setGcpTelemetryCredentials(authJson.toString()); + + String traceId = IdGenerator.random().generateTraceId(); + String spanId = IdGenerator.random().generateSpanId(); + SpanContext parentContext = + SpanContext.create(traceId, spanId, TraceFlags.getSampled(), TraceState.getDefault()); + + try (Scope scope = Span.wrap(parentContext).makeCurrent()) { + try (Connection connection = ds.getConnection(); + Statement statement = connection.createStatement()) { + String query = "SELECT 1;"; + try (ResultSet rs = statement.executeQuery(query)) { + assertTrue(rs.next()); + } + } + } + + Trace trace = verifyAndFetchTrace(traceId); + assertNotNull(trace, "Trace must be found in Cloud Trace API: " + traceId); + } finally { + System.clearProperty("otel.exporter.otlp.protocol"); + } + } + private void verifyTraceDelivery(DataSource ds) throws Exception { ds.setEnableGcpLogExporter(true); ds.setLogLevel("5"); @@ -293,18 +338,42 @@ private Trace verifyAndFetchTrace(String traceId) throws Exception { GoogleCredentials credentials = getCredentials(); - TraceServiceSettings settings = + TraceServiceSettings.Builder settingsBuilder = TraceServiceSettings.newBuilder() - .setCredentialsProvider(FixedCredentialsProvider.create(credentials)) - .build(); + .setCredentialsProvider(FixedCredentialsProvider.create(credentials)); + + DataSource ds = DataSource.fromUrl(CONNECTION_URL); + String proxyHost = ds.getProxyHost(); + String proxyPortStr = ds.getProxyPort(); + if (proxyHost != null && proxyPortStr != null) { + settingsBuilder.setTransportChannelProvider( + createProxyTransportChannelProvider(proxyHost, Integer.parseInt(proxyPortStr))); + } - try (TraceServiceClient traceClient = TraceServiceClient.create(settings)) { + try (TraceServiceClient traceClient = TraceServiceClient.create(settingsBuilder.build())) { Trace trace = fetchTraceWithRetry(traceClient, PROJECT_ID, hexTraceId); assertNotNull(trace, "Trace must be found in Cloud Trace API: " + hexTraceId); return trace; } } + private TransportChannelProvider createProxyTransportChannelProvider(String host, int port) { + return InstantiatingGrpcChannelProvider.newBuilder() + .setChannelConfigurator( + builder -> + builder.proxyDetector( + new ProxyDetector() { + @Override + public ProxiedSocketAddress proxyFor(SocketAddress socketAddress) { + return HttpConnectProxiedSocketAddress.newBuilder() + .setProxyAddress(new InetSocketAddress(host, port)) + .setTargetAddress((InetSocketAddress) socketAddress) + .build(); + } + })) + .build(); + } + private T pollWithRetry(java.util.concurrent.Callable task) throws InterruptedException { int attempts = 0; int maxAttempts = 10;