Skip to content

Commit 0e8ae30

Browse files
xds: Ext proc client drop backend_service metrics label, make eos_without_message conditional on eos field (#13028)
Implementing dropping backend_service metrics label, and the discrepancy with the gRFC in `end_of_stream_without_message` handling. The code treats this field independently, but as per the gRFC it should be only be considered when the `end_of_stream` field is true. This applies both in `ProcessingRequest` and `ProcessingResponse`.
1 parent 61b910b commit 0e8ae30

2 files changed

Lines changed: 654 additions & 94 deletions

File tree

‎xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java‎

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -132,7 +132,7 @@ static synchronized void initMetricInstruments() {
132132
"s",
133133
LATENCY_BUCKETS,
134134
ImmutableList.of("grpc.target"),
135-
ImmutableList.of("grpc.lb.backend_service"),
135+
ImmutableList.of(),
136136
true);
137137

138138
clientHalfCloseDuration = registry.registerDoubleHistogram(
@@ -142,7 +142,7 @@ static synchronized void initMetricInstruments() {
142142
"s",
143143
LATENCY_BUCKETS,
144144
ImmutableList.of("grpc.target"),
145-
ImmutableList.of("grpc.lb.backend_service"),
145+
ImmutableList.of(),
146146
true);
147147

148148
serverHeadersDuration = registry.registerDoubleHistogram(
@@ -152,7 +152,7 @@ static synchronized void initMetricInstruments() {
152152
"s",
153153
LATENCY_BUCKETS,
154154
ImmutableList.of("grpc.target"),
155-
ImmutableList.of("grpc.lb.backend_service"),
155+
ImmutableList.of(),
156156
true);
157157

158158
serverTrailersDuration = registry.registerDoubleHistogram(
@@ -162,7 +162,7 @@ static synchronized void initMetricInstruments() {
162162
"s",
163163
LATENCY_BUCKETS,
164164
ImmutableList.of("grpc.target"),
165-
ImmutableList.of("grpc.lb.backend_service"),
165+
ImmutableList.of(),
166166
true);
167167
}
168168
}
@@ -256,8 +256,7 @@ public void cancel(@Nullable String message, @Nullable Throwable cause) {
256256

257257
DataPlaneClientCall dataPlaneCall = new DataPlaneClientCall(
258258
delayedCall, rawCall, extProcStub, filterConfig, filterConfig.getMutationRulesConfig(),
259-
scheduler, rawMethod, next, metricsRecorder, next.authority(),
260-
callOptions.getOption(XdsNameResolver.CLUSTER_SELECTION_KEY));
259+
scheduler, rawMethod, next, metricsRecorder, next.authority());
261260

262261
return (ClientCall<ReqT, RespT>) (ClientCall<?, ?>) dataPlaneCall;
263262
}
@@ -350,7 +349,6 @@ private static class DataPlaneClientCall
350349
private final Channel channel;
351350
private final MetricRecorder metricsRecorder;
352351
private final String target;
353-
private final String backendService;
354352
private volatile Context callContext = Context.ROOT;
355353

356354
private volatile long clientHeadersStartNanos;
@@ -383,8 +381,7 @@ protected DataPlaneClientCall(
383381
MethodDescriptor<?, ?> method,
384382
Channel channel,
385383
MetricRecorder metricsRecorder,
386-
String target,
387-
String backendService) {
384+
String target) {
388385
super(delayedCall);
389386
this.delayedCall = delayedCall;
390387
this.rawCall = rawCall;
@@ -397,7 +394,6 @@ protected DataPlaneClientCall(
397394
this.channel = channel;
398395
this.metricsRecorder = checkNotNull(metricsRecorder, "metricsRecorder");
399396
this.target = checkNotNull(target, "target");
400-
this.backendService = checkNotNull(backendService, "backendService");
401397
}
402398

403399
private boolean activateCall() {
@@ -429,7 +425,7 @@ private void recordDuration(DoubleHistogramMetricInstrument instrument, long dur
429425
instrument,
430426
durationSecs,
431427
ImmutableList.of(target),
432-
ImmutableList.of(backendService));
428+
ImmutableList.of());
433429
}
434430
}
435431

@@ -1134,6 +1130,7 @@ public void halfClose() {
11341130

11351131
ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
11361132
.setRequestBody(HttpBody.newBuilder()
1133+
.setEndOfStream(true)
11371134
.setEndOfStreamWithoutMessage(true)
11381135
.build());
11391136
mergeAccumulatedWindowUpdates(builder);
@@ -1167,7 +1164,10 @@ private void handleRequestBodyResponse(BodyResponse bodyResponse) {
11671164
BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
11681165
if (mutation.hasStreamedResponse()) {
11691166
StreamedBodyResponse streamed = mutation.getStreamedResponse();
1170-
if (!streamed.getEndOfStreamWithoutMessage()) {
1167+
boolean isEndOfStream = streamed.getEndOfStream();
1168+
boolean isEndOfStreamWithoutMessage =
1169+
isEndOfStream && streamed.getEndOfStreamWithoutMessage();
1170+
if (!isEndOfStreamWithoutMessage) {
11711171
ByteString body = streamed.getBody();
11721172
boolean sendImmediately = false;
11731173
synchronized (streamLock) {
@@ -1184,7 +1184,7 @@ private void handleRequestBodyResponse(BodyResponse bodyResponse) {
11841184
trySendAccumulatedWindowUpdates();
11851185
}
11861186
}
1187-
if (streamed.getEndOfStream() || streamed.getEndOfStreamWithoutMessage()) {
1187+
if (isEndOfStream) {
11881188
synchronized (streamLock) {
11891189
if (pendingUpstreamBodyMessages.isEmpty()) {
11901190
if (requestSideClosed.compareAndSet(false, true)) {

0 commit comments

Comments
 (0)