From 69ec9f6d8e439873ab0eedeedc7b35219362f231 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Wed, 9 Sep 2026 20:27:58 +0000 Subject: [PATCH 1/3] Add proto telemetry header to improve server side visibility --- .../com/google/cloud/pubsub/v1/Publisher.java | 25 +++++- .../cloud/pubsub/v1/PublisherImplTest.java | 77 +++++++++++++++++++ 2 files changed, 101 insertions(+), 1 deletion(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index ea98a8884507..623c42d9eb65 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -49,10 +49,14 @@ import com.google.cloud.pubsub.v1.stub.PublisherStubSettings; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import com.google.protobuf.CodedOutputStream; +import com.google.protobuf.util.Timestamps; import com.google.pubsub.v1.PublishRequest; import com.google.pubsub.v1.PublishResponse; +import com.google.pubsub.v1.PubsubClientTelemetry; import com.google.pubsub.v1.PubsubMessage; import com.google.pubsub.v1.TopicName; import com.google.pubsub.v1.TopicNames; @@ -63,6 +67,7 @@ import java.io.IOException; import java.time.Duration; import java.util.ArrayList; +import java.util.Base64; import java.util.Collections; import java.util.HashMap; import java.util.Iterator; @@ -105,6 +110,9 @@ public class Publisher implements PublisherInterface { private static final Logger logger = Logger.getLogger(Publisher.class.getName()); private LoggingUtil loggingUtil = new LoggingUtil(); + @VisibleForTesting + static final String TELEMETRY_HEADER_KEY = "x-goog-pubsub-client-telemetry"; + private static final String GZIP_COMPRESSION = "gzip"; private static final String OPEN_TELEMETRY_TRACER_NAME = "com.google.cloud.pubsub.v1"; @@ -540,6 +548,18 @@ private void publishAllWithoutInflightForKey(final String orderingKey) { } } + private String createTelemetryHeader(OutstandingBatch outstandingBatch, int attemptNumber) { + PubsubClientTelemetry telemetry = + PubsubClientTelemetry.newBuilder() + .setPublishOperation( + PubsubClientTelemetry.PublishOperation.newBuilder() + .setHedgedAttemptCount(attemptNumber) + .setPublishStartTime(Timestamps.fromMillis(outstandingBatch.creationTime)) + .build()) + .build(); + return Base64.getEncoder().encodeToString(telemetry.toByteArray()); + } + private ApiFuture publishCall(OutstandingBatch outstandingBatch) { return publishCall(outstandingBatch, 0, null); } @@ -561,7 +581,10 @@ private ApiFuture publishCall( outstandingBatch.getMessageWrappers().get(0)); context = context.withRetryableCodes(Collections.emptySet()); } - + String telemetryHeader = createTelemetryHeader(outstandingBatch, attemptNumber); + Map> extraHeaders = + ImmutableMap.of("x-goog-pubsub-client-telemetry", ImmutableList.of(telemetryHeader)); + context = context.withExtraHeaders(extraHeaders); int numMessagesInBatch = outstandingBatch.size(); List pubsubMessagesList = new ArrayList(numMessagesInBatch); List messageWrappers = outstandingBatch.getMessageWrappers(); diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index ddb3b5b42be6..faa5e1d54d6c 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -43,6 +43,8 @@ import com.google.pubsub.v1.PublishRequest; import com.google.pubsub.v1.PublishResponse; import com.google.pubsub.v1.PubsubMessage; +import com.google.pubsub.v1.PubsubClientTelemetry; +import java.util.Base64; import io.grpc.ManagedChannel; import io.grpc.Metadata; import io.grpc.Server; @@ -1737,4 +1739,79 @@ private void shutdownTestPublisher(Publisher publisher) throws InterruptedExcept fakeExecutor.advanceTime(Duration.ofSeconds(10)); assertTrue(publisher.awaitTermination(1, TimeUnit.MINUTES)); } + + private PubsubClientTelemetry extractTelemetryHeader(Metadata headers) throws Exception { + Metadata.Key key = + Metadata.Key.of(Publisher.TELEMETRY_HEADER_KEY, Metadata.ASCII_STRING_MARSHALLER); + String headerValue = headers.get(key); + assertThat(headerValue).isNotNull(); + byte[] decodedBytes = Base64.getDecoder().decode(headerValue); + return PubsubClientTelemetry.parseFrom(decodedBytes); + } + + @Test + public void testTelemetryHeaderOnNormalPublish() throws Exception { + testPublisherServiceImpl.setAutoPublishResponse(true); + Publisher publisher = + getTestPublisherBuilder() + .setBatchingSettings( + Publisher.Builder.DEFAULT_BATCHING_SETTINGS.toBuilder() + .setElementCountThreshold(1L) + .build()) + .build(); + + ApiFuture future = sendTestMessage(publisher, "msg-normal"); + assertEquals("1", future.get(5, TimeUnit.SECONDS)); + + List capturedHeaders = testPublisherServiceImpl.getCapturedHeaders(); + assertThat(capturedHeaders).hasSize(1); + + PubsubClientTelemetry telemetry = extractTelemetryHeader(capturedHeaders.get(0)); + assertThat(telemetry.hasPublishOperation()).isTrue(); + assertThat(telemetry.getPublishOperation().getHedgedAttemptCount()).isEqualTo(0); + assertThat(telemetry.getPublishOperation().getPublishStartTime().getSeconds()) + .isGreaterThan(0); + + shutdownTestPublisher(publisher); + } + + @Test + public void testTelemetryHeaderOnHedgedPublish() throws Exception { + Publisher publisher = getPublisherWithHedge(Duration.ofMillis(100), 0.2f, 20); + fillTokenBucket(publisher, 5); + + // Delay response so hedge fires + testPublisherServiceImpl.setAutoPublishResponse(false); + testPublisherServiceImpl.setPublishResponseDelay(Duration.ofMillis(200)); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("1")); + testPublisherServiceImpl.addPublishResponse(PublishResponse.newBuilder().addMessageIds("2")); + + ApiFuture future = sendTestMessage(publisher, "msg-hedged"); + waitForRequests(testPublisherServiceImpl, 1); + + // Advance past hedge delay to trigger hedge attempt + fakeExecutor.advanceTime(Duration.ofMillis(120)); + waitForRequests(testPublisherServiceImpl, 2); + + // Finish request + fakeExecutor.advanceTime(Duration.ofMillis(100)); + assertEquals("1", future.get(5, TimeUnit.SECONDS)); + + List capturedHeaders = testPublisherServiceImpl.getCapturedHeaders(); + assertThat(capturedHeaders).hasSize(2); + + // Verify Attempt 0 (Original) + PubsubClientTelemetry initialTelemetry = extractTelemetryHeader(capturedHeaders.get(0)); + assertThat(initialTelemetry.getPublishOperation().getHedgedAttemptCount()).isEqualTo(0); + + // Verify Attempt 1 (Hedge) + PubsubClientTelemetry hedgedTelemetry = extractTelemetryHeader(capturedHeaders.get(1)); + assertThat(hedgedTelemetry.getPublishOperation().getHedgedAttemptCount()).isEqualTo(1); + + // Both must share the identical start time + assertThat(hedgedTelemetry.getPublishOperation().getPublishStartTime()) + .isEqualTo(initialTelemetry.getPublishOperation().getPublishStartTime()); + + shutdownTestPublisher(publisher); + } } From 94de134281c8be4a86a4009237209e3fb5db06c2 Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Wed, 9 Sep 2026 21:28:23 +0000 Subject: [PATCH 2/3] Fix lint --- .../main/java/com/google/cloud/pubsub/v1/Publisher.java | 5 ++--- .../java/com/google/cloud/pubsub/v1/PublisherImplTest.java | 7 +++---- 2 files changed, 5 insertions(+), 7 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index 623c42d9eb65..f673d5c48f73 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -110,8 +110,7 @@ public class Publisher implements PublisherInterface { private static final Logger logger = Logger.getLogger(Publisher.class.getName()); private LoggingUtil loggingUtil = new LoggingUtil(); - @VisibleForTesting - static final String TELEMETRY_HEADER_KEY = "x-goog-pubsub-client-telemetry"; + @VisibleForTesting static final String TELEMETRY_HEADER_KEY = "x-goog-pubsub-client-telemetry"; private static final String GZIP_COMPRESSION = "gzip"; @@ -557,7 +556,7 @@ private String createTelemetryHeader(OutstandingBatch outstandingBatch, int atte .setPublishStartTime(Timestamps.fromMillis(outstandingBatch.creationTime)) .build()) .build(); - return Base64.getEncoder().encodeToString(telemetry.toByteArray()); + return Base64.getEncoder().encodeToString(telemetry.toByteArray()); } private ApiFuture publishCall(OutstandingBatch outstandingBatch) { diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index faa5e1d54d6c..f1170b60d3e1 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -42,9 +42,8 @@ import com.google.pubsub.v1.ProjectTopicName; import com.google.pubsub.v1.PublishRequest; import com.google.pubsub.v1.PublishResponse; -import com.google.pubsub.v1.PubsubMessage; import com.google.pubsub.v1.PubsubClientTelemetry; -import java.util.Base64; +import com.google.pubsub.v1.PubsubMessage; import io.grpc.ManagedChannel; import io.grpc.Metadata; import io.grpc.Server; @@ -63,6 +62,7 @@ import io.opentelemetry.sdk.testing.junit4.OpenTelemetryRule; import io.opentelemetry.sdk.trace.data.SpanData; import java.time.Duration; +import java.util.Base64; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; @@ -1769,8 +1769,7 @@ public void testTelemetryHeaderOnNormalPublish() throws Exception { PubsubClientTelemetry telemetry = extractTelemetryHeader(capturedHeaders.get(0)); assertThat(telemetry.hasPublishOperation()).isTrue(); assertThat(telemetry.getPublishOperation().getHedgedAttemptCount()).isEqualTo(0); - assertThat(telemetry.getPublishOperation().getPublishStartTime().getSeconds()) - .isGreaterThan(0); + assertThat(telemetry.getPublishOperation().getPublishStartTime().getSeconds()).isGreaterThan(0); shutdownTestPublisher(publisher); } From a926bdb0330efcabe8d0bf94effb761d8344f5ec Mon Sep 17 00:00:00 2001 From: Tony Cui Date: Fri, 11 Sep 2026 01:30:11 +0000 Subject: [PATCH 3/3] Remove timestamps for now --- .../java/com/google/cloud/pubsub/v1/Publisher.java | 12 ++++-------- .../google/cloud/pubsub/v1/PublisherImplTest.java | 5 ----- 2 files changed, 4 insertions(+), 13 deletions(-) diff --git a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java index f673d5c48f73..807f91c9ea01 100644 --- a/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java +++ b/java-pubsub/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/v1/Publisher.java @@ -53,7 +53,6 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import com.google.protobuf.CodedOutputStream; -import com.google.protobuf.util.Timestamps; import com.google.pubsub.v1.PublishRequest; import com.google.pubsub.v1.PublishResponse; import com.google.pubsub.v1.PubsubClientTelemetry; @@ -547,16 +546,16 @@ private void publishAllWithoutInflightForKey(final String orderingKey) { } } - private String createTelemetryHeader(OutstandingBatch outstandingBatch, int attemptNumber) { + private Map> createTelemetryHeader(int attemptNumber) { PubsubClientTelemetry telemetry = PubsubClientTelemetry.newBuilder() .setPublishOperation( PubsubClientTelemetry.PublishOperation.newBuilder() .setHedgedAttemptCount(attemptNumber) - .setPublishStartTime(Timestamps.fromMillis(outstandingBatch.creationTime)) .build()) .build(); - return Base64.getEncoder().encodeToString(telemetry.toByteArray()); + String encodedHeader = Base64.getEncoder().encodeToString(telemetry.toByteArray()); + return ImmutableMap.of(TELEMETRY_HEADER_KEY, ImmutableList.of(encodedHeader)); } private ApiFuture publishCall(OutstandingBatch outstandingBatch) { @@ -580,10 +579,7 @@ private ApiFuture publishCall( outstandingBatch.getMessageWrappers().get(0)); context = context.withRetryableCodes(Collections.emptySet()); } - String telemetryHeader = createTelemetryHeader(outstandingBatch, attemptNumber); - Map> extraHeaders = - ImmutableMap.of("x-goog-pubsub-client-telemetry", ImmutableList.of(telemetryHeader)); - context = context.withExtraHeaders(extraHeaders); + context = context.withExtraHeaders(createTelemetryHeader(attemptNumber)); int numMessagesInBatch = outstandingBatch.size(); List pubsubMessagesList = new ArrayList(numMessagesInBatch); List messageWrappers = outstandingBatch.getMessageWrappers(); diff --git a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java index f1170b60d3e1..9c2c4fed9807 100644 --- a/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java +++ b/java-pubsub/google-cloud-pubsub/src/test/java/com/google/cloud/pubsub/v1/PublisherImplTest.java @@ -1769,7 +1769,6 @@ public void testTelemetryHeaderOnNormalPublish() throws Exception { PubsubClientTelemetry telemetry = extractTelemetryHeader(capturedHeaders.get(0)); assertThat(telemetry.hasPublishOperation()).isTrue(); assertThat(telemetry.getPublishOperation().getHedgedAttemptCount()).isEqualTo(0); - assertThat(telemetry.getPublishOperation().getPublishStartTime().getSeconds()).isGreaterThan(0); shutdownTestPublisher(publisher); } @@ -1807,10 +1806,6 @@ public void testTelemetryHeaderOnHedgedPublish() throws Exception { PubsubClientTelemetry hedgedTelemetry = extractTelemetryHeader(capturedHeaders.get(1)); assertThat(hedgedTelemetry.getPublishOperation().getHedgedAttemptCount()).isEqualTo(1); - // Both must share the identical start time - assertThat(hedgedTelemetry.getPublishOperation().getPublishStartTime()) - .isEqualTo(initialTelemetry.getPublishOperation().getPublishStartTime()); - shutdownTestPublisher(publisher); } }