• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

grpc / grpc-java / #20413

20 Aug 2026 05:04AM UTC coverage: 89.2% (+0.006%) from 89.194%
#20413

push

github

web-flow
xds: Prevent concurrent cancellations of data plane call in ExternalProcessorClientInterceptor (#12996)

If the ext-proc stream terminates with a non-OK status, the interceptor
cancels the downstream data plane call (on the executor thread).
Concurrently, the application thread may call cancel() on the returned
proxy call (e.g. for user cancellation or cleanup).

Since ClientCallImpl does not synchronize its cancelCalled field,
concurrent cancellations from these two threads resulted in a TSAN data
race (that was discussed in
https://github.com/grpc/grpc-java/pull/12975).

This commit introduces an AtomicBoolean downstreamCancelled in
DataPlaneClientCall to guard all downstream cancellations, ensuring that
only the first cancellation is forwarded to delayedCall/super.cancel().
This deduplicates and serializes cancellation handling from both the
application and the interceptor threads, preventing the data race.

38554 of 43222 relevant lines covered (89.2%)

0.89 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

96.14
/../xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java
1
/*
2
 * Copyright 2024 The gRPC Authors
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 */
16

17
package io.grpc.xds;
18

19
import static com.google.common.base.Preconditions.checkNotNull;
20
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.applyHeaderMutations;
21
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.collectAttributes;
22
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.markDataPlaneCallClosed;
23
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.markExtProcStreamCompleted;
24
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.markExtProcStreamFailed;
25
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.outboundStreamToByteString;
26
import static io.grpc.xds.internal.extproc.ExternalProcessorUtil.toHeaderMap;
27

28
import com.google.common.annotations.VisibleForTesting;
29
import com.google.common.collect.ImmutableList;
30
import com.google.common.collect.ImmutableMap;
31
import com.google.protobuf.ByteString;
32
import com.google.protobuf.Struct;
33
import io.envoyproxy.envoy.extensions.filters.http.ext_proc.v3.ProcessingMode;
34
import io.envoyproxy.envoy.service.ext_proc.v3.BodyMutation;
35
import io.envoyproxy.envoy.service.ext_proc.v3.BodyResponse;
36
import io.envoyproxy.envoy.service.ext_proc.v3.CommonResponse;
37
import io.envoyproxy.envoy.service.ext_proc.v3.ExternalProcessorGrpc;
38
import io.envoyproxy.envoy.service.ext_proc.v3.HttpBody;
39
import io.envoyproxy.envoy.service.ext_proc.v3.HttpHeaders;
40
import io.envoyproxy.envoy.service.ext_proc.v3.HttpTrailers;
41
import io.envoyproxy.envoy.service.ext_proc.v3.ImmediateResponse;
42
import io.envoyproxy.envoy.service.ext_proc.v3.ProcessingRequest;
43
import io.envoyproxy.envoy.service.ext_proc.v3.ProcessingResponse;
44
import io.envoyproxy.envoy.service.ext_proc.v3.ProtocolConfiguration;
45
import io.envoyproxy.envoy.service.ext_proc.v3.StreamedBodyResponse;
46
import io.grpc.CallOptions;
47
import io.grpc.Channel;
48
import io.grpc.ClientCall;
49
import io.grpc.ClientInterceptor;
50
import io.grpc.Context;
51
import io.grpc.Deadline;
52
import io.grpc.DoubleHistogramMetricInstrument;
53
import io.grpc.ForwardingClientCall.SimpleForwardingClientCall;
54
import io.grpc.ForwardingClientCallListener.SimpleForwardingClientCallListener;
55
import io.grpc.ManagedChannel;
56
import io.grpc.Metadata;
57
import io.grpc.MethodDescriptor;
58
import io.grpc.MetricInstrumentRegistry;
59
import io.grpc.MetricRecorder;
60
import io.grpc.Status;
61
import io.grpc.StatusRuntimeException;
62
import io.grpc.internal.DelayedClientCall;
63
import io.grpc.internal.SerializingExecutor;
64
import io.grpc.stub.ClientCallStreamObserver;
65
import io.grpc.stub.ClientResponseObserver;
66
import io.grpc.xds.ExternalProcessorFilter.ExternalProcessorFilterConfig;
67
import io.grpc.xds.Filter.FilterContext;
68
import io.grpc.xds.internal.extproc.DataPlaneCallState;
69
import io.grpc.xds.internal.extproc.EventType;
70
import io.grpc.xds.internal.extproc.ExtProcStreamState;
71
import io.grpc.xds.internal.extproc.KnownLengthInputStream;
72
import io.grpc.xds.internal.grpcservice.CachedChannelManager;
73
import io.grpc.xds.internal.grpcservice.HeaderValue;
74
import io.grpc.xds.internal.headermutations.HeaderMutationDisallowedException;
75
import io.grpc.xds.internal.headermutations.HeaderMutationFilter;
76
import io.grpc.xds.internal.headermutations.HeaderMutationRulesConfig;
77
import io.grpc.xds.internal.headermutations.HeaderMutator;
78
import java.io.IOException;
79
import java.io.InputStream;
80
import java.util.List;
81
import java.util.Optional;
82
import java.util.Queue;
83
import java.util.concurrent.ConcurrentLinkedQueue;
84
import java.util.concurrent.Executor;
85
import java.util.concurrent.ScheduledExecutorService;
86
import java.util.concurrent.ScheduledFuture;
87
import java.util.concurrent.TimeUnit;
88
import java.util.concurrent.atomic.AtomicBoolean;
89
import java.util.concurrent.atomic.AtomicInteger;
90
import java.util.concurrent.atomic.AtomicReference;
91
import javax.annotation.Nullable;
92

93
/**
94
 * Client-side interceptor for external processing filter.
95
 */
96
final class ExternalProcessorClientInterceptor implements ClientInterceptor {
97

98
  @VisibleForTesting
99
  static DoubleHistogramMetricInstrument clientHeadersDuration;
100
  @VisibleForTesting
101
  static DoubleHistogramMetricInstrument clientHalfCloseDuration;
102
  @VisibleForTesting
103
  static DoubleHistogramMetricInstrument serverHeadersDuration;
104
  @VisibleForTesting
105
  static DoubleHistogramMetricInstrument serverTrailersDuration;
106

107
  // Copied from io.grpc.opentelemetry.internal.OpenTelemetryConstants.LATENCY_BUCKETS
108
  private static final List<Double> LATENCY_BUCKETS = ImmutableList.of(
1 ✔
109
      0d,     0.00001d, 0.00005d, 0.0001d, 0.0003d, 0.0006d, 0.0008d, 0.001d, 0.002d,
1 ✔
110
      0.003d, 0.004d,   0.005d,   0.006d,  0.008d,  0.01d,   0.013d,  0.016d, 0.02d,
1 ✔
111
      0.025d, 0.03d,    0.04d,    0.05d,   0.065d,  0.08d,   0.1d,    0.13d,  0.16d,
1 ✔
112
      0.2d,   0.25d,    0.3d,     0.4d,    0.5d,    0.65d,   0.8d,    1d,     2d,
1 ✔
113
      5d,     10d,      20d,      50d,     100d);
1 ✔
114

115
  static {
116
    initMetricInstruments();
1 ✔
117
  }
1 ✔
118

119
  static synchronized void initMetricInstruments() {
120
    if (io.grpc.internal.GrpcUtil.getFlag("GRPC_EXPERIMENTAL_XDS_EXT_PROC_ON_CLIENT", false)) {
1 ✔
121
      if (clientHeadersDuration == null) {
1 ✔
122
        MetricInstrumentRegistry registry = MetricInstrumentRegistry.getDefaultRegistry();
1 ✔
123

124
        clientHeadersDuration = registry.registerDoubleHistogram(
1 ✔
125
            "grpc.client_ext_proc.client_headers_duration",
126
            "Time between when the ext_proc filter sees the client's headers and when "
127
                + "it allows those headers to continue on to the next filter",
128
            "s",
129
            LATENCY_BUCKETS,
130
            ImmutableList.of("grpc.target"),
1 ✔
131
            ImmutableList.of("grpc.lb.backend_service"),
1 ✔
132
            true);
133

134
        clientHalfCloseDuration = registry.registerDoubleHistogram(
1 ✔
135
            "grpc.client_ext_proc.client_half_close_duration",
136
            "Time between when the ext_proc filter sees the client's half-close and when "
137
                + "it allows that half-close to continue on to the next filter",
138
            "s",
139
            LATENCY_BUCKETS,
140
            ImmutableList.of("grpc.target"),
1 ✔
141
            ImmutableList.of("grpc.lb.backend_service"),
1 ✔
142
            true);
143

144
        serverHeadersDuration = registry.registerDoubleHistogram(
1 ✔
145
            "grpc.client_ext_proc.server_headers_duration",
146
            "Time between when the ext_proc filter sees the server's headers and when "
147
                + "it allows those headers to continue on to the next filter",
148
            "s",
149
            LATENCY_BUCKETS,
150
            ImmutableList.of("grpc.target"),
1 ✔
151
            ImmutableList.of("grpc.lb.backend_service"),
1 ✔
152
            true);
153

154
        serverTrailersDuration = registry.registerDoubleHistogram(
1 ✔
155
            "grpc.client_ext_proc.server_trailers_duration",
156
            "Time between when the ext_proc filter sees the server's trailers and when "
157
                + "it allows those trailers to continue on to the next filter",
158
            "s",
159
            LATENCY_BUCKETS,
160
            ImmutableList.of("grpc.target"),
1 ✔
161
            ImmutableList.of("grpc.lb.backend_service"),
1 ✔
162
            true);
163
      }
164
    }
165
  }
1 ✔
166

167
  private final ExternalProcessorFilterConfig filterConfig;
168
  private final ScheduledExecutorService scheduler;
169
  private final MetricRecorder metricsRecorder;
170
  private final ManagedChannel extProcChannel;
171

172
  @VisibleForTesting
173
  ExternalProcessorClientInterceptor(ExternalProcessorFilterConfig filterConfig,
174
      CachedChannelManager cachedChannelManager,
175
      ScheduledExecutorService scheduler,
176
      FilterContext context) {
1 ✔
177
    this.filterConfig = filterConfig;
1 ✔
178
    checkNotNull(cachedChannelManager, "cachedChannelManager");
1 ✔
179
    this.scheduler = checkNotNull(scheduler, "scheduler");
1 ✔
180
    this.metricsRecorder = checkNotNull(context.metricsRecorder(), "metricsRecorder");
1 ✔
181
    this.extProcChannel = cachedChannelManager.getChannel(filterConfig.getGrpcServiceConfig());
1 ✔
182
  }
1 ✔
183

184
  @VisibleForTesting
185
  ExternalProcessorFilterConfig getFilterConfig() {
186
    return filterConfig;
1 ✔
187
  }
188

189
  @VisibleForTesting
190
  ManagedChannel getExtProcChannel() {
191
    return extProcChannel;
×
192
  }
193

194
  @Override
195
  @SuppressWarnings("unchecked")
196
  public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
197
      MethodDescriptor<ReqT, RespT> method,
198
      CallOptions callOptions,
199
      Channel next) {
200
    SerializingExecutor serializingExecutor = new SerializingExecutor(callOptions.getExecutor());
1 ✔
201
    
202
    ExternalProcessorGrpc.ExternalProcessorStub extProcStub = ExternalProcessorGrpc.newStub(
1 ✔
203
        extProcChannel)
204
        .withExecutor(serializingExecutor);
1 ✔
205
    
206
    if (filterConfig.getGrpcServiceConfig().timeout().isPresent()) {
1 ✔
207
      long timeoutNanos = filterConfig.getGrpcServiceConfig().timeout().get().toNanos();
1 ✔
208
      if (timeoutNanos > 0) {
1 ✔
209
        extProcStub = extProcStub.withDeadlineAfter(timeoutNanos, TimeUnit.NANOSECONDS);
1 ✔
210
      }
211
    }
212
    if (filterConfig.getGrpcServiceConfig().initialMetadata() != null
1 ✔
213
        && !filterConfig.getGrpcServiceConfig().initialMetadata().isEmpty()) {
1 ✔
214
      Metadata extraHeaders = new Metadata();
1 ✔
215
      for (HeaderValue headerValue : filterConfig.getGrpcServiceConfig().initialMetadata()) {
1 ✔
216
        String key = headerValue.key();
1 ✔
217
        if (key.endsWith(Metadata.BINARY_HEADER_SUFFIX)) {
1 ✔
218
          if (headerValue.rawValue().isPresent()) {
1 ✔
219
            Metadata.Key<byte[]> metadataKey =
1 ✔
220
                Metadata.Key.of(key, Metadata.BINARY_BYTE_MARSHALLER);
1 ✔
221
            extraHeaders.put(metadataKey, headerValue.rawValue().get().toByteArray());
1 ✔
222
          }
1 ✔
223
        } else {
224
          if (headerValue.value().isPresent()) {
1 ✔
225
            Metadata.Key<String> metadataKey =
1 ✔
226
                Metadata.Key.of(key, Metadata.ASCII_STRING_MARSHALLER);
1 ✔
227
            extraHeaders.put(metadataKey, headerValue.value().get());
1 ✔
228
          }
229
        }
230
      }
1 ✔
231
      extProcStub = extProcStub.withInterceptors(
1 ✔
232
          io.grpc.stub.MetadataUtils.newAttachHeadersInterceptor(extraHeaders));
1 ✔
233
    }
234

235

236
    // The filter chain is preceded by RawMessageClientInterceptor, so ReqT and RespT are
237
    // InputStream.
238
    MethodDescriptor<InputStream, InputStream> rawMethod =
1 ✔
239
        (MethodDescriptor<InputStream, InputStream>) (MethodDescriptor<?, ?>) method;
240
    ClientCall<InputStream, InputStream> rawCall =
1 ✔
241
        (ClientCall<InputStream, InputStream>) (ClientCall<?, ?>)
242
            next.newCall(method, callOptions);
1 ✔
243

244
    // Create a local subclass instance to buffer outbound actions
245
    DataPlaneDelayedCall<InputStream, InputStream> delayedCall =
1 ✔
246
        new DataPlaneDelayedCall<>(
247
            serializingExecutor, scheduler, callOptions.getDeadline());
1 ✔
248

249
    DataPlaneClientCall dataPlaneCall = new DataPlaneClientCall(
1 ✔
250
        delayedCall, rawCall, extProcStub, filterConfig, filterConfig.getMutationRulesConfig(),
1 ✔
251
        scheduler, rawMethod, next, metricsRecorder, next.authority(),
1 ✔
252
        callOptions.getOption(XdsNameResolver.CLUSTER_SELECTION_KEY));
1 ✔
253

254
    return (ClientCall<ReqT, RespT>) (ClientCall<?, ?>) dataPlaneCall;
1 ✔
255
  }
256

257
  // --- SHARED UTILITY METHODS ---
258

259
  /**
260
   * A local subclass to expose the protected constructor of DelayedClientCall.
261
   */
262
  private static class DataPlaneDelayedCall<ReqT, RespT> extends DelayedClientCall<ReqT, RespT> {
263
    DataPlaneDelayedCall(
264
        Executor executor, ScheduledExecutorService scheduler, @Nullable Deadline deadline) {
265
      super("ext_proc", executor, scheduler, deadline);
1 ✔
266
    }
1 ✔
267
  }
268

269
  /**
270
   * Handles the bidirectional stream with the External Processor.
271
   * Buffers the actual RPC start until the Ext Proc header response is received.
272
   */
273
  private static class DataPlaneClientCall 
274
      extends SimpleForwardingClientCall<InputStream, InputStream> {
275

276
    private final ExternalProcessorGrpc.ExternalProcessorStub stub;
277
    private final ExternalProcessorFilterConfig config;
278
    private final ClientCall<InputStream, InputStream> rawCall;
279
    private final DataPlaneDelayedCall<InputStream, InputStream> delayedCall;
280
    private final ScheduledExecutorService scheduler;
281
    private final Object streamLock = new Object();
1 ✔
282
    @Nullable private volatile EventType expectedRequestResponse;
283
    @Nullable private volatile EventType expectedResponseResponse;
284
    @Nullable private volatile ClientCallStreamObserver<ProcessingRequest>
285
        extProcClientCallRequestObserver;
286
    private final Queue<InputStream> pendingDrainingMessages =
1 ✔
287
        new ConcurrentLinkedQueue<>();
288
    @Nullable private volatile DataPlaneListener wrappedListener;
289
    private final HeaderMutationFilter mutationFilter;
290
    private final HeaderMutator mutator = HeaderMutator.create();
1 ✔
291
    private final AtomicInteger pendingRequests = new AtomicInteger(0);
1 ✔
292
    private final ProcessingMode currentProcessingMode;
293
    private final MethodDescriptor<?, ?> method;
294
    private final Channel channel;
295
    private final MetricRecorder metricsRecorder;
296
    private final String target;
297
    private final String backendService;
298
    private volatile Context callContext = Context.ROOT;
1 ✔
299

300
    private long clientHeadersStartNanos;
301
    private long clientHalfCloseStartNanos;
302
    private long serverHeadersStartNanos;
303
    private long serverTrailersStartNanos;
304

305
    private boolean protocolConfigSent = false;
1 ✔
306
    private ImmutableMap<String, Struct> collectedAttributes;
307
    private boolean requestAttributesSent = false;
1 ✔
308
    @Nullable private volatile Metadata requestHeaders;
309
    final AtomicReference<DataPlaneCallState> dataPlaneCallState =
1 ✔
310
        new AtomicReference<>(DataPlaneCallState.IDLE);
311
    final AtomicReference<ExtProcStreamState> extProcStreamState =
1 ✔
312
        new AtomicReference<>(ExtProcStreamState.ACTIVE);
313
    final AtomicBoolean passThroughMode = new AtomicBoolean(false);
1 ✔
314
    final AtomicBoolean requestSideClosed = new AtomicBoolean(false);
1 ✔
315
    final AtomicBoolean isProcessingTrailers = new AtomicBoolean(false);
1 ✔
316
    final AtomicBoolean pendingHalfClose = new AtomicBoolean(false);
1 ✔
317
    final AtomicBoolean bodyMessageSentToExtProc = new AtomicBoolean(false);
1 ✔
318
    private final AtomicBoolean downstreamCancelled = new AtomicBoolean(false);
1 ✔
319

320
    protected DataPlaneClientCall(
321
        DataPlaneDelayedCall<InputStream, InputStream> delayedCall,
322
        ClientCall<InputStream, InputStream> rawCall,
323
        ExternalProcessorGrpc.ExternalProcessorStub stub,
324
        ExternalProcessorFilterConfig config,
325
        Optional<HeaderMutationRulesConfig> mutationRulesConfig,
326
        ScheduledExecutorService scheduler,
327
        MethodDescriptor<?, ?> method,
328
        Channel channel,
329
        MetricRecorder metricsRecorder,
330
        String target,
331
        String backendService) {
332
      super(delayedCall);
1 ✔
333
      this.delayedCall = delayedCall;
1 ✔
334
      this.rawCall = rawCall;
1 ✔
335
      this.stub = stub;
1 ✔
336
      this.config = config;
1 ✔
337
      this.currentProcessingMode = config.getExternalProcessor().getProcessingMode();
1 ✔
338
      this.mutationFilter = new HeaderMutationFilter(mutationRulesConfig);
1 ✔
339
      this.scheduler = scheduler;
1 ✔
340
      this.method = method;
1 ✔
341
      this.channel = channel;
1 ✔
342
      this.metricsRecorder = checkNotNull(metricsRecorder, "metricsRecorder");
1 ✔
343
      this.target = checkNotNull(target, "target");
1 ✔
344
      this.backendService = checkNotNull(backendService, "backendService");
1 ✔
345
    }
1 ✔
346

347

348

349
    private void activateCall() {
350
      if ((extProcStreamState.get() == ExtProcStreamState.FAILED
1 ✔
351
              && !config.getFailureModeAllow()
1 ✔
352
              && !config.getObservabilityMode())
1 ✔
353
          || !dataPlaneCallState.compareAndSet(
1 ✔
354
              DataPlaneCallState.IDLE, DataPlaneCallState.ACTIVE)) {
355
        return;
1 ✔
356
      }
357
      if (clientHeadersStartNanos > 0) {
1 ✔
358
        long durationNanos = System.nanoTime() - clientHeadersStartNanos;
1 ✔
359
        recordDuration(clientHeadersDuration, durationNanos);
1 ✔
360
        clientHeadersStartNanos = 0;
1 ✔
361
      }
362
      Runnable toRun = delayedCall.setCall(rawCall);
1 ✔
363
      if (toRun != null) {
1 ✔
364
        callContext.run(toRun);
1 ✔
365
      }
366
      drainPendingRequests();
1 ✔
367
      onReadyNotify();
1 ✔
368
    }
1 ✔
369

370
    private void recordDuration(DoubleHistogramMetricInstrument instrument, long durationNanos) {
371
      if (instrument != null) {
1 ✔
372
        double durationSecs = (double) durationNanos / 1_000_000_000.0;
1 ✔
373
        metricsRecorder.recordDoubleHistogram(
1 ✔
374
            instrument,
375
            durationSecs,
376
            ImmutableList.of(target),
1 ✔
377
            ImmutableList.of(backendService));
1 ✔
378
      }
379
    }
1 ✔
380

381
    /**
382
     * Validates whether the body response uses unsupported gRPC message compression.
383
     * If compression is unsupported, this method will cancel the call, transition the
384
     * stream to a failed state, send an error to the external processor, and return false.
385
     *
386
     * @param bodyResponse the response to validate
387
     * @return true if validation passes (compression is supported or not used),
388
     *     false if validation fails
389
     */
390
    private boolean validateCompressionSupport(BodyResponse bodyResponse) {
391
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1 ✔
392
        BodyMutation mutation = 
1 ✔
393
            bodyResponse.getResponse().getBodyMutation();
1 ✔
394
        if (mutation.hasStreamedResponse()
1 ✔
395
            && mutation.getStreamedResponse().getGrpcMessageCompressed()) {
1 ✔
396
          StatusRuntimeException ex = Status.UNAVAILABLE
1 ✔
397
              .withDescription("gRPC message compression not supported in ext_proc")
1 ✔
398
              .asRuntimeException();
1 ✔
399
          synchronized (streamLock) {
1 ✔
400
            if (!extProcStreamState.get().isCompleted()
1 ✔
401
                && extProcClientCallRequestObserver != null) {
402
              extProcClientCallRequestObserver.onError(ex);
1 ✔
403
            }
404
          }
1 ✔
405
          activateCall();
1 ✔
406
          markExtProcStreamFailed(extProcStreamState);
1 ✔
407
          cancelDownstream("gRPC message compression not supported in ext_proc", ex);
1 ✔
408
          closeExtProcStream();
1 ✔
409
          return false;
1 ✔
410
        }
411
      }
412
      return true;
1 ✔
413
    }
414

415

416

417
    @Override
418
    public void start(Listener<InputStream> responseListener, Metadata headers) {
419
      this.callContext = Context.current();
1 ✔
420
      clientHeadersStartNanos = System.nanoTime();
1 ✔
421
      this.requestHeaders = headers;
1 ✔
422
      this.wrappedListener = new DataPlaneListener(responseListener, rawCall, this);
1 ✔
423

424
      // DelayedClientCall.start will buffer the listener and headers until setCall is called.
425
      super.start(wrappedListener, headers);
1 ✔
426

427
      stub.process(new ClientResponseObserver<ProcessingRequest, ProcessingResponse>() {
1 ✔
428
        @Override
429
        public void beforeStart(ClientCallStreamObserver<ProcessingRequest> requestStream) {
430
          synchronized (streamLock) {
1 ✔
431
            extProcClientCallRequestObserver = requestStream;
1 ✔
432
          }
1 ✔
433
          requestStream.setOnReadyHandler(DataPlaneClientCall.this::onExtProcStreamReady);
1 ✔
434
        }
1 ✔
435

436
        @Override
437
        public void onNext(ProcessingResponse response) {
438
          try {
439
            if (config.getObservabilityMode()) {
1 ✔
440
              return;
1 ✔
441
            }
442

443
            if (response.hasImmediateResponse()) {
1 ✔
444
              if (config.getDisableImmediateResponse()) {
1 ✔
445
                internalOnError(Status.UNAVAILABLE
1 ✔
446
                    .withDescription(
1 ✔
447
                        "Immediate response is disabled but received from external processor")
448
                    .asRuntimeException());
1 ✔
449
                return;
1 ✔
450
              }
451
              handleImmediateResponse(response.getImmediateResponse(), wrappedListener);
1 ✔
452
              return;
1 ✔
453
            }
454

455
            if (response.hasRequestHeaders()) {
1 ✔
456
              EventType expected = expectedRequestResponse;
1 ✔
457
              if (expected == null || expected != EventType.REQUEST_HEADERS) {
1 ✔
458
                internalOnError(Status.UNAVAILABLE
×
459
                    .withDescription("Protocol error: received response out of order. Expected: " 
×
460
                        + expected + ", Received: REQUEST_HEADERS")
461
                    .asRuntimeException());
×
462
                return;
×
463
              }
464
              expectedRequestResponse = null;
1 ✔
465
            } else if (response.hasResponseHeaders()) {
1 ✔
466
              EventType expected = expectedResponseResponse;
1 ✔
467
              if (expected == null || expected != EventType.RESPONSE_HEADERS) {
1 ✔
468
                internalOnError(Status.UNAVAILABLE
1 ✔
469
                    .withDescription("Protocol error: received response out of order. Expected: " 
1 ✔
470
                        + expected + ", Received: RESPONSE_HEADERS")
471
                    .asRuntimeException());
1 ✔
472
                return;
1 ✔
473
              }
474
              expectedResponseResponse = null;
1 ✔
475
            } else if (response.hasResponseTrailers()) {
1 ✔
476
              EventType expected = expectedResponseResponse;
1 ✔
477
              if (expected == null || expected != EventType.RESPONSE_TRAILERS) {
1 ✔
478
                internalOnError(Status.UNAVAILABLE
1 ✔
479
                    .withDescription("Protocol error: received response out of order. Expected: " 
1 ✔
480
                        + expected + ", Received: RESPONSE_TRAILERS")
481
                    .asRuntimeException());
1 ✔
482
                return;
1 ✔
483
              }
484
              expectedResponseResponse = null;
1 ✔
485
            } else if (response.hasRequestBody()) {
1 ✔
486
              EventType expected = expectedRequestResponse;
1 ✔
487
              if (expected == EventType.REQUEST_HEADERS) {
1 ✔
488
                internalOnError(Status.UNAVAILABLE
1 ✔
489
                    .withDescription(
1 ✔
490
                        "Protocol error: received request_body before request_headers response.")
491
                    .asRuntimeException());
1 ✔
492
                return;
1 ✔
493
              }
494
            } else if (response.hasResponseBody()) {
1 ✔
495
              EventType expected = expectedResponseResponse;
1 ✔
496
              if (expected == EventType.RESPONSE_HEADERS) {
1 ✔
497
                internalOnError(Status.UNAVAILABLE
1 ✔
498
                    .withDescription(
1 ✔
499
                        "Protocol error: received response_body before headers response.")
500
                    .asRuntimeException());
1 ✔
501
                return;
1 ✔
502
              }
503
            }
504

505
            if (response.getRequestDrain()) {
1 ✔
506
              extProcStreamState.set(ExtProcStreamState.DRAINING);
1 ✔
507
              halfCloseExtProcStream();
1 ✔
508
              activateCall();
1 ✔
509
            }
510

511
            // 1. Client Headers
512
            if (response.hasRequestHeaders()) {
1 ✔
513
              if (response.getRequestHeaders().hasResponse()) {
1 ✔
514
                if (response.getRequestHeaders().getResponse().getStatus()
1 ✔
515
                    == CommonResponse.ResponseStatus.CONTINUE_AND_REPLACE) {
516
                  internalOnError(Status.UNAVAILABLE
1 ✔
517
                      .withDescription("CONTINUE_AND_REPLACE is not supported")
1 ✔
518
                      .asRuntimeException());
1 ✔
519
                  return;
1 ✔
520
                }
521
                applyHeaderMutations(
1 ✔
522
                    requestHeaders,
1 ✔
523
                    response.getRequestHeaders().getResponse().getHeaderMutation(),
1 ✔
524
                    mutationFilter,
1 ✔
525
                    mutator);
1 ✔
526
              }
527
              activateCall();
1 ✔
528
            }
529
            // 2. Client Message (Request Body)
530
            else if (response.hasRequestBody()) {
1 ✔
531
              if (validateCompressionSupport(response.getRequestBody())) {
1 ✔
532
                handleRequestBodyResponse(response.getRequestBody());
1 ✔
533
              }
534
            }
535
            // 4. Server Headers
536
            else if (response.hasResponseHeaders()) {
1 ✔
537
              if (response.getResponseHeaders().hasResponse()) {
1 ✔
538
                if (response.getResponseHeaders().getResponse().getStatus()
1 ✔
539
                    == CommonResponse.ResponseStatus.CONTINUE_AND_REPLACE) {
540
                  internalOnError(Status.UNAVAILABLE
1 ✔
541
                      .withDescription("CONTINUE_AND_REPLACE is not supported")
1 ✔
542
                      .asRuntimeException());
1 ✔
543
                  return;
1 ✔
544
                }
545
                Metadata target = wrappedListener.isTrailersOnly()
1 ✔
546
                    ? wrappedListener.getSavedTrailers() : wrappedListener.getSavedHeaders();
1 ✔
547
                applyHeaderMutations(
1 ✔
548
                    target,
549
                    response.getResponseHeaders().getResponse().getHeaderMutation(),
1 ✔
550
                    mutationFilter,
1 ✔
551
                    mutator);
1 ✔
552
              }
553
              if (wrappedListener.isTrailersOnly()) {
1 ✔
554
                wrappedListener.proceedWithClose();
1 ✔
555
              } else {
556
                wrappedListener.proceedWithHeaders();
1 ✔
557
              }
558
            }
559
            // 5. Server Message (Response Body)
560
            else if (response.hasResponseBody()) {
1 ✔
561
              if (validateCompressionSupport(response.getResponseBody())) {
1 ✔
562
                handleResponseBodyResponse(response.getResponseBody(), wrappedListener);
1 ✔
563
              }
564
            }
565
            // 6. Response Trailers
566
            else if (response.hasResponseTrailers()) {
1 ✔
567
              if (response.getResponseTrailers().hasHeaderMutation()) {
1 ✔
568
                applyHeaderMutations(
1 ✔
569
                    wrappedListener.getSavedTrailers(),
1 ✔
570
                    response.getResponseTrailers().getHeaderMutation(),
1 ✔
571
                    mutationFilter,
1 ✔
572
                    mutator);
1 ✔
573
              }
574
              wrappedListener.proceedWithClose();
1 ✔
575
            }
576

577
            checkEndOfStream(response);
1 ✔
578
          } catch (Throwable t) {
1 ✔
579
            internalOnError(t);
1 ✔
580
          }
1 ✔
581
        }
1 ✔
582

583
        @Override
584
        public void onError(Throwable t) {
585
          if (markExtProcStreamFailed(extProcStreamState)) {
1 ✔
586
            synchronized (streamLock) {
1 ✔
587
              extProcClientCallRequestObserver = null;
1 ✔
588
            }
1 ✔
589
            if (config.getObservabilityMode()
1 ✔
590
                || (config.getFailureModeAllow() && !bodyMessageSentToExtProc.get())) {
1 ✔
591
              handleFailOpen(wrappedListener);
1 ✔
592
            } else {
593
              String message = "External processor stream failed";
1 ✔
594
              cancelDownstream(message, t);
1 ✔
595
              wrappedListener.proceedWithClose();
1 ✔
596
            }
597
          }
598
        }
1 ✔
599

600
        @Override
601
        public void onCompleted() {
602
          if (markExtProcStreamCompleted(extProcStreamState)) {
1 ✔
603
            handleFailOpen(wrappedListener);
1 ✔
604
          }
605
        }
1 ✔
606
      });
607

608
      this.collectedAttributes = collectAttributes(
1 ✔
609
          config.getRequestAttributes(), method, channel.authority(), headers);
1 ✔
610

611
      boolean sendRequestHeaders =
1 ✔
612
          currentProcessingMode.getRequestHeaderMode() == ProcessingMode.HeaderSendMode.SEND
1 ✔
613
          || currentProcessingMode.getRequestHeaderMode()
1 ✔
614
              == ProcessingMode.HeaderSendMode.DEFAULT;
615

616
      if (sendRequestHeaders) {
1 ✔
617
        sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
618
            .setRequestHeaders(HttpHeaders.newBuilder()
1 ✔
619
                .setHeaders(toHeaderMap(headers, config.getForwardRulesConfig()))
1 ✔
620
                .setEndOfStream(false)
1 ✔
621
                .build())
1 ✔
622
            .build());
1 ✔
623
      }
624

625
      if (config.getObservabilityMode() || !sendRequestHeaders) {
1 ✔
626
        activateCall();
1 ✔
627
      }
628
    }
1 ✔
629

630
    private void sendToExtProc(ProcessingRequest request) {
631
      synchronized (streamLock) {
1 ✔
632
        if (extProcStreamState.get().isCompleted()) {
1 ✔
633
          return;
×
634
        }
635
        
636
        if (request.hasRequestHeaders()) {
1 ✔
637
          expectedRequestResponse = EventType.REQUEST_HEADERS;
1 ✔
638
        } else if (request.hasResponseHeaders()) {
1 ✔
639
          expectedResponseResponse = EventType.RESPONSE_HEADERS;
1 ✔
640
        } else if (request.hasResponseTrailers()) {
1 ✔
641
          expectedResponseResponse = EventType.RESPONSE_TRAILERS;
1 ✔
642
        }
643

644
        ProcessingRequest requestToSend = request;
1 ✔
645
        if (!protocolConfigSent) {
1 ✔
646
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1 ✔
647
              .setProtocolConfig(ProtocolConfiguration.newBuilder()
1 ✔
648
                  .setRequestBodyMode(currentProcessingMode.getRequestBodyMode())
1 ✔
649
                  .setResponseBodyMode(currentProcessingMode.getResponseBodyMode())
1 ✔
650
                  .build())
1 ✔
651
              .build();
1 ✔
652
          protocolConfigSent = true;
1 ✔
653
        }
654

655
        boolean isClientServerMessage =
1 ✔
656
            requestToSend.hasRequestHeaders() || requestToSend.hasRequestBody();
1 ✔
657
        if (isClientServerMessage
1 ✔
658
            && !requestAttributesSent
659
            && collectedAttributes != null
660
            && !collectedAttributes.isEmpty()) {
1 ✔
661
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1 ✔
662
              .putAllAttributes(collectedAttributes)
1 ✔
663
              .build();
1 ✔
664
          requestAttributesSent = true;
1 ✔
665
        }
666

667
        if (config.getObservabilityMode()) {
1 ✔
668
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1 ✔
669
              .setObservabilityMode(true)
1 ✔
670
              .build();
1 ✔
671
        }
672

673
        extProcClientCallRequestObserver.onNext(requestToSend);
1 ✔
674
      }
1 ✔
675
    }
1 ✔
676

677
    private void onExtProcStreamReady() {
678
      drainPendingRequests();
1 ✔
679
      onReadyNotify();
1 ✔
680
    }
1 ✔
681

682
    private void drainPendingRequests() {
683
      int toRequest = pendingRequests.getAndSet(0);
1 ✔
684
      if (toRequest > 0) {
1 ✔
685
        super.request(toRequest);
1 ✔
686
      }
687
    }
1 ✔
688

689
    private void closeExtProcStream() {
690
      synchronized (streamLock) {
1 ✔
691
        if (markExtProcStreamCompleted(extProcStreamState)) {
1 ✔
692
          if (extProcClientCallRequestObserver != null) {
1 ✔
693
            extProcClientCallRequestObserver.onCompleted();
1 ✔
694
          }
695
        }
696
      }
1 ✔
697
    }
1 ✔
698

699
    private void internalOnError(Throwable t) {
700
      if (markExtProcStreamFailed(extProcStreamState)) {
1 ✔
701
        synchronized (streamLock) {
1 ✔
702
          if (extProcClientCallRequestObserver != null) {
1 ✔
703
            try {
704
              extProcClientCallRequestObserver.onError(t);
1 ✔
705
            } catch (Throwable ignored) {
×
706
              // Ignore exceptions during cancel/onError propagation
707
            }
1 ✔
708
            extProcClientCallRequestObserver = null;
1 ✔
709
          }
710
        }
1 ✔
711
        if (config.getObservabilityMode()
1 ✔
712
            || (config.getFailureModeAllow() && !bodyMessageSentToExtProc.get())) {
1 ✔
713
          handleFailOpen(wrappedListener);
1 ✔
714
        } else {
715
          String message = "External processor stream failed";
1 ✔
716
          cancelDownstream(message, t);
1 ✔
717
          wrappedListener.proceedWithClose();
1 ✔
718
        }
719
      }
720
    }
1 ✔
721

722
    private void halfCloseExtProcStream() {
723
      synchronized (streamLock) {
1 ✔
724
        if (!extProcStreamState.get().isCompleted() && extProcClientCallRequestObserver != null) {
1 ✔
725
          extProcClientCallRequestObserver.onCompleted();
1 ✔
726
        }
727
      }
1 ✔
728
    }
1 ✔
729

730
    private void onReadyNotify() {
731
      wrappedListener.onReadyNotify();
1 ✔
732
    }
1 ✔
733

734
    private boolean isSidecarReady() {
735
      ExtProcStreamState state = extProcStreamState.get();
1 ✔
736
      if (state.isCompleted()) {
1 ✔
737
        return true;
×
738
      }
739
      if (state.isDraining()) {
1 ✔
740
        return false;
1 ✔
741
      }
742
      synchronized (streamLock) {
1 ✔
743
        ClientCallStreamObserver<ProcessingRequest> observer = extProcClientCallRequestObserver;
1 ✔
744
        return observer != null && observer.isReady();
1 ✔
745
      }
746
    }
747

748
    @Override
749
    public boolean isReady() {
750
      if (passThroughMode.get()) {
1 ✔
751
        return super.isReady();
1 ✔
752
      }
753
      if (extProcStreamState.get().isCompleted()) {
1 ✔
754
        return super.isReady();
×
755
      }
756
      if (dataPlaneCallState.get() == DataPlaneCallState.IDLE && !config.getObservabilityMode()) {
1 ✔
757
        return false;
1 ✔
758
      }
759
      boolean sidecarReady = isSidecarReady();
1 ✔
760
      if (config.getObservabilityMode()) {
1 ✔
761
        return super.isReady() && sidecarReady;
1 ✔
762
      }
763
      return sidecarReady;
1 ✔
764
    }
765

766
    @Override
767
    public void request(int numMessages) {
768
      if (passThroughMode.get() || extProcStreamState.get().isCompleted()) {
1 ✔
769
        super.request(numMessages);
1 ✔
770
        return;
1 ✔
771
      }
772
      if (!config.getObservabilityMode() 
1 ✔
773
          && currentProcessingMode.getResponseBodyMode() != ProcessingMode.BodySendMode.GRPC) {
1 ✔
774
        super.request(numMessages);
1 ✔
775
        return;
1 ✔
776
      }
777
      if (!isSidecarReady()) {
1 ✔
778
        pendingRequests.addAndGet(numMessages);
1 ✔
779
        return;
1 ✔
780
      }
781
      super.request(numMessages);
1 ✔
782
    }
1 ✔
783

784
    @Override
785
    public void sendMessage(InputStream message) {
786
      if (requestSideClosed.get()) {
1 ✔
787
        // External processor already closed the request stream. Discard further messages.
788
        return;
1 ✔
789
      }
790

791
      if (passThroughMode.get()) {
1 ✔
792
        super.sendMessage(message);
1 ✔
793
        return;
1 ✔
794
      }
795

796
      synchronized (streamLock) {
1 ✔
797
        if (passThroughMode.get()) {
1 ✔
798
          super.sendMessage(message);
×
799
          return;
×
800
        }
801

802
        ExtProcStreamState state = extProcStreamState.get();
1 ✔
803
        if (state.isDraining() || state.isCompleted()) {
1 ✔
804
          if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1 ✔
805
            super.sendMessage(message);
1 ✔
806
            return;
1 ✔
807
          }
808
          try {
809
            ByteString copiedBody = ByteString.readFrom(message);
1 ✔
810
            pendingDrainingMessages.add(new KnownLengthInputStream(copiedBody));
1 ✔
811
          } catch (IOException e) {
×
812
            rawCall.cancel("Failed to copy outbound message for buffering", e);
×
813
          }
1 ✔
814
          return;
1 ✔
815
        }
816
      }
1 ✔
817

818
      if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1 ✔
819
        super.sendMessage(message);
1 ✔
820
        return;
1 ✔
821
      }
822

823
      // Mode is GRPC
824
      try {
825
        ByteString bodyByteString = outboundStreamToByteString(message);
1 ✔
826
        sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
827
            .setRequestBody(HttpBody.newBuilder()
1 ✔
828
                .setBody(bodyByteString)
1 ✔
829
                .setEndOfStream(false)
1 ✔
830
                .build())
1 ✔
831
            .build());
1 ✔
832
        bodyMessageSentToExtProc.set(true);
1 ✔
833

834
        if (config.getObservabilityMode()) {
1 ✔
835
          super.sendMessage(new KnownLengthInputStream(bodyByteString));
×
836
        }
837
      } catch (IOException e) {
×
838
        rawCall.cancel("Failed to serialize message for External Processor", e);
×
839
      }
1 ✔
840
    }
1 ✔
841

842
    private void proceedWithHalfClose() {
843
      if (clientHalfCloseStartNanos > 0) {
1 ✔
844
        long durationNanos = System.nanoTime() - clientHalfCloseStartNanos;
1 ✔
845
        recordDuration(clientHalfCloseDuration, durationNanos);
1 ✔
846
        clientHalfCloseStartNanos = 0;
1 ✔
847
      }
848
      super.halfClose();
1 ✔
849
    }
1 ✔
850

851
    @Override
852
    public void halfClose() {
853
      clientHalfCloseStartNanos = System.nanoTime();
1 ✔
854
      if (passThroughMode.get()) {
1 ✔
855
        if (requestSideClosed.compareAndSet(false, true)) {
1 ✔
856
          proceedWithHalfClose();
1 ✔
857
        }
858
        return;
1 ✔
859
      }
860

861
      pendingHalfClose.set(true);
1 ✔
862

863
      if (extProcStreamState.get().isCompleted()) {
1 ✔
864
        if (passThroughMode.get()) {
1 ✔
865
          if (requestSideClosed.compareAndSet(false, true)) {
×
866
            proceedWithHalfClose();
×
867
          }
868
        }
869
        return;
1 ✔
870
      }
871

872
      if (extProcStreamState.get().isDraining()) {
1 ✔
873
        boolean canProceed = false;
1 ✔
874
        synchronized (streamLock) {
1 ✔
875
          if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE
1 ✔
876
              || (!bodyMessageSentToExtProc.get() && pendingDrainingMessages.isEmpty())) {
1 ✔
877
            canProceed = true;
1 ✔
878
          }
879
        }
1 ✔
880
        if (canProceed) {
1 ✔
881
          if (requestSideClosed.compareAndSet(false, true)) {
1 ✔
882
            proceedWithHalfClose();
1 ✔
883
          }
884
        }
885
        return;
1 ✔
886
      }
887

888
      if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1 ✔
889
        if (requestSideClosed.compareAndSet(false, true)) {
1 ✔
890
          proceedWithHalfClose();
1 ✔
891
        }
892
        return;
1 ✔
893
      }
894

895
      // Mode is GRPC
896
      sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
897
          .setRequestBody(HttpBody.newBuilder()
1 ✔
898
              .setEndOfStreamWithoutMessage(true)
1 ✔
899
              .build())
1 ✔
900
          .build());
1 ✔
901
    }
1 ✔
902

903
    private void cancelDownstream(@Nullable String message, @Nullable Throwable cause) {
904
      if (downstreamCancelled.compareAndSet(false, true)) {
1 ✔
905
        delayedCall.cancel(message, cause);
1 ✔
906
      }
907
    }
1 ✔
908

909
    @Override
910
    public void cancel(@Nullable String message, @Nullable Throwable cause) {
911
      synchronized (streamLock) {
1 ✔
912
        if (!extProcStreamState.get().isCompleted() && extProcClientCallRequestObserver != null) {
1 ✔
913
          extProcClientCallRequestObserver.onError(
1 ✔
914
              Status.CANCELLED
915
                  .withDescription(message)
1 ✔
916
                  .withCause(cause)
1 ✔
917
                  .asRuntimeException());
1 ✔
918
        }
919
      }
1 ✔
920
      cancelDownstream(message, cause);
1 ✔
921
    }
1 ✔
922

923
    private void handleRequestBodyResponse(BodyResponse bodyResponse) {
924
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1 ✔
925
        BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
1 ✔
926
        if (mutation.hasStreamedResponse()) {
1 ✔
927
          StreamedBodyResponse streamed = mutation.getStreamedResponse();
1 ✔
928
          if (!streamed.getEndOfStreamWithoutMessage()) {
1 ✔
929
            super.sendMessage(new KnownLengthInputStream(streamed.getBody()));
1 ✔
930
          }
931
          if (streamed.getEndOfStream() || streamed.getEndOfStreamWithoutMessage()) {
1 ✔
932
            if (requestSideClosed.compareAndSet(false, true)) {
1 ✔
933
              proceedWithHalfClose();
1 ✔
934
            }
935
          }
936
        }
937
      }
938
    }
1 ✔
939

940
    private void handleResponseBodyResponse(
941
        BodyResponse bodyResponse, DataPlaneListener listener) {
942
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1 ✔
943
        BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
1 ✔
944
        if (mutation.hasStreamedResponse()) {
1 ✔
945
          StreamedBodyResponse streamed = mutation.getStreamedResponse();
1 ✔
946
          listener.onExternalBody(streamed.getBody());
1 ✔
947
        }
948
      }
949
    }
1 ✔
950

951
    private void handleImmediateResponse(ImmediateResponse immediate, DataPlaneListener listener)
952
        throws HeaderMutationDisallowedException {
953
      Status status = Status.fromCodeValue(immediate.getGrpcStatus().getStatus());
1 ✔
954
      if (!immediate.getDetails().isEmpty()) {
1 ✔
955
        status = status.withDescription(immediate.getDetails());
1 ✔
956
      }
957

958
      Metadata trailers = new Metadata();
1 ✔
959
      if (immediate.hasHeaders()) {
1 ✔
960
        applyHeaderMutations(trailers, immediate.getHeaders(), mutationFilter, mutator);
1 ✔
961
      }
962

963
      listener.setImmediateResponse(status, trailers);
1 ✔
964

965
      if (isProcessingTrailers.get()) {
1 ✔
966
        // If sent in response to a server trailers event, sets the status and optionally
967
        // headers to be included in the trailers.
968
        listener.unblockAfterStreamComplete();
1 ✔
969
      } else {
970
        // If sent in response to any other event, it will cause the data plane RPC to
971
        // immediately fail with the specified status as if it were an out-of-band
972
        // cancellation.
973
        rawCall.cancel(status.getDescription(), null);
1 ✔
974
        listener.unblockAfterStreamComplete();
1 ✔
975
      }
976
      closeExtProcStream();
1 ✔
977
    }
1 ✔
978

979
    private void drainPendingDrainingMessages() {
980
      synchronized (streamLock) {
1 ✔
981
        InputStream msg;
982
        while ((msg = pendingDrainingMessages.poll()) != null) {
1 ✔
983
          super.sendMessage(msg);
1 ✔
984
        }
985
        passThroughMode.set(true);
1 ✔
986
        if (pendingHalfClose.get()) {
1 ✔
987
          if (requestSideClosed.compareAndSet(false, true)) {
1 ✔
988
            proceedWithHalfClose();
1 ✔
989
          }
990
        }
991
      }
1 ✔
992
    }
1 ✔
993

994
    private void handleFailOpen(DataPlaneListener listener) {
995
      activateCall();
1 ✔
996
      drainPendingRequests();
1 ✔
997
      listener.unblockAfterStreamComplete();
1 ✔
998
      closeExtProcStream();
1 ✔
999
    }
1 ✔
1000

1001
    private void checkEndOfStream(ProcessingResponse response) {
1002
      boolean terminal = false;
1 ✔
1003
      if (response.hasResponseTrailers()) {
1 ✔
1004
        terminal = true;
1 ✔
1005
      } else if (response.hasResponseHeaders() && wrappedListener.isTrailersOnly()) {
1 ✔
1006
        terminal = true;
1 ✔
1007
      }
1008

1009
      if (terminal) {
1 ✔
1010
        wrappedListener.unblockAfterStreamComplete();
1 ✔
1011
        closeExtProcStream();
1 ✔
1012
      }
1013
    }
1 ✔
1014

1015
    long getServerHeadersStartNanos() {
1016
      return serverHeadersStartNanos;
1 ✔
1017
    }
1018

1019
    void setServerHeadersStartNanos(long serverHeadersStartNanos) {
1020
      this.serverHeadersStartNanos = serverHeadersStartNanos;
1 ✔
1021
    }
1 ✔
1022

1023
    long getServerTrailersStartNanos() {
1024
      return serverTrailersStartNanos;
1 ✔
1025
    }
1026

1027
    void setServerTrailersStartNanos(long serverTrailersStartNanos) {
1028
      this.serverTrailersStartNanos = serverTrailersStartNanos;
1 ✔
1029
    }
1 ✔
1030

1031
    AtomicReference<ExtProcStreamState> getExtProcStreamState() {
1032
      return extProcStreamState;
1 ✔
1033
    }
1034

1035
    ProcessingMode getCurrentProcessingMode() {
1036
      return currentProcessingMode;
1 ✔
1037
    }
1038

1039
    AtomicBoolean getPassThroughMode() {
1040
      return passThroughMode;
1 ✔
1041
    }
1042

1043
    ExternalProcessorFilterConfig getConfig() {
1044
      return config;
1 ✔
1045
    }
1046

1047
    Context getCallContext() {
1048
      return callContext;
1 ✔
1049
    }
1050

1051
    ScheduledExecutorService getScheduler() {
1052
      return scheduler;
1 ✔
1053
    }
1054

1055
    AtomicBoolean getIsProcessingTrailers() {
1056
      return isProcessingTrailers;
1 ✔
1057
    }
1058
  }
1059

1060
  private static class DataPlaneListener extends SimpleForwardingClientCallListener<InputStream> {
1061
    private final ClientCall<?, ?> rawCall;
1062
    private final DataPlaneClientCall dataPlaneClientCall;
1063
    private final Queue<InputStream> savedMessages = new ConcurrentLinkedQueue<>();
1 ✔
1064
    private boolean inboundPassThrough = false;
1 ✔
1065
    @Nullable private volatile Metadata savedHeaders;
1066
    @Nullable private volatile Metadata savedTrailers;
1067
    @Nullable private volatile Status savedStatus;
1068
    private final AtomicBoolean terminationTriggered = new AtomicBoolean(false);
1 ✔
1069
    private final AtomicBoolean responseHeadersSent = new AtomicBoolean(false);
1 ✔
1070
    private final AtomicBoolean trailersOnly = new AtomicBoolean(false);
1 ✔
1071

1072
    protected DataPlaneListener(
1073
        ClientCall.Listener<InputStream> delegate,
1074
        ClientCall<?, ?> rawCall,
1075
        DataPlaneClientCall dataPlaneClientCall) {
1076
      super(delegate);
1 ✔
1077
      this.rawCall = rawCall;
1 ✔
1078
      this.dataPlaneClientCall = dataPlaneClientCall;
1 ✔
1079
    }
1 ✔
1080

1081
    boolean isTrailersOnly() {
1082
      return trailersOnly.get();
1 ✔
1083
    }
1084

1085
    Metadata getSavedHeaders() {
1086
      return savedHeaders;
1 ✔
1087
    }
1088

1089
    Metadata getSavedTrailers() {
1090
      return savedTrailers;
1 ✔
1091
    }
1092

1093
    void setImmediateResponse(Status status, Metadata trailers) {
1094
      this.savedStatus = status;
1 ✔
1095
      this.savedTrailers = trailers;
1 ✔
1096
    }
1 ✔
1097

1098
    @Override
1099
    public void onReady() {
1100
      dataPlaneClientCall.drainPendingRequests();
1 ✔
1101
      onReadyNotify();
1 ✔
1102
    }
1 ✔
1103

1104
    @Override
1105
    public void onHeaders(Metadata headers) {
1106
      dataPlaneClientCall.setServerHeadersStartNanos(System.nanoTime());
1 ✔
1107
      responseHeadersSent.set(true);
1 ✔
1108
      boolean sendResponseHeaders =
1 ✔
1109
          dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1 ✔
1110
              == ProcessingMode.HeaderSendMode.SEND
1111
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1 ✔
1112
              == ProcessingMode.HeaderSendMode.DEFAULT;
1113

1114
      if (dataPlaneClientCall.getExtProcStreamState().get().isDraining() && sendResponseHeaders) {
1 ✔
1115
        this.savedHeaders = headers;
1 ✔
1116
        return;
1 ✔
1117
      }
1118

1119
      if (dataPlaneClientCall.getPassThroughMode().get() 
1 ✔
1120
          || dataPlaneClientCall.getExtProcStreamState().get().isCompleted() 
1 ✔
1121
          || !sendResponseHeaders) {
1122
        proceedWithHeaders(headers);
1 ✔
1123
        return;
1 ✔
1124
      }
1125

1126
      this.savedHeaders = headers;
1 ✔
1127
      dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
1128
          .setResponseHeaders(HttpHeaders.newBuilder()
1 ✔
1129
              .setHeaders(
1 ✔
1130
                  toHeaderMap(headers, dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1 ✔
1131
              .build())
1 ✔
1132
          .build());
1 ✔
1133

1134
      if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1 ✔
1135
        proceedWithHeaders();
×
1136
      }
1137
    }
1 ✔
1138

1139
    @Override
1140
    public void onMessage(InputStream message) {
1141
      synchronized (savedMessages) {
1 ✔
1142
        if (inboundPassThrough) {
1 ✔
1143
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
1 ✔
1144
          return;
1 ✔
1145
        }
1146

1147
        boolean checkDrain = dataPlaneClientCall.getExtProcStreamState().get().isDraining()
1 ✔
1148
            && dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1 ✔
1149
                == ProcessingMode.BodySendMode.GRPC;
1150

1151
        if (savedHeaders != null || checkDrain) {
1 ✔
1152
          try {
1153
            ByteString copiedBody = ByteString.readFrom(message);
1 ✔
1154
            savedMessages.add(new KnownLengthInputStream(copiedBody));
1 ✔
1155
          } catch (IOException e) {
×
1156
            rawCall.cancel("Failed to copy inbound message for buffering", e);
×
1157
          }
1 ✔
1158
          return;
1 ✔
1159
        }
1160
      }
1 ✔
1161

1162
      if (dataPlaneClientCall.getPassThroughMode().get()) {
1 ✔
1163
        dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
×
1164
        return;
×
1165
      }
1166

1167
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1 ✔
1168
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1 ✔
1169
              != ProcessingMode.BodySendMode.GRPC) {
1170
        dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
1 ✔
1171
        return;
1 ✔
1172
      }
1173

1174
      try {
1175
        ByteString bodyByteString = ByteString.readFrom(message);
1 ✔
1176
        sendResponseBodyToExtProc(bodyByteString, false);
1 ✔
1177
        dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
1 ✔
1178

1179
        if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1 ✔
1180
          // If needed, downstream reading can be made more optimal by creating a wrapped
1181
          // Inputstream wraps the underlying bytestring and that implements HasByteBuffer,
1182
          // Detachable, KnownLength
1183
          dataPlaneClientCall.getCallContext().run(
×
1184
              () -> delegate().onMessage(bodyByteString.newInput()));
×
1185
        }
1186
      } catch (IOException e) {
×
1187
        rawCall.cancel("Failed to read server response", e);
×
1188
      }
1 ✔
1189
    }
1 ✔
1190

1191
    @Override
1192
    public void onClose(Status status, Metadata trailers) {
1193
      dataPlaneClientCall.setServerTrailersStartNanos(System.nanoTime());
1 ✔
1194
      ExtProcStreamState extProcStreamState =
1 ✔
1195
          dataPlaneClientCall.getExtProcStreamState().get();
1 ✔
1196
      if (extProcStreamState.isFailed()
1 ✔
1197
          && !dataPlaneClientCall.getConfig().getObservabilityMode()
1 ✔
1198
          && (!dataPlaneClientCall.getConfig().getFailureModeAllow()
1 ✔
1199
              || dataPlaneClientCall.bodyMessageSentToExtProc.get())) {
1 ✔
1200
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1 ✔
1201
          proceedWithClose(Status.INTERNAL.withDescription("External processor stream failed")
1 ✔
1202
              .withCause(status.getCause()), new Metadata());
1 ✔
1203
        }
1204
        return;
1 ✔
1205
      }
1206
      if (dataPlaneClientCall.getPassThroughMode().get()) {
1 ✔
1207
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1 ✔
1208
          proceedWithClose(status, trailers);
1 ✔
1209
        }
1210
        return;
1 ✔
1211
      }
1212

1213
      if (this.savedStatus == null) {
1 ✔
1214
        this.savedStatus = status;
1 ✔
1215
        this.savedTrailers = trailers;
1 ✔
1216
      }
1217

1218
      // If we are still waiting for the external processor to validate response headers,
1219
      // buffer the close status/trailers and defer the close until headers are processed.
1220
      if (savedHeaders != null) {
1 ✔
1221
        return;
1 ✔
1222
      }
1223

1224
      boolean sendResponseTrailers =
1 ✔
1225
          dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1 ✔
1226
              == ProcessingMode.HeaderSendMode.SEND;
1227

1228
      if (dataPlaneClientCall.getExtProcStreamState().get().isDraining() && sendResponseTrailers) {
1 ✔
1229
        return;
1 ✔
1230
      }
1231

1232
      if (!responseHeadersSent.get()) {
1 ✔
1233
        trailersOnly.set(true);
1 ✔
1234
      }
1235

1236
      triggerCloseHandshake();
1 ✔
1237
    }
1 ✔
1238

1239
    void onReadyNotify() {
1240
      dataPlaneClientCall.getCallContext().run(() -> delegate().onReady());
1 ✔
1241
    }
1 ✔
1242

1243
    void proceedWithHeaders() {
1244
      if (savedHeaders != null) {
1 ✔
1245
        proceedWithHeaders(savedHeaders);
1 ✔
1246
        synchronized (savedMessages) {
1 ✔
1247
          savedHeaders = null;
1 ✔
1248
          if (!dataPlaneClientCall.getExtProcStreamState().get().isDraining()) {
1 ✔
1249
            InputStream msg;
1250
            while ((msg = savedMessages.poll()) != null) {
1 ✔
1251
              onMessage(msg);
1 ✔
1252
            }
1253
          }
1254
        }
1 ✔
1255
        onReadyNotify();
1 ✔
1256
        if (savedStatus != null) {
1 ✔
1257
          triggerCloseHandshake();
1 ✔
1258
        }
1259
      }
1260
    }
1 ✔
1261

1262
    private void proceedWithHeaders(Metadata headers) {
1263
      if (dataPlaneClientCall.getServerHeadersStartNanos() > 0) {
1 ✔
1264
        long durationNanos = System.nanoTime() - dataPlaneClientCall.getServerHeadersStartNanos();
1 ✔
1265
        dataPlaneClientCall.recordDuration(serverHeadersDuration, durationNanos);
1 ✔
1266
        dataPlaneClientCall.setServerHeadersStartNanos(0);
1 ✔
1267
      }
1268
      dataPlaneClientCall.getCallContext().run(() -> delegate().onHeaders(headers));
1 ✔
1269
    }
1 ✔
1270

1271
    void proceedWithClose() {
1272
      if (savedStatus != null) {
1 ✔
1273
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1 ✔
1274
          proceedWithClose(savedStatus, savedTrailers);
1 ✔
1275
        }
1276
        savedStatus = null;
1 ✔
1277
        savedTrailers = null;
1 ✔
1278
      }
1279
    }
1 ✔
1280

1281
    private void proceedWithClose(Status status, Metadata trailers) {
1282
      if (dataPlaneClientCall.getServerTrailersStartNanos() > 0) {
1 ✔
1283
        long durationNanos = System.nanoTime() - dataPlaneClientCall.getServerTrailersStartNanos();
1 ✔
1284
        dataPlaneClientCall.recordDuration(serverTrailersDuration, durationNanos);
1 ✔
1285
        dataPlaneClientCall.setServerTrailersStartNanos(0);
1 ✔
1286
      }
1287
      dataPlaneClientCall.getCallContext().run(() -> delegate().onClose(status, trailers));
1 ✔
1288
    }
1 ✔
1289

1290
    void onExternalBody(ByteString body) {
1291
      // If needed, downstream reading can be made more optimal by creating a wrapped
1292
      // Inputstream wraps the underlying bytestring and that implements HasByteBuffer,
1293
      // Detachable, KnownLength
1294
      dataPlaneClientCall.getCallContext().run(
1 ✔
1295
          () -> delegate().onMessage(body.newInput()));
1 ✔
1296
    }
1 ✔
1297

1298
    void unblockAfterStreamComplete() {
1299
      proceedWithHeaders();
1 ✔
1300
      proceedWithSavedMessages();
1 ✔
1301
      dataPlaneClientCall.drainPendingDrainingMessages();
1 ✔
1302
      proceedWithClose();
1 ✔
1303
    }
1 ✔
1304

1305
    private void proceedWithSavedMessages() {
1306
      synchronized (savedMessages) {
1 ✔
1307
        InputStream msg;
1308
        while ((msg = savedMessages.poll()) != null) {
1 ✔
1309
          final InputStream finalMsg = msg;
1 ✔
1310
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(finalMsg));
1 ✔
1311
        }
1 ✔
1312
        inboundPassThrough = true;
1 ✔
1313
      }
1 ✔
1314
    }
1 ✔
1315

1316
    private void triggerCloseHandshake() {
1317
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1 ✔
1318
          || !terminationTriggered.compareAndSet(false, true)) {
1 ✔
1319
        return;
1 ✔
1320
      }
1321

1322
      boolean sendResponseHeaders =
1 ✔
1323
          dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1 ✔
1324
              == ProcessingMode.HeaderSendMode.SEND
1325
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1 ✔
1326
              == ProcessingMode.HeaderSendMode.DEFAULT;
1327

1328
      boolean sendResponseTrailers =
1 ✔
1329
          dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1 ✔
1330
              == ProcessingMode.HeaderSendMode.SEND;
1331

1332
      if (trailersOnly.get()) {
1 ✔
1333
        if (sendResponseHeaders) {
1 ✔
1334
          dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
1335
              .setResponseHeaders(HttpHeaders.newBuilder()
1 ✔
1336
                  .setHeaders(
1 ✔
1337
                      toHeaderMap(
1 ✔
1338
                          savedTrailers,
1339
                          dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1 ✔
1340
                  .setEndOfStream(true)
1 ✔
1341
                  .build())
1 ✔
1342
              .build());
1 ✔
1343
        } else {
1344
          proceedWithClose();
1 ✔
1345
          if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1 ✔
1346
            dataPlaneClientCall.closeExtProcStream();
1 ✔
1347
          }
1348
        }
1349
      } else if (sendResponseTrailers) {
1 ✔
1350
        dataPlaneClientCall.getIsProcessingTrailers().set(true);
1 ✔
1351
        dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
1352
            .setResponseTrailers(HttpTrailers.newBuilder()
1 ✔
1353
                .setTrailers(
1 ✔
1354
                    toHeaderMap(
1 ✔
1355
                        savedTrailers,
1356
                        dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1 ✔
1357
                .build())
1 ✔
1358
            .build());
1 ✔
1359
      } else {
1360
        proceedWithClose();
1 ✔
1361
        if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1 ✔
1362
          dataPlaneClientCall.closeExtProcStream();
1 ✔
1363
        }
1364
      }
1365

1366
      if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1 ✔
1367
        proceedWithClose();
1 ✔
1368
        @SuppressWarnings("unused")
1369
        ScheduledFuture<?> unused = dataPlaneClientCall.getScheduler().schedule(
1 ✔
1370
            dataPlaneClientCall::closeExtProcStream,
1 ✔
1371
            dataPlaneClientCall.getConfig().getDeferredCloseTimeoutNanos(),
1 ✔
1372
            TimeUnit.NANOSECONDS);
1373
      }
1374
    }
1 ✔
1375

1376
    private void sendResponseBodyToExtProc(
1377
        @Nullable ByteString bodyByteString, boolean endOfStream) {
1378
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1 ✔
1379
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1 ✔
1380
              != ProcessingMode.BodySendMode.GRPC) {
1381
        return;
×
1382
      }
1383

1384
      HttpBody.Builder bodyBuilder =
1385
          HttpBody.newBuilder();
1 ✔
1386
      if (bodyByteString != null) {
1 ✔
1387
        bodyBuilder.setBody(bodyByteString);
1 ✔
1388
      }
1389
      bodyBuilder.setEndOfStream(endOfStream);
1 ✔
1390

1391
      dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1 ✔
1392
          .setResponseBody(bodyBuilder.build())
1 ✔
1393
          .build());
1 ✔
1394
    }
1 ✔
1395
  }
1396
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc