Skip to content
This repository was archived by the owner on May 8, 2026. It is now read-only.

Commit 06f15f7

Browse files
format
Change-Id: Id9dfd28ca4a3f9e98474b33c98895e05b830b410
1 parent 89703d5 commit 06f15f7

10 files changed

Lines changed: 164 additions & 170 deletions

File tree

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/MetadataExtractorInterceptor.java

Lines changed: 136 additions & 135 deletions
Original file line numberDiff line numberDiff line change
@@ -33,160 +33,161 @@
3333
import io.grpc.MethodDescriptor;
3434
import io.grpc.Status;
3535
import io.grpc.alts.AltsContextUtil;
36-
37-
import javax.annotation.Nullable;
3836
import java.util.Base64;
3937
import java.util.regex.Matcher;
4038
import java.util.regex.Pattern;
39+
import javax.annotation.Nullable;
4140

4241
@InternalApi
4342
public class MetadataExtractorInterceptor implements ClientInterceptor {
44-
private final SidebandData sidebandData = new SidebandData();
45-
46-
public GrpcCallContext injectInto(GrpcCallContext ctx) {
47-
return ctx
48-
.withChannel(ClientInterceptors.intercept(ctx.getChannel(), this))
49-
.withCallOptions(ctx.getCallOptions().withOption(SidebandData.KEY, sidebandData));
43+
private final SidebandData sidebandData = new SidebandData();
44+
45+
public GrpcCallContext injectInto(GrpcCallContext ctx) {
46+
return ctx.withChannel(ClientInterceptors.intercept(ctx.getChannel(), this))
47+
.withCallOptions(ctx.getCallOptions().withOption(SidebandData.KEY, sidebandData));
48+
}
49+
50+
@Override
51+
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
52+
MethodDescriptor<ReqT, RespT> methodDescriptor, CallOptions callOptions, Channel channel) {
53+
return new ForwardingClientCall.SimpleForwardingClientCall<ReqT, RespT>(
54+
channel.newCall(methodDescriptor, callOptions)) {
55+
@Override
56+
public void start(Listener<RespT> responseListener, Metadata headers) {
57+
sidebandData.reset();
58+
59+
super.start(
60+
new ForwardingClientCallListener.SimpleForwardingClientCallListener<RespT>(
61+
responseListener) {
62+
@Override
63+
public void onHeaders(Metadata headers) {
64+
sidebandData.onResponseHeaders(headers, getAttributes());
65+
super.onHeaders(headers);
66+
}
67+
68+
@Override
69+
public void onClose(Status status, Metadata trailers) {
70+
sidebandData.onClose(status, trailers);
71+
super.onClose(status, trailers);
72+
}
73+
},
74+
headers);
75+
}
76+
};
77+
}
78+
79+
public SidebandData getSidebandData() {
80+
return sidebandData;
81+
}
82+
83+
public static class SidebandData {
84+
private static final CallOptions.Key<SidebandData> KEY =
85+
CallOptions.Key.create("bigtable-sideband");
86+
87+
private static final Metadata.Key<String> SERVER_TIMING_HEADER_KEY =
88+
Metadata.Key.of("server-timing", Metadata.ASCII_STRING_MARSHALLER);
89+
private static final Pattern SERVER_TIMING_HEADER_PATTERN =
90+
Pattern.compile(".*dur=(?<dur>\\d+)");
91+
private static final Metadata.Key<byte[]> LOCATION_METADATA_KEY =
92+
Metadata.Key.of("x-goog-ext-425905942-bin", Metadata.BINARY_BYTE_MARSHALLER);
93+
private static final Metadata.Key<String> PEER_INFO_KEY =
94+
Metadata.Key.of("bigtable-peer-info", Metadata.ASCII_STRING_MARSHALLER);
95+
96+
@Nullable private volatile ResponseParams responseParams;
97+
@Nullable private volatile PeerInfo peerInfo;
98+
@Nullable private volatile Long gfeTiming;
99+
100+
@Nullable
101+
public ResponseParams getResponseParams() {
102+
return responseParams;
50103
}
51104

52-
@Override
53-
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(MethodDescriptor<ReqT, RespT> methodDescriptor, CallOptions callOptions, Channel channel) {
54-
return new ForwardingClientCall.SimpleForwardingClientCall<ReqT, RespT>(channel.newCall(methodDescriptor, callOptions)) {
55-
@Override
56-
public void start(Listener<RespT> responseListener, Metadata headers) {
57-
sidebandData.reset();
58-
59-
super.start(new ForwardingClientCallListener.SimpleForwardingClientCallListener<RespT>(responseListener) {
60-
@Override
61-
public void onHeaders(Metadata headers) {
62-
sidebandData.onResponseHeaders(headers, getAttributes());
63-
super.onHeaders(headers);
64-
}
65-
66-
@Override
67-
public void onClose(Status status, Metadata trailers) {
68-
sidebandData.onClose(status, trailers);
69-
super.onClose(status, trailers);
70-
}
71-
}, headers);
72-
}
73-
};
105+
@Nullable
106+
public PeerInfo getPeerInfo() {
107+
return peerInfo;
74108
}
75109

76-
public SidebandData getSidebandData() {
77-
return sidebandData;
110+
@Nullable
111+
public Long getGfeTiming() {
112+
return gfeTiming;
78113
}
79114

80-
public static class SidebandData {
81-
private static final CallOptions.Key<SidebandData> KEY = CallOptions.Key.create("bigtable-sideband");
82-
83-
private static final Metadata.Key<String> SERVER_TIMING_HEADER_KEY =
84-
Metadata.Key.of("server-timing", Metadata.ASCII_STRING_MARSHALLER);
85-
private static final Pattern SERVER_TIMING_HEADER_PATTERN = Pattern.compile(".*dur=(?<dur>\\d+)");
86-
private static final Metadata.Key<byte[]> LOCATION_METADATA_KEY =
87-
Metadata.Key.of("x-goog-ext-425905942-bin", Metadata.BINARY_BYTE_MARSHALLER);
88-
private static final Metadata.Key<String> PEER_INFO_KEY =
89-
Metadata.Key.of("bigtable-peer-info", Metadata.ASCII_STRING_MARSHALLER);
90-
91-
@Nullable
92-
private volatile ResponseParams responseParams;
93-
@Nullable
94-
private volatile PeerInfo peerInfo;
95-
@Nullable
96-
private volatile Long gfeTiming;
97-
98-
@Nullable
99-
public ResponseParams getResponseParams() {
100-
return responseParams;
101-
}
102-
103-
@Nullable
104-
public PeerInfo getPeerInfo() {
105-
return peerInfo;
106-
}
115+
private void reset() {
116+
responseParams = null;
117+
peerInfo = null;
118+
gfeTiming = null;
119+
}
107120

108-
@Nullable
109-
public Long getGfeTiming() {
110-
return gfeTiming;
111-
}
121+
void onResponseHeaders(Metadata md, Attributes attributes) {
122+
responseParams = extractResponseParams(md);
123+
gfeTiming = extractGfeLatency(md);
124+
peerInfo = extractPeerInfo(md, gfeTiming, attributes);
125+
}
112126

127+
void onClose(Status status, Metadata trailers) {
128+
if (responseParams == null) {
129+
responseParams = extractResponseParams(trailers);
130+
}
131+
}
113132

114-
private void reset() {
115-
responseParams = null;
116-
peerInfo = null;
117-
gfeTiming = null;
118-
}
119-
void onResponseHeaders(Metadata md, Attributes attributes) {
120-
responseParams = extractResponseParams(md);
121-
gfeTiming = extractGfeLatency(md);
122-
peerInfo = extractPeerInfo(md, gfeTiming, attributes);
123-
}
124-
void onClose(Status status, Metadata trailers) {
125-
if (responseParams == null) {
126-
responseParams = extractResponseParams(trailers);
127-
}
128-
}
133+
@Nullable
134+
private static Long extractGfeLatency(Metadata metadata) {
135+
String serverTiming = metadata.get(SERVER_TIMING_HEADER_KEY);
136+
if (serverTiming == null) {
137+
return null;
138+
}
139+
Matcher matcher = SERVER_TIMING_HEADER_PATTERN.matcher(serverTiming);
140+
// this should always be true
141+
if (matcher.find()) {
142+
return Long.parseLong(matcher.group("dur"));
143+
}
144+
return null;
145+
}
129146

130-
@Nullable
131-
private static Long extractGfeLatency(Metadata metadata) {
132-
String serverTiming = metadata.get(SERVER_TIMING_HEADER_KEY);
133-
if (serverTiming == null) {
134-
return null;
135-
}
136-
Matcher matcher = SERVER_TIMING_HEADER_PATTERN.matcher(serverTiming);
137-
// this should always be true
138-
if (matcher.find()) {
139-
return Long.parseLong(matcher.group("dur"));
140-
}
141-
return null;
147+
@Nullable
148+
private static PeerInfo extractPeerInfo(
149+
Metadata metadata, Long gfeTiming, Attributes attributes) {
150+
String encodedStr = metadata.get(PEER_INFO_KEY);
151+
if (Strings.isNullOrEmpty(encodedStr)) {
152+
return null;
153+
}
154+
155+
try {
156+
byte[] decoded = Base64.getUrlDecoder().decode(encodedStr);
157+
PeerInfo peerInfo = PeerInfo.parseFrom(decoded);
158+
PeerInfo.TransportType effectiveTransport = peerInfo.getTransportType();
159+
160+
if (effectiveTransport == PeerInfo.TransportType.TRANSPORT_TYPE_UNKNOWN) {
161+
boolean isAlts = AltsContextUtil.check(attributes);
162+
if (isAlts) {
163+
effectiveTransport = PeerInfo.TransportType.TRANSPORT_TYPE_DIRECT_ACCESS;
164+
} else if (gfeTiming != null) {
165+
effectiveTransport = PeerInfo.TransportType.TRANSPORT_TYPE_CLOUD_PATH;
166+
}
142167
}
143-
144-
@Nullable
145-
private static PeerInfo extractPeerInfo(Metadata metadata, Long gfeTiming, Attributes attributes) {
146-
String encodedStr = metadata.get(PEER_INFO_KEY);
147-
if (Strings.isNullOrEmpty(encodedStr)) {
148-
return null;
149-
}
150-
151-
try {
152-
byte[] decoded = Base64.getUrlDecoder().decode(encodedStr);
153-
PeerInfo peerInfo = PeerInfo.parseFrom(decoded);
154-
PeerInfo.TransportType effectiveTransport = peerInfo.getTransportType();
155-
156-
if (effectiveTransport == PeerInfo.TransportType.TRANSPORT_TYPE_UNKNOWN) {
157-
boolean isAlts = AltsContextUtil.check(attributes);
158-
if (isAlts) {
159-
effectiveTransport = PeerInfo.TransportType.TRANSPORT_TYPE_DIRECT_ACCESS;
160-
} else if (gfeTiming != null) {
161-
effectiveTransport = PeerInfo.TransportType.TRANSPORT_TYPE_CLOUD_PATH;
162-
}
163-
}
164-
if (effectiveTransport != PeerInfo.TransportType.TRANSPORT_TYPE_UNKNOWN) {
165-
peerInfo = peerInfo.toBuilder()
166-
.setTransportType(effectiveTransport)
167-
.build();
168-
}
169-
return peerInfo;
170-
} catch (Exception e) {
171-
throw new IllegalArgumentException(
172-
"Failed to parse "
173-
+ PEER_INFO_KEY.name()
174-
+ " from the response header value: "
175-
+ encodedStr);
176-
}
168+
if (effectiveTransport != PeerInfo.TransportType.TRANSPORT_TYPE_UNKNOWN) {
169+
peerInfo = peerInfo.toBuilder().setTransportType(effectiveTransport).build();
177170
}
171+
return peerInfo;
172+
} catch (Exception e) {
173+
throw new IllegalArgumentException(
174+
"Failed to parse "
175+
+ PEER_INFO_KEY.name()
176+
+ " from the response header value: "
177+
+ encodedStr);
178+
}
179+
}
178180

179-
180-
@Nullable
181-
private static ResponseParams extractResponseParams(Metadata metadata) {
182-
byte[] responseParams = metadata.get(LOCATION_METADATA_KEY);
183-
if (responseParams != null) {
184-
try {
185-
return ResponseParams.parseFrom(responseParams);
186-
} catch (InvalidProtocolBufferException e) {
187-
}
188-
}
189-
return null;
181+
@Nullable
182+
private static ResponseParams extractResponseParams(Metadata metadata) {
183+
byte[] responseParams = metadata.get(LOCATION_METADATA_KEY);
184+
if (responseParams != null) {
185+
try {
186+
return ResponseParams.parseFrom(responseParams);
187+
} catch (InvalidProtocolBufferException e) {
190188
}
189+
}
190+
return null;
191191
}
192+
}
192193
}

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableGrpcStreamTracer.java

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,10 +15,8 @@
1515
*/
1616
package com.google.cloud.bigtable.data.v2.stub.metrics;
1717

18-
import com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsTracer.TransportAttrs;
1918
import io.grpc.ClientStreamTracer;
2019
import io.grpc.Metadata;
21-
import io.grpc.Status;
2220

2321
/**
2422
* Records the time a request is enqueued in a grpc channel queue. This a bridge between gRPC stream

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableTracer.java

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121
import com.google.api.gax.tracing.ApiTracer;
2222
import com.google.api.gax.tracing.BaseApiTracer;
2323
import com.google.cloud.bigtable.data.v2.stub.MetadataExtractorInterceptor;
24-
2524
import java.time.Duration;
2625
import javax.annotation.Nullable;
2726

@@ -76,6 +75,7 @@ public int getAttempt() {
7675
* Record the latency between Google's network receives the RPC and reads back the first byte of
7776
* the response from server-timing header. If server-timing header is missing, increment the
7877
* missing header count.
78+
*
7979
* @deprecated Use {@link #setSidebandData(MetadataExtractorInterceptor.SidebandData)}
8080
*/
8181
@Deprecated
@@ -91,17 +91,19 @@ public void batchRequestThrottled(long throttledTimeMs) {
9191
public void setSidebandData(MetadataExtractorInterceptor.SidebandData sidebandData) {
9292
// noop
9393
}
94+
9495
/**
9596
* Set the Bigtable zone and cluster so metrics can be tagged with location information. This will
9697
* be called in BuiltinMetricsTracer.
98+
*
9799
* @deprecated Use {@link #setSidebandData(MetadataExtractorInterceptor.SidebandData)}
98100
*/
99101
@Deprecated
100-
public void setLocations(String zone, String cluster) {
101-
}
102+
public void setLocations(String zone, String cluster) {}
102103

103104
/**
104105
* Set the underlying transport used to process the attempt
106+
*
105107
* @deprecated Use {@link #setSidebandData(MetadataExtractorInterceptor.SidebandData)}
106108
*/
107109
@Deprecated

google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableTracerStreamingCallable.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -26,8 +26,6 @@
2626
import com.google.cloud.bigtable.data.v2.stub.SafeResponseObserver;
2727
import com.google.common.base.Preconditions;
2828
import com.google.common.base.Stopwatch;
29-
import io.grpc.ClientInterceptors;
30-
3129
import java.util.concurrent.TimeUnit;
3230
import javax.annotation.Nonnull;
3331

@@ -63,15 +61,17 @@ public void call(
6361
MetadataExtractorInterceptor metadataExtractor = new MetadataExtractorInterceptor();
6462
grpcCtx = metadataExtractor.injectInto(grpcCtx);
6563

66-
6764
// tracer should always be an instance of bigtable tracer
6865
if (context.getTracer() instanceof BigtableTracer) {
6966
BigtableTracer tracer = (BigtableTracer) context.getTracer();
7067
grpcCtx.withCallOptions(
71-
grpcCtx.getCallOptions().withStreamTracerFactory(new BigtableGrpcStreamTracer.Factory(tracer)));
68+
grpcCtx
69+
.getCallOptions()
70+
.withStreamTracerFactory(new BigtableGrpcStreamTracer.Factory(tracer)));
7271

7372
BigtableTracerResponseObserver<ResponseT> innerObserver =
74-
new BigtableTracerResponseObserver<>(responseObserver, tracer, metadataExtractor.getSidebandData());
73+
new BigtableTracerResponseObserver<>(
74+
responseObserver, tracer, metadataExtractor.getSidebandData());
7575
if (context.getRetrySettings() != null) {
7676
tracer.setTotalTimeoutDuration(context.getRetrySettings().getTotalTimeoutDuration());
7777
}

0 commit comments

Comments
 (0)