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

grpc / grpc-java / #20485

17 Sep 2026 05:22AM UTC coverage: 89.313% (+0.009%) from 89.304%
#20485

push

github

jdcormie
binder: deliver, even in the case where a prefix-only txn arrives last

Inbound's queuedTransactionData holds message fragments from the peer
that haven't yet been assembled and delivered to the application.
Binder transactions are sent in index order and must contain at least
one of: a prefix, part/all of a message, and a suffix. When a
transaction with message data arrives before its predecessors (according
to index), Inbound's enqueueTransactionData() reserves slots for those
predecessors in queuedTransactionData. After each predecessor
transaction trickles in, enqueueTransactionData() checks whether it
completes the message, which can then be delivered. There's an edge
case, though, where a late-arrival transaction contains just the prefix
with no message data at all. In that case, Inbound remove()s the
queuedTransactionData slot it previously reserved and carries on without
considering that this prefix-only transaction might have completed the
first message! And if that was the final transaction in the stream,
Inbound will never call lookForCompleteMessage() again, leaving the
stream stuck forever with a complete but undelivered message.

39146 of 43830 relevant lines covered (89.31%)

0.89 hits per line

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

95.69
/../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.GrpcUtil;
64
import io.grpc.internal.SerializingExecutor;
65
import io.grpc.stub.ClientCallStreamObserver;
66
import io.grpc.stub.ClientResponseObserver;
67
import io.grpc.stub.MetadataUtils;
68
import io.grpc.xds.ExternalProcessorFilter.ExternalProcessorFilterConfig;
69
import io.grpc.xds.Filter.FilterContext;
70
import io.grpc.xds.internal.extproc.DataPlaneCallState;
71
import io.grpc.xds.internal.extproc.EventType;
72
import io.grpc.xds.internal.extproc.ExtProcStreamState;
73
import io.grpc.xds.internal.extproc.KnownLengthInputStream;
74
import io.grpc.xds.internal.grpcservice.CachedChannelManager;
75
import io.grpc.xds.internal.grpcservice.HeaderValue;
76
import io.grpc.xds.internal.headermutations.HeaderMutationDisallowedException;
77
import io.grpc.xds.internal.headermutations.HeaderMutationFilter;
78
import io.grpc.xds.internal.headermutations.HeaderMutationRulesConfig;
79
import io.grpc.xds.internal.headermutations.HeaderMutator;
80
import java.io.IOException;
81
import java.io.InputStream;
82
import java.util.ArrayList;
83
import java.util.List;
84
import java.util.Optional;
85
import java.util.Queue;
86
import java.util.concurrent.ConcurrentLinkedQueue;
87
import java.util.concurrent.Executor;
88
import java.util.concurrent.ScheduledExecutorService;
89
import java.util.concurrent.ScheduledFuture;
90
import java.util.concurrent.TimeUnit;
91
import java.util.concurrent.atomic.AtomicBoolean;
92
import java.util.concurrent.atomic.AtomicInteger;
93
import java.util.concurrent.atomic.AtomicReference;
94
import javax.annotation.Nullable;
95
import javax.annotation.concurrent.GuardedBy;
96

97
/**
98
 * Client-side interceptor for external processing filter.
99
 */
100
final class ExternalProcessorClientInterceptor implements ClientInterceptor {
101

102
  @VisibleForTesting
103
  static DoubleHistogramMetricInstrument clientHeadersDuration;
104
  @VisibleForTesting
105
  static DoubleHistogramMetricInstrument clientHalfCloseDuration;
106
  @VisibleForTesting
107
  static DoubleHistogramMetricInstrument serverHeadersDuration;
108
  @VisibleForTesting
109
  static DoubleHistogramMetricInstrument serverTrailersDuration;
110

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

119
  static {
120
    initMetricInstruments();
1✔
121
  }
1✔
122

123
  static synchronized void initMetricInstruments() {
124
    if (GrpcUtil.getFlag("GRPC_EXPERIMENTAL_XDS_EXT_PROC_ON_CLIENT", false)) {
1✔
125
      if (clientHeadersDuration == null) {
1✔
126
        MetricInstrumentRegistry registry = MetricInstrumentRegistry.getDefaultRegistry();
1✔
127

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

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

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

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

171
  private final ExternalProcessorFilterConfig filterConfig;
172
  private final ScheduledExecutorService scheduler;
173
  private final MetricRecorder metricsRecorder;
174
  private final ManagedChannel extProcChannel;
175

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

188
  @VisibleForTesting
189
  ExternalProcessorFilterConfig getFilterConfig() {
190
    return filterConfig;
1✔
191
  }
192

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

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

244
          @Override
245
          public void cancel(@Nullable String message, @Nullable Throwable cause) {
246
            if (cancelled.compareAndSet(false, true)) {
1✔
247
              super.cancel(message, cause);
1✔
248
            }
249
          }
1✔
250
        };
251

252
    // Create a local subclass instance to buffer outbound actions
253
    DataPlaneDelayedCall<InputStream, InputStream> delayedCall =
1✔
254
        new DataPlaneDelayedCall<>(
255
            serializingExecutor, scheduler, callOptions.getDeadline());
1✔
256

257
    DataPlaneClientCall dataPlaneCall = new DataPlaneClientCall(
1✔
258
        delayedCall, rawCall, extProcStub, filterConfig, filterConfig.getMutationRulesConfig(),
1✔
259
        scheduler, rawMethod, next, metricsRecorder, next.authority());
1✔
260

261
    return (ClientCall<ReqT, RespT>) (ClientCall<?, ?>) dataPlaneCall;
1✔
262
  }
263

264
  // --- SHARED UTILITY METHODS ---
265

266
  /**
267
   * A local subclass to expose the protected constructor of DelayedClientCall.
268
   */
269
  private static class DataPlaneDelayedCall<ReqT, RespT> extends DelayedClientCall<ReqT, RespT> {
270
    DataPlaneDelayedCall(
271
        Executor executor, ScheduledExecutorService scheduler, @Nullable Deadline deadline) {
272
      super("ext_proc", executor, scheduler, deadline);
1✔
273
    }
1✔
274
  }
275

276
  /**
277
   * Handles the bidirectional stream with the External Processor.
278
   * Buffers the actual RPC start until the Ext Proc header response is received.
279
   */
280
  private static class DataPlaneClientCall 
281
      extends SimpleForwardingClientCall<InputStream, InputStream> {
282

283
    private final ExternalProcessorGrpc.ExternalProcessorStub stub;
284
    private final ExternalProcessorFilterConfig config;
285
    private final ClientCall<InputStream, InputStream> rawCall;
286
    private final DataPlaneDelayedCall<InputStream, InputStream> delayedCall;
287
    private final ScheduledExecutorService scheduler;
288
    final Object streamLock = new Object();
1✔
289
    @Nullable private volatile EventType expectedRequestResponse;
290
    @Nullable private volatile EventType expectedResponseResponse;
291
    @Nullable private volatile ClientCallStreamObserver<ProcessingRequest>
292
        extProcClientCallRequestObserver;
293
    @GuardedBy("streamLock")
1✔
294
    private final Queue<InputStream> pendingDrainingMessages =
295
        new ConcurrentLinkedQueue<>();
296
    @Nullable private volatile DataPlaneListener wrappedListener;
297
    private final HeaderMutationFilter mutationFilter;
298
    private final HeaderMutator mutator = HeaderMutator.create();
1✔
299
    private final AtomicInteger pendingRequests = new AtomicInteger(0);
1✔
300
    private final ProcessingMode currentProcessingMode;
301

302
    // Default initial window size
303
    private static final long DEFAULT_INITIAL_WINDOW_SIZE = 65536;
304

305
    // Outbound (sending) windows
306
    @GuardedBy("streamLock")
1✔
307
    private long downstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
308
    @GuardedBy("streamLock")
1✔
309
    private long upstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
310

311
    // Inbound (receiving) windows
312
    @GuardedBy("streamLock")
1✔
313
    private long sidestreamToUpstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
314
    @GuardedBy("streamLock")
1✔
315
    private long sidestreamToDownstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
316

317
    // Threshold to trigger standalone client window updates
318
    private static final long WINDOW_UPDATE_THRESHOLD = DEFAULT_INITIAL_WINDOW_SIZE / 2;
319

320
    // Path 1: Pending/buffered request body messages from downstream
321
    @GuardedBy("streamLock")
1✔
322
    private final Queue<ByteString> pendingRequestBodyMessages = new ConcurrentLinkedQueue<>();
323
    // Deferred half-close flag for upstream direction
324
    private final AtomicBoolean pendingUpstreamHalfClose = new AtomicBoolean(false);
1✔
325

326
    // Path 2: Buffered request body messages from ext_proc server to forward upstream
327
    @GuardedBy("streamLock")
1✔
328
    private final Queue<ByteString> pendingUpstreamBodyMessages =
329
        new ConcurrentLinkedQueue<>();
330
    // Path 4: Outstanding requests from downstream for pulling responses
331
    @GuardedBy("streamLock")
1✔
332
    private int downstreamRequestsPending = 0;
333
    // Buffered mutated response bodies from ext_proc server
334
    @GuardedBy("streamLock")
1✔
335
    private final Queue<ByteString> pendingMutatedResponseBodies =
336
        new ConcurrentLinkedQueue<>();
337

338
    // Accumulated client window updates to send to ext_proc
339
    @GuardedBy("streamLock")
1✔
340
    private long accumulatedWindowUpdateSidestreamToUpstream = 0;
341
    @GuardedBy("streamLock")
1✔
342
    private long accumulatedWindowUpdateSidestreamToDownstream = 0;
343

344
    // Flag to track if FlowControlInit was sent in the initial message
345
    @GuardedBy("streamLock")
1✔
346
    private boolean flowControlInitSent = false;
347

348
    private final MethodDescriptor<?, ?> method;
349
    private final Channel channel;
350
    private final MetricRecorder metricsRecorder;
351
    private final String target;
352
    private volatile Context callContext = Context.ROOT;
1✔
353

354
    private volatile long clientHeadersStartNanos;
355
    private volatile long clientHalfCloseStartNanos;
356
    private volatile long serverHeadersStartNanos;
357
    private volatile long serverTrailersStartNanos;
358

359
    private boolean protocolConfigSent = false;
1✔
360
    private ImmutableMap<String, Struct> collectedAttributes;
361
    private boolean requestAttributesSent = false;
1✔
362
    @Nullable private volatile Metadata requestHeaders;
363
    final AtomicReference<DataPlaneCallState> dataPlaneCallState =
1✔
364
        new AtomicReference<>(DataPlaneCallState.IDLE);
365
    final AtomicReference<ExtProcStreamState> extProcStreamState =
1✔
366
        new AtomicReference<>(ExtProcStreamState.ACTIVE);
367
    final AtomicBoolean passThroughMode = new AtomicBoolean(false);
1✔
368
    final AtomicBoolean requestSideClosed = new AtomicBoolean(false);
1✔
369
    final AtomicBoolean appHalfClosed = new AtomicBoolean(false);
1✔
370
    final AtomicBoolean isProcessingTrailers = new AtomicBoolean(false);
1✔
371
    final AtomicBoolean pendingHalfClose = new AtomicBoolean(false);
1✔
372
    final AtomicBoolean bodyMessageSentToExtProc = new AtomicBoolean(false);
1✔
373

374
    protected DataPlaneClientCall(
375
        DataPlaneDelayedCall<InputStream, InputStream> delayedCall,
376
        ClientCall<InputStream, InputStream> rawCall,
377
        ExternalProcessorGrpc.ExternalProcessorStub stub,
378
        ExternalProcessorFilterConfig config,
379
        Optional<HeaderMutationRulesConfig> mutationRulesConfig,
380
        ScheduledExecutorService scheduler,
381
        MethodDescriptor<?, ?> method,
382
        Channel channel,
383
        MetricRecorder metricsRecorder,
384
        String target) {
385
      super(delayedCall);
1✔
386
      this.delayedCall = delayedCall;
1✔
387
      this.rawCall = rawCall;
1✔
388
      this.stub = stub;
1✔
389
      this.config = config;
1✔
390
      this.currentProcessingMode = config.getExternalProcessor().getProcessingMode();
1✔
391
      this.mutationFilter = new HeaderMutationFilter(mutationRulesConfig);
1✔
392
      this.scheduler = scheduler;
1✔
393
      this.method = method;
1✔
394
      this.channel = channel;
1✔
395
      this.metricsRecorder = checkNotNull(metricsRecorder, "metricsRecorder");
1✔
396
      this.target = checkNotNull(target, "target");
1✔
397
    }
1✔
398

399
    private boolean activateCall() {
400
      if ((extProcStreamState.get() == ExtProcStreamState.FAILED
1✔
401
              && !config.getFailureModeAllow()
1✔
402
              && !config.getObservabilityMode())
1✔
403
          || !dataPlaneCallState.compareAndSet(
1✔
404
              DataPlaneCallState.IDLE, DataPlaneCallState.ACTIVE)) {
405
        return false;
1✔
406
      }
407
      if (clientHeadersStartNanos > 0) {
1✔
408
        long durationNanos = System.nanoTime() - clientHeadersStartNanos;
1✔
409
        recordDuration(clientHeadersDuration, durationNanos);
1✔
410
        clientHeadersStartNanos = 0;
1✔
411
      }
412
      Runnable toRun = delayedCall.setCall(rawCall);
1✔
413
      if (toRun != null) {
1✔
414
        callContext.run(toRun);
1✔
415
      }
416
      drainPendingRequests();
1✔
417
      onReadyNotify();
1✔
418
      return true;
1✔
419
    }
420

421
    private void recordDuration(DoubleHistogramMetricInstrument instrument, long durationNanos) {
422
      if (instrument != null) {
1✔
423
        double durationSecs = (double) durationNanos / 1_000_000_000.0;
1✔
424
        metricsRecorder.recordDoubleHistogram(
1✔
425
            instrument,
426
            durationSecs,
427
            ImmutableList.of(target),
1✔
428
            ImmutableList.of());
1✔
429
      }
430
    }
1✔
431

432
    /**
433
     * Validates whether the body response uses unsupported gRPC message compression.
434
     * If compression is unsupported, this method will cancel the call, transition the
435
     * stream to a failed state, send an error to the external processor, and return false.
436
     *
437
     * @param bodyResponse the response to validate
438
     * @return true if validation passes (compression is supported or not used),
439
     *     false if validation fails
440
     */
441
    private boolean validateCompressionSupport(BodyResponse bodyResponse) {
442
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1✔
443
        BodyMutation mutation = 
1✔
444
            bodyResponse.getResponse().getBodyMutation();
1✔
445
        if (mutation.hasStreamedResponse()
1✔
446
            && mutation.getStreamedResponse().getGrpcMessageCompressed()) {
1✔
447
          StatusRuntimeException ex = Status.UNAVAILABLE
1✔
448
              .withDescription("gRPC message compression not supported in ext_proc")
1✔
449
              .asRuntimeException();
1✔
450
          synchronized (streamLock) {
1✔
451
            if (markExtProcStreamFailed(extProcStreamState)) {
1✔
452
              if (extProcClientCallRequestObserver != null) {
1✔
453
                extProcClientCallRequestObserver.onError(ex);
1✔
454
                extProcClientCallRequestObserver = null;
1✔
455
              }
456
            }
457
          }
1✔
458
          activateCall();
1✔
459
          cancelDownstream("gRPC message compression not supported in ext_proc", ex);
1✔
460
          closeExtProcStream();
1✔
461
          return false;
1✔
462
        }
463
      }
464
      return true;
1✔
465
    }
466

467
    @Override
468
    public void start(Listener<InputStream> responseListener, Metadata headers) {
469
      this.callContext = Context.current();
1✔
470
      clientHeadersStartNanos = System.nanoTime();
1✔
471
      this.requestHeaders = headers;
1✔
472
      this.wrappedListener = new DataPlaneListener(responseListener, this);
1✔
473

474
      // DelayedClientCall.start will buffer the listener and headers until setCall is called.
475
      super.start(wrappedListener, headers);
1✔
476

477
      stub.process(new ClientResponseObserver<ProcessingRequest, ProcessingResponse>() {
1✔
478
        @Override
479
        public void beforeStart(ClientCallStreamObserver<ProcessingRequest> requestStream) {
480
          synchronized (streamLock) {
1✔
481
            extProcClientCallRequestObserver = requestStream;
1✔
482
          }
1✔
483
          requestStream.setOnReadyHandler(DataPlaneClientCall.this::onExtProcStreamReady);
1✔
484
        }
1✔
485

486
        @Override
487
        public void onNext(ProcessingResponse response) {
488
          try {
489
            if (config.getObservabilityMode()) {
1✔
490
              return;
1✔
491
            }
492

493
            if (response.hasServerWindowUpdate()) {
1✔
494
              ProcessingResponse.ServerWindowUpdate update = response.getServerWindowUpdate();
1✔
495
              boolean wasReady;
496
              synchronized (streamLock) {
1✔
497
                wasReady = isReady();
1✔
498
                downstreamToSidestreamWindow += update.getWindowIncrementDownstreamToSidestream();
1✔
499
                upstreamToSidestreamWindow += update.getWindowIncrementUpstreamToSidestream();
1✔
500
                drainPendingRequestBodyMessages();
1✔
501
                drainPendingRequests();
1✔
502
                if (wrappedListener != null) {
1✔
503
                  wrappedListener.drainSavedMessages();
1✔
504
                }
505
              }
1✔
506
              // If isReady() becomes true (depends on updated downstreamToSidestreamWindow),
507
              // notify the client application via onReadyNotify() (runs unlocked).
508
              if (!wasReady && isReady()) {
1✔
509
                onReadyNotify();
1✔
510
              }
511
            }
512

513
            if (response.hasImmediateResponse()) {
1✔
514
              if (config.getDisableImmediateResponse()) {
1✔
515
                internalOnError(Status.UNAVAILABLE
1✔
516
                    .withDescription(
1✔
517
                        "Immediate response is disabled but received from external processor")
518
                    .asRuntimeException());
1✔
519
                return;
1✔
520
              }
521
              handleImmediateResponse(response.getImmediateResponse(), wrappedListener);
1✔
522
              return;
1✔
523
            }
524

525
            if (response.hasRequestHeaders()) {
1✔
526
              EventType expected = expectedRequestResponse;
1✔
527
              if (expected == null || expected != EventType.REQUEST_HEADERS) {
1✔
528
                internalOnError(Status.UNAVAILABLE
×
529
                    .withDescription("Protocol error: received response out of order. Expected: " 
×
530
                        + expected + ", Received: REQUEST_HEADERS")
531
                    .asRuntimeException());
×
532
                return;
×
533
              }
534
              expectedRequestResponse = null;
1✔
535
            } else if (response.hasResponseHeaders()) {
1✔
536
              EventType expected = expectedResponseResponse;
1✔
537
              if (expected == null || expected != EventType.RESPONSE_HEADERS) {
1✔
538
                internalOnError(Status.UNAVAILABLE
1✔
539
                    .withDescription("Protocol error: received response out of order. Expected: " 
1✔
540
                        + expected + ", Received: RESPONSE_HEADERS")
541
                    .asRuntimeException());
1✔
542
                return;
1✔
543
              }
544
              expectedResponseResponse = null;
1✔
545
            } else if (response.hasResponseTrailers()) {
1✔
546
              EventType expected = expectedResponseResponse;
1✔
547
              if (expected == null || expected != EventType.RESPONSE_TRAILERS) {
1✔
548
                internalOnError(Status.UNAVAILABLE
1✔
549
                    .withDescription("Protocol error: received response out of order. Expected: " 
1✔
550
                        + expected + ", Received: RESPONSE_TRAILERS")
551
                    .asRuntimeException());
1✔
552
                return;
1✔
553
              }
554
              expectedResponseResponse = null;
1✔
555
            } else if (response.hasRequestBody()) {
1✔
556
              EventType expected = expectedRequestResponse;
1✔
557
              if (expected == EventType.REQUEST_HEADERS) {
1✔
558
                internalOnError(Status.UNAVAILABLE
1✔
559
                    .withDescription(
1✔
560
                        "Protocol error: received request_body before request_headers response.")
561
                    .asRuntimeException());
1✔
562
                return;
1✔
563
              }
564
            } else if (response.hasResponseBody()) {
1✔
565
              EventType expected = expectedResponseResponse;
1✔
566
              if (expected == EventType.RESPONSE_HEADERS) {
1✔
567
                internalOnError(Status.UNAVAILABLE
1✔
568
                    .withDescription(
1✔
569
                        "Protocol error: received response_body before headers response.")
570
                    .asRuntimeException());
1✔
571
                return;
1✔
572
              }
573
            }
574

575
            if (response.getRequestDrain()) {
1✔
576
              extProcStreamState.set(ExtProcStreamState.DRAINING);
1✔
577
              halfCloseExtProcStream();
1✔
578
              activateCall();
1✔
579
            }
580

581
            // 1. Client Headers
582
            if (response.hasRequestHeaders()) {
1✔
583
              if (response.getRequestHeaders().hasResponse()) {
1✔
584
                if (response.getRequestHeaders().getResponse().getStatus()
1✔
585
                    == CommonResponse.ResponseStatus.CONTINUE_AND_REPLACE) {
586
                  internalOnError(Status.UNAVAILABLE
1✔
587
                      .withDescription("CONTINUE_AND_REPLACE is not supported")
1✔
588
                      .asRuntimeException());
1✔
589
                  return;
1✔
590
                }
591
                applyHeaderMutations(
1✔
592
                    requestHeaders,
1✔
593
                    response.getRequestHeaders().getResponse().getHeaderMutation(),
1✔
594
                    mutationFilter,
1✔
595
                    mutator);
1✔
596
              }
597
              activateCall();
1✔
598
            }
599
            // 2. Client Message (Request Body)
600
            else if (response.hasRequestBody()) {
1✔
601
              if (validateCompressionSupport(response.getRequestBody())) {
1✔
602
                handleRequestBodyResponse(response.getRequestBody());
1✔
603
              }
604
            }
605
            // 4. Server Headers
606
            else if (response.hasResponseHeaders()) {
1✔
607
              if (response.getResponseHeaders().hasResponse()) {
1✔
608
                if (response.getResponseHeaders().getResponse().getStatus()
1✔
609
                    == CommonResponse.ResponseStatus.CONTINUE_AND_REPLACE) {
610
                  internalOnError(Status.UNAVAILABLE
1✔
611
                      .withDescription("CONTINUE_AND_REPLACE is not supported")
1✔
612
                      .asRuntimeException());
1✔
613
                  return;
1✔
614
                }
615
                Metadata target = wrappedListener.isTrailersOnly()
1✔
616
                    ? wrappedListener.getSavedTrailers() : wrappedListener.getSavedHeaders();
1✔
617
                applyHeaderMutations(
1✔
618
                    target,
619
                    response.getResponseHeaders().getResponse().getHeaderMutation(),
1✔
620
                    mutationFilter,
1✔
621
                    mutator);
1✔
622
              }
623
              if (wrappedListener.isTrailersOnly()) {
1✔
624
                wrappedListener.proceedWithClose();
1✔
625
              } else {
626
                wrappedListener.proceedWithHeaders();
1✔
627
              }
628
            }
629
            // 5. Server Message (Response Body)
630
            else if (response.hasResponseBody()) {
1✔
631
              if (validateCompressionSupport(response.getResponseBody())) {
1✔
632
                handleResponseBodyResponse(response.getResponseBody(), wrappedListener);
1✔
633
              }
634
            }
635
            // 6. Response Trailers
636
            else if (response.hasResponseTrailers()) {
1✔
637
              if (response.getResponseTrailers().hasHeaderMutation()) {
1✔
638
                applyHeaderMutations(
1✔
639
                    wrappedListener.getSavedTrailers(),
1✔
640
                    response.getResponseTrailers().getHeaderMutation(),
1✔
641
                    mutationFilter,
1✔
642
                    mutator);
1✔
643
              }
644
              wrappedListener.proceedWithClose();
1✔
645
            }
646

647
            checkEndOfStream(response);
1✔
648
          } catch (Throwable t) {
1✔
649
            internalOnError(t);
1✔
650
          }
1✔
651
        }
1✔
652

653
        @Override
654
        public void onError(Throwable t) {
655
          if (markExtProcStreamFailed(extProcStreamState)) {
1✔
656
            synchronized (streamLock) {
1✔
657
              extProcClientCallRequestObserver = null;
1✔
658
            }
1✔
659
            if (config.getObservabilityMode()
1✔
660
                || (config.getFailureModeAllow() && !bodyMessageSentToExtProc.get())) {
1✔
661
              handleFailOpen(wrappedListener);
1✔
662
            } else {
663
              String message = "External processor stream failed";
1✔
664
              cancelDownstream(message, t);
1✔
665
              wrappedListener.proceedWithClose();
1✔
666
            }
667
          }
668
        }
1✔
669

670
        @Override
671
        public void onCompleted() {
672
          if (markExtProcStreamCompleted(extProcStreamState)) {
1✔
673
            synchronized (streamLock) {
1✔
674
              extProcClientCallRequestObserver = null;
1✔
675
            }
1✔
676
            handleFailOpen(wrappedListener);
1✔
677
          }
678
        }
1✔
679
      });
680

681
      this.collectedAttributes = collectAttributes(
1✔
682
          config.getRequestAttributes(), method, channel.authority(), headers);
1✔
683

684
      boolean sendRequestHeaders =
1✔
685
          currentProcessingMode.getRequestHeaderMode() == ProcessingMode.HeaderSendMode.SEND
1✔
686
          || currentProcessingMode.getRequestHeaderMode()
1✔
687
              == ProcessingMode.HeaderSendMode.DEFAULT;
688

689
      if (sendRequestHeaders) {
1✔
690
        sendToExtProc(ProcessingRequest.newBuilder()
1✔
691
            .setRequestHeaders(HttpHeaders.newBuilder()
1✔
692
                .setHeaders(toHeaderMap(headers, config.getForwardRulesConfig()))
1✔
693
                .setEndOfStream(false)
1✔
694
                .build())
1✔
695
            .build());
1✔
696
      }
697

698
      if (config.getObservabilityMode() || !sendRequestHeaders) {
1✔
699
        activateCall();
1✔
700
      }
701
    }
1✔
702

703
    private void sendToExtProc(ProcessingRequest request) {
704
      synchronized (streamLock) {
1✔
705
        if (extProcStreamState.get().isCompleted() || extProcClientCallRequestObserver == null) {
1✔
706
          return;
1✔
707
        }
708
        
709
        if (request.hasRequestHeaders()) {
1✔
710
          expectedRequestResponse = EventType.REQUEST_HEADERS;
1✔
711
        } else if (request.hasResponseHeaders()) {
1✔
712
          expectedResponseResponse = EventType.RESPONSE_HEADERS;
1✔
713
        } else if (request.hasResponseTrailers()) {
1✔
714
          expectedResponseResponse = EventType.RESPONSE_TRAILERS;
1✔
715
        }
716

717
        ProcessingRequest requestToSend = request;
1✔
718
        if (!protocolConfigSent) {
1✔
719
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
720
              .setProtocolConfig(ProtocolConfiguration.newBuilder()
1✔
721
                  .setRequestBodyMode(currentProcessingMode.getRequestBodyMode())
1✔
722
                  .setResponseBodyMode(currentProcessingMode.getResponseBodyMode())
1✔
723
                  .build())
1✔
724
              .build();
1✔
725
          protocolConfigSent = true;
1✔
726
        }
727

728
        boolean isClientServerMessage =
1✔
729
            requestToSend.hasRequestHeaders() || requestToSend.hasRequestBody();
1✔
730
        if (isClientServerMessage
1✔
731
            && !requestAttributesSent
732
            && collectedAttributes != null
733
            && !collectedAttributes.isEmpty()) {
1✔
734
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
735
              .putAllAttributes(collectedAttributes)
1✔
736
              .build();
1✔
737
          requestAttributesSent = true;
1✔
738
        }
739

740
        if (config.getObservabilityMode()) {
1✔
741
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
742
              .setObservabilityMode(true)
1✔
743
              .build();
1✔
744
        } else if (!flowControlInitSent) {
1✔
745
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
746
              .setFlowControlInit(ProcessingRequest.FlowControlInit.newBuilder()
1✔
747
                  .setInitialWindowDownstreamToSidestream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
748
                  .setInitialWindowSidestreamToUpstream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
749
                  .setInitialWindowUpstreamToSidestream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
750
                  .setInitialWindowSidestreamToDownstream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
751
                  .build())
1✔
752
              .build();
1✔
753
          flowControlInitSent = true;
1✔
754
        }
755

756
        extProcClientCallRequestObserver.onNext(requestToSend);
1✔
757
      }
1✔
758
    }
1✔
759

760
    // Note: This method not only modifies the builder, but has the side effect of modifying
761
    // the window update bookkeeping.
762
    @GuardedBy("streamLock")
763
    void mergeAccumulatedWindowUpdates(ProcessingRequest.Builder requestBuilder) {
764
      long incrementUpstream = accumulatedWindowUpdateSidestreamToUpstream;
1✔
765
      long incrementDownstream = accumulatedWindowUpdateSidestreamToDownstream;
1✔
766

767
      if (incrementUpstream > 0 || incrementDownstream > 0) {
1✔
768
        requestBuilder.setClientWindowUpdate(
1✔
769
            ProcessingRequest.ClientWindowUpdate.newBuilder()
1✔
770
                .setWindowIncrementSidestreamToUpstream(incrementUpstream)
1✔
771
                .setWindowIncrementSidestreamToDownstream(incrementDownstream)
1✔
772
                .build());
1✔
773
        accumulatedWindowUpdateSidestreamToUpstream -= incrementUpstream;
1✔
774
        accumulatedWindowUpdateSidestreamToDownstream -= incrementDownstream;
1✔
775
        sidestreamToUpstreamWindow += incrementUpstream;
1✔
776
        sidestreamToDownstreamWindow += incrementDownstream;
1✔
777
      }
778
    }
1✔
779

780
    private void trySendAccumulatedWindowUpdates() {
781
      synchronized (streamLock) {
1✔
782
        if (extProcStreamState.get().isCompleted()) {
1✔
783
          return;
1✔
784
        }
785
        long incrementUpstream = accumulatedWindowUpdateSidestreamToUpstream;
1✔
786
        long incrementDownstream = accumulatedWindowUpdateSidestreamToDownstream;
1✔
787

788
        boolean shouldSend = (incrementUpstream > 0 || incrementDownstream > 0) && (
1✔
789
            (incrementUpstream >= WINDOW_UPDATE_THRESHOLD)
790
            || (incrementDownstream >= WINDOW_UPDATE_THRESHOLD)
791
            || (sidestreamToUpstreamWindow <= 0 && accumulatedWindowUpdateSidestreamToUpstream > 0)
792
            || (sidestreamToDownstreamWindow <= 0
793
                && accumulatedWindowUpdateSidestreamToDownstream > 0)
794
        );
795

796
        if (shouldSend) {
1✔
797
          accumulatedWindowUpdateSidestreamToUpstream -= incrementUpstream;
1✔
798
          accumulatedWindowUpdateSidestreamToDownstream -= incrementDownstream;
1✔
799
          sidestreamToUpstreamWindow += incrementUpstream;
1✔
800
          sidestreamToDownstreamWindow += incrementDownstream;
1✔
801

802
          sendToExtProc(ProcessingRequest.newBuilder()
1✔
803
              .setClientWindowUpdate(ProcessingRequest.ClientWindowUpdate.newBuilder()
1✔
804
                  .setWindowIncrementSidestreamToUpstream(incrementUpstream)
1✔
805
                  .setWindowIncrementSidestreamToDownstream(incrementDownstream)
1✔
806
                  .build())
1✔
807
              .build());
1✔
808
        }
809
      }
1✔
810
    }
1✔
811

812
    private void onExtProcStreamReady() {
813
      drainPendingRequests();
1✔
814
      onReadyNotify();
1✔
815
    }
1✔
816

817
    void drainPendingRequests() {
818
      synchronized (streamLock) {
1✔
819
        if (config.getObservabilityMode()
1✔
820
            || currentProcessingMode.getResponseBodyMode() != ProcessingMode.BodySendMode.GRPC
1✔
821
            || extProcStreamState.get().isCompleted()) {
1✔
822
          int toRequest = pendingRequests.getAndSet(0);
1✔
823
          if (toRequest > 0) {
1✔
824
            super.request(toRequest);
1✔
825
          }
826
          return;
1✔
827
        }
828

829
        // Normal mode flow control: pull 1 message at a time
830
        if (isSidecarReady() && upstreamToSidestreamWindow > 0 && pendingRequests.get() > 0) {
1✔
831
          super.request(1);
1✔
832
          pendingRequests.decrementAndGet();
1✔
833
        }
834
      }
1✔
835
    }
1✔
836

837
    private void closeExtProcStream() {
838
      synchronized (streamLock) {
1✔
839
        if (markExtProcStreamCompleted(extProcStreamState)) {
1✔
840
          if (extProcClientCallRequestObserver != null) {
1✔
841
            extProcClientCallRequestObserver.onCompleted();
1✔
842
            extProcClientCallRequestObserver = null;
1✔
843
          }
844
        }
845
      }
1✔
846
    }
1✔
847

848
    private void internalOnError(Throwable t) {
849
      if (markExtProcStreamFailed(extProcStreamState)) {
1✔
850
        synchronized (streamLock) {
1✔
851
          if (extProcClientCallRequestObserver != null) {
1✔
852
            extProcClientCallRequestObserver.onError(t);
1✔
853
            extProcClientCallRequestObserver = null;
1✔
854
          }
855
        }
1✔
856
        if (config.getObservabilityMode()
1✔
857
            || (config.getFailureModeAllow() && !bodyMessageSentToExtProc.get())) {
1✔
858
          handleFailOpen(wrappedListener);
1✔
859
        } else {
860
          String message = "External processor stream failed";
1✔
861
          cancelDownstream(message, t);
1✔
862
          wrappedListener.proceedWithClose();
1✔
863
        }
864
      }
865
    }
1✔
866

867
    private void halfCloseExtProcStream() {
868
      synchronized (streamLock) {
1✔
869
        if (!extProcStreamState.get().isCompleted() && extProcClientCallRequestObserver != null) {
1✔
870
          extProcClientCallRequestObserver.onCompleted();
1✔
871
        }
872
      }
1✔
873
    }
1✔
874

875
    private void onReadyNotify() {
876
      wrappedListener.onReadyNotify();
1✔
877
    }
1✔
878

879
    void onReady() {
880
      boolean isPassThrough;
881
      boolean isCompleted;
882
      boolean isDraining;
883

884
      synchronized (streamLock) {
1✔
885
        isPassThrough = passThroughMode.get();
1✔
886
        ExtProcStreamState state = extProcStreamState.get();
1✔
887
        isCompleted = state.isCompleted();
1✔
888
        isDraining = state.isDraining();
1✔
889
      }
1✔
890

891
      if (isPassThrough) {
1✔
892
        onReadyNotify();
×
893
        return;
×
894
      }
895

896
      if (isCompleted) {
1✔
897
        drainPendingDrainingMessages();
1✔
898
        return;
1✔
899
      }
900

901
      // Normal or Draining operation
902
      drainPendingUpstreamBodyMessages();
1✔
903
      if (!isDraining) {
1✔
904
        trySendAccumulatedWindowUpdates();
1✔
905
      }
906
      drainPendingRequests();
1✔
907
      onReadyNotify();
1✔
908
    }
1✔
909

910
    @GuardedBy("streamLock")
911
    boolean isSidecarReady() {
912
      ExtProcStreamState state = extProcStreamState.get();
1✔
913
      if (state.isCompleted()) {
1✔
914
        return true;
×
915
      }
916
      if (state.isDraining()) {
1✔
917
        return false;
1✔
918
      }
919
      ClientCallStreamObserver<ProcessingRequest> observer = extProcClientCallRequestObserver;
1✔
920
      return observer != null && observer.isReady();
1✔
921
    }
922

923
    @Override
924
    public boolean isReady() {
925
      if (passThroughMode.get()) {
1✔
926
        return super.isReady();
1✔
927
      }
928
      if (extProcStreamState.get().isCompleted()) {
1✔
929
        return super.isReady();
1✔
930
      }
931
      if (dataPlaneCallState.get() == DataPlaneCallState.IDLE && !config.getObservabilityMode()) {
1✔
932
        return false;
1✔
933
      }
934
      synchronized (streamLock) {
1✔
935
        boolean sidecarReady = isSidecarReady();
1✔
936
        if (config.getObservabilityMode()) {
1✔
937
          return super.isReady() && sidecarReady;
1✔
938
        }
939
        return downstreamToSidestreamWindow > 0 && sidecarReady
1✔
940
            && pendingRequestBodyMessages.isEmpty();
1✔
941
      }
942
    }
943

944
    @Override
945
    public void request(int numMessages) {
946
      if (passThroughMode.get() || extProcStreamState.get().isCompleted()) {
1✔
947
        super.request(numMessages);
1✔
948
        return;
1✔
949
      }
950
      if (!config.getObservabilityMode()
1✔
951
          && currentProcessingMode.getResponseBodyMode() != ProcessingMode.BodySendMode.GRPC) {
1✔
952
        super.request(numMessages);
1✔
953
        return;
1✔
954
      }
955
      synchronized (streamLock) {
1✔
956
        // We send response bodies to ext_proc server (either in normal GRPC mode or
957
        // observability mode).
958
        // Gated by ext_proc server readiness.
959
        // i.e. normal GRPC response body mode
960
        boolean normalFlowControl = !config.getObservabilityMode();
1✔
961

962
        if (normalFlowControl) {
1✔
963
          pendingRequests.addAndGet(numMessages);
1✔
964
          downstreamRequestsPending += numMessages;
1✔
965
          drainPendingMutatedResponseBodies();
1✔
966
          if (isSidecarReady()) {
1✔
967
            drainPendingRequests();
1✔
968
          }
969
        } else {
970
          // Observability mode: gate on readiness but pull all at once
971
          if (isSidecarReady()) {
1✔
972
            super.request(numMessages);
1✔
973
          } else {
974
            pendingRequests.addAndGet(numMessages);
1✔
975
          }
976
        }
977
      }
1✔
978
    }
1✔
979

980
    @Override
981
    public void sendMessage(InputStream message) {
982
      if (requestSideClosed.get()) {
1✔
983
        // External processor already closed the request stream. Discard further messages.
984
        return;
1✔
985
      }
986

987
      if (passThroughMode.get()) {
1✔
988
        super.sendMessage(message);
1✔
989
        return;
1✔
990
      }
991

992
      synchronized (streamLock) {
1✔
993
        if (passThroughMode.get()) {
1✔
994
          super.sendMessage(message);
×
995
          return;
×
996
        }
997

998
        ExtProcStreamState state = extProcStreamState.get();
1✔
999
        if (state.isDraining() || state.isCompleted()) {
1✔
1000
          if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1✔
1001
            super.sendMessage(message);
1✔
1002
            return;
1✔
1003
          }
1004
          try {
1005
            ByteString copiedBody = ByteString.readFrom(message);
1✔
1006
            pendingDrainingMessages.add(new KnownLengthInputStream(copiedBody));
1✔
1007
          } catch (IOException e) {
×
1008
            cancelDownstream("Failed to copy outbound message for buffering", e);
×
1009
          }
1✔
1010
          return;
1✔
1011
        }
1012

1013
        if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1✔
1014
          super.sendMessage(message);
1✔
1015
          return;
1✔
1016
        }
1017

1018
        // Mode is GRPC
1019
        try {
1020
          ByteString bodyByteString = outboundStreamToByteString(message);
1✔
1021
          if (config.getObservabilityMode()) {
1✔
1022
            sendToExtProc(ProcessingRequest.newBuilder()
×
1023
                .setRequestBody(HttpBody.newBuilder()
×
1024
                    .setBody(bodyByteString)
×
1025
                    .setEndOfStream(false)
×
1026
                    .build())
×
1027
                .build());
×
1028
            bodyMessageSentToExtProc.set(true);
×
1029
            super.sendMessage(new KnownLengthInputStream(bodyByteString));
×
1030
          } else {
1031
            if (downstreamToSidestreamWindow <= 0 || !pendingRequestBodyMessages.isEmpty()) {
1✔
1032
              pendingRequestBodyMessages.add(bodyByteString);
1✔
1033
            } else {
1034
              sendRequestBodyToExtProc(bodyByteString);
1✔
1035
            }
1036
          }
1037
        } catch (IOException e) {
×
1038
          cancelDownstream("Failed to serialize message for External Processor", e);
×
1039
        }
1✔
1040
      }
1✔
1041
    }
1✔
1042

1043
    @GuardedBy("streamLock")
1044
    private void sendRequestBodyToExtProc(ByteString body) {
1045
      downstreamToSidestreamWindow -= body.size();
1✔
1046
      ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1✔
1047
          .setRequestBody(HttpBody.newBuilder()
1✔
1048
              .setBody(body)
1✔
1049
              .setEndOfStream(false)
1✔
1050
              .build());
1✔
1051
      mergeAccumulatedWindowUpdates(builder);
1✔
1052
      sendToExtProc(builder.build());
1✔
1053
      bodyMessageSentToExtProc.set(true);
1✔
1054
    }
1✔
1055

1056
    @GuardedBy("streamLock")
1057
    private void drainPendingRequestBodyMessages() {
1058
      while (downstreamToSidestreamWindow > 0 && !pendingRequestBodyMessages.isEmpty()) {
1✔
1059
        ByteString body = pendingRequestBodyMessages.poll();
1✔
1060
        sendRequestBodyToExtProc(body);
1✔
1061
      }
1✔
1062
      if (pendingRequestBodyMessages.isEmpty() && pendingHalfClose.compareAndSet(true, false)) {
1✔
1063
        halfClose();
1✔
1064
      }
1065
    }
1✔
1066

1067
    private void proceedWithHalfClose() {
1068
      if (clientHalfCloseStartNanos > 0) {
1✔
1069
        long durationNanos = System.nanoTime() - clientHalfCloseStartNanos;
1✔
1070
        recordDuration(clientHalfCloseDuration, durationNanos);
1✔
1071
        clientHalfCloseStartNanos = 0;
1✔
1072
      }
1073
      super.halfClose();
1✔
1074
    }
1✔
1075

1076
    @Override
1077
    public void halfClose() {
1078
      if (appHalfClosed.compareAndSet(false, true)) {
1✔
1079
        clientHalfCloseStartNanos = System.nanoTime();
1✔
1080
      }
1081
      if (passThroughMode.get()) {
1✔
1082
        if (requestSideClosed.compareAndSet(false, true)) {
1✔
1083
          proceedWithHalfClose();
1✔
1084
        }
1085
        return;
1✔
1086
      }
1087

1088
      if (extProcStreamState.get().isCompleted()) {
1✔
1089
        if (passThroughMode.get()) {
1✔
1090
          if (requestSideClosed.compareAndSet(false, true)) {
×
1091
            proceedWithHalfClose();
×
1092
          }
1093
        } else {
1094
          pendingHalfClose.set(true);
1✔
1095
        }
1096
        return;
1✔
1097
      }
1098

1099
      if (extProcStreamState.get().isDraining()) {
1✔
1100
        boolean canProceed = false;
1✔
1101
        synchronized (streamLock) {
1✔
1102
          if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE
1✔
1103
              || (!bodyMessageSentToExtProc.get() && pendingDrainingMessages.isEmpty())) {
1✔
1104
            canProceed = true;
1✔
1105
          }
1106
        }
1✔
1107
        if (canProceed) {
1✔
1108
          if (requestSideClosed.compareAndSet(false, true)) {
1✔
1109
            proceedWithHalfClose();
1✔
1110
          }
1111
        } else {
1112
          pendingHalfClose.set(true);
1✔
1113
        }
1114
        return;
1✔
1115
      }
1116

1117
      if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1✔
1118
        if (requestSideClosed.compareAndSet(false, true)) {
1✔
1119
          proceedWithHalfClose();
1✔
1120
        }
1121
        return;
1✔
1122
      }
1123

1124
      // Mode is GRPC
1125
      synchronized (streamLock) {
1✔
1126
        if (!pendingRequestBodyMessages.isEmpty()) {
1✔
1127
          pendingHalfClose.set(true);
1✔
1128
          return;
1✔
1129
        }
1130

1131
        ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1✔
1132
            .setRequestBody(HttpBody.newBuilder()
1✔
1133
                .setEndOfStream(true)
1✔
1134
                .setEndOfStreamWithoutMessage(true)
1✔
1135
                .build());
1✔
1136
        mergeAccumulatedWindowUpdates(builder);
1✔
1137
        sendToExtProc(builder.build());
1✔
1138
      }
1✔
1139
    }
1✔
1140

1141
    void cancelDownstream(@Nullable String message, @Nullable Throwable cause) {
1142
      delayedCall.cancel(message, cause);
1✔
1143
    }
1✔
1144

1145
    @Override
1146
    public void cancel(@Nullable String message, @Nullable Throwable cause) {
1147
      synchronized (streamLock) {
1✔
1148
        if (markExtProcStreamFailed(extProcStreamState)) {
1✔
1149
          if (extProcClientCallRequestObserver != null) {
1✔
1150
            extProcClientCallRequestObserver.onError(
1✔
1151
                Status.CANCELLED
1152
                    .withDescription(message)
1✔
1153
                    .withCause(cause)
1✔
1154
                    .asRuntimeException());
1✔
1155
            extProcClientCallRequestObserver = null;
1✔
1156
          }
1157
        }
1158
      }
1✔
1159
      cancelDownstream(message, cause);
1✔
1160
    }
1✔
1161

1162
    private void handleRequestBodyResponse(BodyResponse bodyResponse) {
1163
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1✔
1164
        BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
1✔
1165
        if (mutation.hasStreamedResponse()) {
1✔
1166
          StreamedBodyResponse streamed = mutation.getStreamedResponse();
1✔
1167
          boolean isEndOfStream = streamed.getEndOfStream();
1✔
1168
          boolean isEndOfStreamWithoutMessage =
1✔
1169
              isEndOfStream && streamed.getEndOfStreamWithoutMessage();
1✔
1170
          if (!isEndOfStreamWithoutMessage) {
1✔
1171
            ByteString body = streamed.getBody();
1✔
1172
            boolean sendImmediately = false;
1✔
1173
            synchronized (streamLock) {
1✔
1174
              sidestreamToUpstreamWindow -= body.size();
1✔
1175
              if (pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
1✔
1176
                sendImmediately = true;
1✔
1177
                accumulatedWindowUpdateSidestreamToUpstream += body.size();
1✔
1178
              } else {
1179
                pendingUpstreamBodyMessages.add(body);
1✔
1180
              }
1181
            }
1✔
1182
            if (sendImmediately) {
1✔
1183
              super.sendMessage(new KnownLengthInputStream(body));
1✔
1184
              trySendAccumulatedWindowUpdates();
1✔
1185
            }
1186
          }
1187
          if (isEndOfStream) {
1✔
1188
            synchronized (streamLock) {
1✔
1189
              if (pendingUpstreamBodyMessages.isEmpty()) {
1✔
1190
                if (requestSideClosed.compareAndSet(false, true)) {
1✔
1191
                  proceedWithHalfClose();
1✔
1192
                }
1193
              } else {
1194
                pendingUpstreamHalfClose.set(true);
1✔
1195
              }
1196
            }
1✔
1197
          }
1198
        }
1199
      }
1200
    }
1✔
1201

1202
    private void handleResponseBodyResponse(
1203
        BodyResponse bodyResponse, DataPlaneListener listener) {
1204
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1✔
1205
        BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
1✔
1206
        if (mutation.hasStreamedResponse()) {
1✔
1207
          StreamedBodyResponse streamed = mutation.getStreamedResponse();
1✔
1208
          ByteString body = streamed.getBody();
1✔
1209
          final int bodySize = body.size();
1✔
1210
          synchronized (streamLock) {
1✔
1211
            sidestreamToDownstreamWindow -= bodySize;
1✔
1212
          }
1✔
1213
          deliverResponseBody(body, listener);
1✔
1214
        }
1215
      }
1216
    }
1✔
1217

1218
    private void deliverResponseBody(ByteString body, DataPlaneListener listener) {
1219
      boolean shouldDeliver = false;
1✔
1220
      synchronized (streamLock) {
1✔
1221
        if (downstreamRequestsPending > 0) {
1✔
1222
          downstreamRequestsPending--;
1✔
1223
          shouldDeliver = true;
1✔
1224
        } else {
1225
          pendingMutatedResponseBodies.add(body);
1✔
1226
        }
1227
      }
1✔
1228
      if (shouldDeliver) {
1✔
1229
        final int bodySize = body.size();
1✔
1230
        callContext.run(() -> {
1✔
1231
          try {
1232
            listener.onExternalBody(body);
1✔
1233
          } finally {
1234
            synchronized (streamLock) {
1✔
1235
              accumulatedWindowUpdateSidestreamToDownstream += bodySize;
1✔
1236
            }
1✔
1237
            trySendAccumulatedWindowUpdates();
1✔
1238
          }
1239
        });
1✔
1240
      }
1241
    }
1✔
1242

1243
    private void drainPendingMutatedResponseBodies() {
1244
      List<ByteString> toDeliver = new ArrayList<>();
1✔
1245
      synchronized (streamLock) {
1✔
1246
        while (downstreamRequestsPending > 0 && !pendingMutatedResponseBodies.isEmpty()) {
1✔
1247
          ByteString body = pendingMutatedResponseBodies.poll();
1✔
1248
          downstreamRequestsPending--;
1✔
1249
          pendingRequests.decrementAndGet();
1✔
1250
          toDeliver.add(body);
1✔
1251
        }
1✔
1252
      }
1✔
1253
      for (ByteString body : toDeliver) {
1✔
1254
        final int bodySize = body.size();
1✔
1255
        callContext.run(() -> {
1✔
1256
          try {
1257
            wrappedListener.onExternalBody(body);
1✔
1258
          } finally {
1259
            synchronized (streamLock) {
1✔
1260
              accumulatedWindowUpdateSidestreamToDownstream += bodySize;
1✔
1261
            }
1✔
1262
            trySendAccumulatedWindowUpdates();
1✔
1263
          }
1264
        });
1✔
1265
      }
1✔
1266
    }
1✔
1267

1268
    // Used to immediately flush any mutated response chunks that we already received and buffered
1269
    // before the stream failed, ensuring the application receives them in the correct order
1270
    void drainPendingMutatedResponseBodiesDirect(DataPlaneListener listener) {
1271
      List<ByteString> toDeliver = new ArrayList<>();
1✔
1272
      synchronized (streamLock) {
1✔
1273
        ByteString body;
1274
        while ((body = pendingMutatedResponseBodies.poll()) != null) {
1✔
1275
          toDeliver.add(body);
×
1276
        }
1277
      }
1✔
1278
      for (ByteString body : toDeliver) {
1✔
1279
        listener.onExternalBody(body);
×
1280
      }
×
1281
    }
1✔
1282

1283
    void drainPendingUpstreamBodyMessages() {
1284
      while (true) {
1285
        ByteString body = null;
1✔
1286
        boolean triggerHalfClose = false;
1✔
1287
        synchronized (streamLock) {
1✔
1288
          if (!pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
1✔
1289
            body = pendingUpstreamBodyMessages.poll();
1✔
1290
            accumulatedWindowUpdateSidestreamToUpstream += body.size();
1✔
1291
            if (pendingUpstreamBodyMessages.isEmpty()
1✔
1292
                && pendingUpstreamHalfClose.compareAndSet(true, false)) {
1✔
1293
              triggerHalfClose = true;
1✔
1294
            }
1295
          }
1296
        }
1✔
1297
        if (body == null) {
1✔
1298
          break;
1✔
1299
        }
1300
        super.sendMessage(new KnownLengthInputStream(body));
1✔
1301
        trySendAccumulatedWindowUpdates();
1✔
1302
        if (triggerHalfClose) {
1✔
1303
          if (requestSideClosed.compareAndSet(false, true)) {
1✔
1304
            proceedWithHalfClose();
1✔
1305
          }
1306
        }
1307
      }
1✔
1308
    }
1✔
1309

1310
    private void handleImmediateResponse(ImmediateResponse immediate, DataPlaneListener listener)
1311
        throws HeaderMutationDisallowedException {
1312
      Status status = Status.fromCodeValue(immediate.getGrpcStatus().getStatus());
1✔
1313
      if (!immediate.getDetails().isEmpty()) {
1✔
1314
        status = status.withDescription(immediate.getDetails());
1✔
1315
      }
1316

1317
      Metadata trailers = new Metadata();
1✔
1318
      if (immediate.hasHeaders()) {
1✔
1319
        applyHeaderMutations(trailers, immediate.getHeaders(), mutationFilter, mutator);
1✔
1320
      }
1321

1322
      listener.setImmediateResponse(status, trailers);
1✔
1323

1324
      if (isProcessingTrailers.get()) {
1✔
1325
        // If sent in response to a server trailers event, sets the status and optionally
1326
        // headers to be included in the trailers.
1327
        listener.unblockAfterStreamComplete();
1✔
1328
      } else {
1329
        // If sent in response to any other event, it will cause the data plane RPC to
1330
        // immediately fail with the specified status as if it were an out-of-band
1331
        // cancellation.
1332
        cancelDownstream(status.getDescription(), null);
1✔
1333
        listener.unblockAfterStreamComplete();
1✔
1334
      }
1335
      closeExtProcStream();
1✔
1336
    }
1✔
1337

1338
    private void drainPendingDrainingMessages() {
1339
      while (true) {
1340
        Object msg = null; // Can be ByteString or InputStream
1✔
1341
        boolean isMutated = false;
1✔
1342
        boolean triggerHalfClose = false;
1✔
1343

1344
        synchronized (streamLock) {
1✔
1345
          if (!pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
1✔
1346
            msg = pendingUpstreamBodyMessages.poll();
1✔
1347
            isMutated = true;
1✔
1348
          } else if (pendingUpstreamBodyMessages.isEmpty()
1✔
1349
              && !pendingRequestBodyMessages.isEmpty() && super.isReady()) {
1✔
1350
            msg = pendingRequestBodyMessages.poll();
1✔
1351
            isMutated = true;
1✔
1352
          } else if (pendingUpstreamBodyMessages.isEmpty()
1✔
1353
              && pendingRequestBodyMessages.isEmpty()
1✔
1354
              && !pendingDrainingMessages.isEmpty() && super.isReady()) {
1✔
1355
            msg = pendingDrainingMessages.poll();
1✔
1356
            isMutated = false;
1✔
1357
          }
1358

1359
          if (msg == null) {
1✔
1360
            if (pendingUpstreamBodyMessages.isEmpty()
1✔
1361
                && pendingRequestBodyMessages.isEmpty()
1✔
1362
                && pendingDrainingMessages.isEmpty()) {
1✔
1363
              passThroughMode.set(true);
1✔
1364
              if (appHalfClosed.get()) {
1✔
1365
                triggerHalfClose = true;
1✔
1366
              }
1367
            }
1368
          }
1369
        }
1✔
1370

1371
        if (msg == null) {
1✔
1372
          if (triggerHalfClose) {
1✔
1373
            if (requestSideClosed.compareAndSet(false, true)) {
1✔
1374
              proceedWithHalfClose();
1✔
1375
            }
1376
          }
1377
          break;
1378
        }
1379

1380
        if (isMutated) {
1✔
1381
          super.sendMessage(new KnownLengthInputStream((ByteString) msg));
1✔
1382
        } else {
1383
          super.sendMessage((InputStream) msg);
1✔
1384
        }
1385
      }
1✔
1386
    }
1✔
1387

1388
    private void handleFailOpen(DataPlaneListener listener) {
1389
      if (!activateCall()) {
1✔
1390
        drainPendingRequests();
1✔
1391
      }
1392
      listener.unblockAfterStreamComplete();
1✔
1393
      closeExtProcStream();
1✔
1394
    }
1✔
1395

1396
    private void checkEndOfStream(ProcessingResponse response) {
1397
      boolean terminal = false;
1✔
1398
      if (response.hasResponseTrailers()) {
1✔
1399
        terminal = true;
1✔
1400
      } else if (response.hasResponseHeaders() && wrappedListener.isTrailersOnly()) {
1✔
1401
        terminal = true;
1✔
1402
      }
1403

1404
      if (terminal) {
1✔
1405
        wrappedListener.unblockAfterStreamComplete();
1✔
1406
        closeExtProcStream();
1✔
1407
      }
1408
    }
1✔
1409

1410
    long getServerHeadersStartNanos() {
1411
      return serverHeadersStartNanos;
1✔
1412
    }
1413

1414
    void setServerHeadersStartNanos(long serverHeadersStartNanos) {
1415
      this.serverHeadersStartNanos = serverHeadersStartNanos;
1✔
1416
    }
1✔
1417

1418
    long getServerTrailersStartNanos() {
1419
      return serverTrailersStartNanos;
1✔
1420
    }
1421

1422
    void setServerTrailersStartNanos(long serverTrailersStartNanos) {
1423
      this.serverTrailersStartNanos = serverTrailersStartNanos;
1✔
1424
    }
1✔
1425

1426
    AtomicReference<ExtProcStreamState> getExtProcStreamState() {
1427
      return extProcStreamState;
1✔
1428
    }
1429

1430
    ProcessingMode getCurrentProcessingMode() {
1431
      return currentProcessingMode;
1✔
1432
    }
1433

1434
    AtomicBoolean getPassThroughMode() {
1435
      return passThroughMode;
1✔
1436
    }
1437

1438
    ExternalProcessorFilterConfig getConfig() {
1439
      return config;
1✔
1440
    }
1441

1442
    Context getCallContext() {
1443
      return callContext;
1✔
1444
    }
1445

1446
    ScheduledExecutorService getScheduler() {
1447
      return scheduler;
1✔
1448
    }
1449

1450
    AtomicBoolean getIsProcessingTrailers() {
1451
      return isProcessingTrailers;
1✔
1452
    }
1453
  }
1454

1455
  private static class DataPlaneListener extends SimpleForwardingClientCallListener<InputStream> {
1456
    private final DataPlaneClientCall dataPlaneClientCall;
1457
    // Path 3: Upstream response bodies queued because upstream to sidestream window not available,
1458
    // response headers not cleared by ext_proc or ext_proc stream draining
1459
    private final Queue<InputStream> savedMessages = new ConcurrentLinkedQueue<>();
1✔
1460
    private boolean inboundPassThrough = false;
1✔
1461
    @Nullable private volatile Metadata savedHeaders;
1462
    @Nullable private volatile Metadata savedTrailers;
1463
    @Nullable private volatile Status savedStatus;
1464
    private final AtomicBoolean terminationTriggered = new AtomicBoolean(false);
1✔
1465
    private final AtomicBoolean responseHeadersSent = new AtomicBoolean(false);
1✔
1466
    private final AtomicBoolean trailersOnly = new AtomicBoolean(false);
1✔
1467

1468
    protected DataPlaneListener(
1469
        ClientCall.Listener<InputStream> delegate,
1470
        DataPlaneClientCall dataPlaneClientCall) {
1471
      super(delegate);
1✔
1472
      this.dataPlaneClientCall = dataPlaneClientCall;
1✔
1473
    }
1✔
1474

1475
    boolean isTrailersOnly() {
1476
      return trailersOnly.get();
1✔
1477
    }
1478

1479
    Metadata getSavedHeaders() {
1480
      return savedHeaders;
1✔
1481
    }
1482

1483
    Metadata getSavedTrailers() {
1484
      return savedTrailers;
1✔
1485
    }
1486

1487
    void setImmediateResponse(Status status, Metadata trailers) {
1488
      this.savedStatus = status;
1✔
1489
      this.savedTrailers = trailers;
1✔
1490
    }
1✔
1491

1492
    @Override
1493
    public void onReady() {
1494
      dataPlaneClientCall.onReady();
1✔
1495
    }
1✔
1496

1497
    @Override
1498
    public void onHeaders(Metadata headers) {
1499
      dataPlaneClientCall.setServerHeadersStartNanos(System.nanoTime());
1✔
1500
      responseHeadersSent.set(true);
1✔
1501
      boolean sendResponseHeaders =
1✔
1502
          dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1503
              == ProcessingMode.HeaderSendMode.SEND
1504
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1505
              == ProcessingMode.HeaderSendMode.DEFAULT;
1506

1507
      if (dataPlaneClientCall.getExtProcStreamState().get().isDraining() && sendResponseHeaders) {
1✔
1508
        this.savedHeaders = headers;
1✔
1509
        return;
1✔
1510
      }
1511

1512
      if (dataPlaneClientCall.getPassThroughMode().get()
1✔
1513
          || dataPlaneClientCall.getExtProcStreamState().get().isCompleted() 
1✔
1514
          || !sendResponseHeaders) {
1515
        proceedWithHeaders(headers);
1✔
1516
        return;
1✔
1517
      }
1518

1519
      this.savedHeaders = headers;
1✔
1520
      dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1✔
1521
          .setResponseHeaders(HttpHeaders.newBuilder()
1✔
1522
              .setHeaders(
1✔
1523
                  toHeaderMap(headers, dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1✔
1524
              .build())
1✔
1525
          .build());
1✔
1526

1527
      if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1528
        proceedWithHeaders();
1✔
1529
      }
1530
    }
1✔
1531

1532
    @Override
1533
    public void onMessage(InputStream message) {
1534
      synchronized (dataPlaneClientCall.streamLock) {
1✔
1535
        if (inboundPassThrough) {
1✔
1536
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
1✔
1537
          return;
1✔
1538
        }
1539

1540
        boolean checkDrain = dataPlaneClientCall.getExtProcStreamState().get().isDraining()
1✔
1541
            && dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1✔
1542
                == ProcessingMode.BodySendMode.GRPC;
1543

1544
        if (savedHeaders != null || checkDrain) {
1✔
1545
          try {
1546
            ByteString copiedBody = ByteString.readFrom(message);
1✔
1547
            savedMessages.add(new KnownLengthInputStream(copiedBody));
1✔
1548
          } catch (IOException e) {
×
1549
            dataPlaneClientCall.cancelDownstream("Failed to copy inbound message for buffering", e);
×
1550
          }
1✔
1551
          return;
1✔
1552
        }
1553

1554
        if (dataPlaneClientCall.getPassThroughMode().get()) {
1✔
1555
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
×
1556
          return;
×
1557
        }
1558

1559
        if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1✔
1560
            || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1✔
1561
                != ProcessingMode.BodySendMode.GRPC) {
1562
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
1✔
1563
          return;
1✔
1564
        }
1565

1566
        try {
1567
          ByteString bodyByteString = ByteString.readFrom(message);
1✔
1568
          // TODO: Consider having separate classes handling normal mode and observability mode
1569
          if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1570
            sendResponseBodyToExtProc(bodyByteString, false);
×
1571
            dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
×
1572
            dataPlaneClientCall.getCallContext().run(
×
1573
                () -> delegate().onMessage(bodyByteString.newInput()));
×
1574
          } else {
1575
            if (dataPlaneClientCall.upstreamToSidestreamWindow <= 0 || !savedMessages.isEmpty()) {
1✔
1576
              savedMessages.add(new KnownLengthInputStream(bodyByteString));
1✔
1577
            } else {
1578
              dataPlaneClientCall.upstreamToSidestreamWindow -= bodyByteString.size();
1✔
1579
              sendResponseBodyToExtProc(bodyByteString, false);
1✔
1580
              dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
1✔
1581
            }
1582
            dataPlaneClientCall.drainPendingRequests();
1✔
1583
          }
1584
        } catch (IOException e) {
×
1585
          dataPlaneClientCall.cancelDownstream("Failed to read server response", e);
×
1586
        }
1✔
1587
      }
1✔
1588
    }
1✔
1589

1590
    void drainSavedMessages() {
1591
      synchronized (dataPlaneClientCall.streamLock) {
1✔
1592
        while (dataPlaneClientCall.isSidecarReady()
1✔
1593
            && dataPlaneClientCall.upstreamToSidestreamWindow > 0
1✔
1594
            && !savedMessages.isEmpty()) {
1✔
1595
          InputStream msg = savedMessages.poll();
1✔
1596
          if (msg != null) {
1✔
1597
            try {
1598
              ByteString bodyByteString = ByteString.readFrom(msg);
1✔
1599
              dataPlaneClientCall.upstreamToSidestreamWindow -= bodyByteString.size();
1✔
1600
              sendResponseBodyToExtProc(bodyByteString, false);
1✔
1601
              dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
1✔
1602
            } catch (IOException e) {
×
1603
              dataPlaneClientCall.cancelDownstream("Failed to read buffered response body", e);
×
1604
            }
1✔
1605
          }
1606
        }
1✔
1607
        dataPlaneClientCall.drainPendingRequests();
1✔
1608
      }
1✔
1609
    }
1✔
1610

1611
    @Override
1612
    public void onClose(Status status, Metadata trailers) {
1613
      dataPlaneClientCall.setServerTrailersStartNanos(System.nanoTime());
1✔
1614
      ExtProcStreamState extProcStreamState =
1✔
1615
          dataPlaneClientCall.getExtProcStreamState().get();
1✔
1616
      if (extProcStreamState.isFailed()
1✔
1617
          && !dataPlaneClientCall.getConfig().getObservabilityMode()
1✔
1618
          && (!dataPlaneClientCall.getConfig().getFailureModeAllow()
1✔
1619
              || dataPlaneClientCall.bodyMessageSentToExtProc.get())) {
1✔
1620
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1✔
1621
          proceedWithClose(Status.INTERNAL.withDescription("External processor stream failed")
1✔
1622
              .withCause(status.getCause()), new Metadata());
1✔
1623
        }
1624
        return;
1✔
1625
      }
1626
      if (dataPlaneClientCall.getPassThroughMode().get()) {
1✔
1627
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1✔
1628
          proceedWithClose(status, trailers);
1✔
1629
        }
1630
        return;
1✔
1631
      }
1632

1633
      if (this.savedStatus == null) {
1✔
1634
        this.savedStatus = status;
1✔
1635
        this.savedTrailers = trailers;
1✔
1636
      }
1637

1638
      // If we are still waiting for the external processor to validate response headers,
1639
      // buffer the close status/trailers and defer the close until headers are processed.
1640
      if (savedHeaders != null) {
1✔
1641
        return;
1✔
1642
      }
1643

1644
      boolean sendResponseTrailers =
1✔
1645
          dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1✔
1646
              == ProcessingMode.HeaderSendMode.SEND;
1647

1648
      if (dataPlaneClientCall.getExtProcStreamState().get().isDraining() && sendResponseTrailers) {
1✔
1649
        return;
×
1650
      }
1651

1652
      if (!responseHeadersSent.get()) {
1✔
1653
        trailersOnly.set(true);
1✔
1654
      }
1655

1656
      triggerCloseHandshake();
1✔
1657
    }
1✔
1658

1659
    void onReadyNotify() {
1660
      dataPlaneClientCall.getCallContext().run(() -> delegate().onReady());
1✔
1661
    }
1✔
1662

1663
    void proceedWithHeaders() {
1664
      if (savedHeaders != null) {
1✔
1665
        proceedWithHeaders(savedHeaders);
1✔
1666
        synchronized (dataPlaneClientCall.streamLock) {
1✔
1667
          savedHeaders = null;
1✔
1668
          if (!dataPlaneClientCall.getExtProcStreamState().get().isDraining()) {
1✔
1669
            InputStream msg;
1670
            while ((msg = savedMessages.poll()) != null) {
1✔
1671
              onMessage(msg);
1✔
1672
            }
1673
          }
1674
        }
1✔
1675
        onReadyNotify();
1✔
1676
        if (savedStatus != null) {
1✔
1677
          triggerCloseHandshake();
1✔
1678
        }
1679
      }
1680
    }
1✔
1681

1682
    private void proceedWithHeaders(Metadata headers) {
1683
      if (dataPlaneClientCall.getServerHeadersStartNanos() > 0) {
1✔
1684
        long durationNanos = System.nanoTime() - dataPlaneClientCall.getServerHeadersStartNanos();
1✔
1685
        dataPlaneClientCall.recordDuration(serverHeadersDuration, durationNanos);
1✔
1686
        dataPlaneClientCall.setServerHeadersStartNanos(0);
1✔
1687
      }
1688
      dataPlaneClientCall.getCallContext().run(() -> delegate().onHeaders(headers));
1✔
1689
    }
1✔
1690

1691
    void proceedWithClose() {
1692
      if (savedStatus != null) {
1✔
1693
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1✔
1694
          proceedWithClose(savedStatus, savedTrailers);
1✔
1695
        }
1696
        savedStatus = null;
1✔
1697
        savedTrailers = null;
1✔
1698
      }
1699
    }
1✔
1700

1701
    private void proceedWithClose(Status status, Metadata trailers) {
1702
      if (dataPlaneClientCall.getServerTrailersStartNanos() > 0) {
1✔
1703
        long durationNanos = System.nanoTime() - dataPlaneClientCall.getServerTrailersStartNanos();
1✔
1704
        dataPlaneClientCall.recordDuration(serverTrailersDuration, durationNanos);
1✔
1705
        dataPlaneClientCall.setServerTrailersStartNanos(0);
1✔
1706
      }
1707
      dataPlaneClientCall.getCallContext().run(() -> delegate().onClose(status, trailers));
1✔
1708
    }
1✔
1709

1710
    void onExternalBody(ByteString body) {
1711
      // If needed, downstream reading can be made more optimal by creating a wrapped
1712
      // Inputstream wraps the underlying bytestring and that implements HasByteBuffer,
1713
      // Detachable, KnownLength
1714
      dataPlaneClientCall.getCallContext().run(
1✔
1715
          () -> delegate().onMessage(body.newInput()));
1✔
1716
    }
1✔
1717

1718
    void unblockAfterStreamComplete() {
1719
      proceedWithHeaders();
1✔
1720
      // 1. Drain mutated responses first
1721
      dataPlaneClientCall.drainPendingMutatedResponseBodiesDirect(this);
1✔
1722
      // 2. Drain raw responses
1723
      proceedWithSavedMessages();
1✔
1724
      // 3. Drain outbound requests
1725
      dataPlaneClientCall.drainPendingDrainingMessages();
1✔
1726
      proceedWithClose();
1✔
1727
    }
1✔
1728

1729
    private void proceedWithSavedMessages() {
1730
      synchronized (dataPlaneClientCall.streamLock) {
1✔
1731
        InputStream msg;
1732
        while ((msg = savedMessages.poll()) != null) {
1✔
1733
          final InputStream finalMsg = msg;
×
1734
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(finalMsg));
×
1735
        }
×
1736
        inboundPassThrough = true;
1✔
1737
      }
1✔
1738
    }
1✔
1739

1740
    private void triggerCloseHandshake() {
1741
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1✔
1742
          || !terminationTriggered.compareAndSet(false, true)) {
1✔
1743
        return;
1✔
1744
      }
1745

1746
      boolean sendResponseHeaders =
1✔
1747
          dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1748
              == ProcessingMode.HeaderSendMode.SEND
1749
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1750
              == ProcessingMode.HeaderSendMode.DEFAULT;
1751

1752
      boolean sendResponseTrailers =
1✔
1753
          dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1✔
1754
              == ProcessingMode.HeaderSendMode.SEND;
1755

1756
      if (trailersOnly.get()) {
1✔
1757
        if (sendResponseHeaders) {
1✔
1758
          dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1✔
1759
              .setResponseHeaders(HttpHeaders.newBuilder()
1✔
1760
                  .setHeaders(
1✔
1761
                      toHeaderMap(
1✔
1762
                          savedTrailers,
1763
                          dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1✔
1764
                  .setEndOfStream(true)
1✔
1765
                  .build())
1✔
1766
              .build());
1✔
1767
        } else {
1768
          proceedWithClose();
1✔
1769
          if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1770
            dataPlaneClientCall.closeExtProcStream();
1✔
1771
          }
1772
        }
1773
      } else if (sendResponseTrailers) {
1✔
1774
        dataPlaneClientCall.getIsProcessingTrailers().set(true);
1✔
1775
        dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1✔
1776
            .setResponseTrailers(HttpTrailers.newBuilder()
1✔
1777
                .setTrailers(
1✔
1778
                    toHeaderMap(
1✔
1779
                        savedTrailers,
1780
                        dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1✔
1781
                .build())
1✔
1782
            .build());
1✔
1783
      } else {
1784
        proceedWithClose();
1✔
1785
        if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1786
          dataPlaneClientCall.closeExtProcStream();
1✔
1787
        }
1788
      }
1789

1790
      if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1791
        proceedWithClose();
1✔
1792
        @SuppressWarnings("unused")
1793
        ScheduledFuture<?> unused = dataPlaneClientCall.getScheduler().schedule(
1✔
1794
            dataPlaneClientCall::closeExtProcStream,
1✔
1795
            dataPlaneClientCall.getConfig().getDeferredCloseTimeoutNanos(),
1✔
1796
            TimeUnit.NANOSECONDS);
1797
      }
1798
    }
1✔
1799

1800
    @GuardedBy("dataPlaneClientCall.streamLock")
1801
    private void sendResponseBodyToExtProc(
1802
        @Nullable ByteString bodyByteString, boolean endOfStream) {
1803
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1✔
1804
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1✔
1805
              != ProcessingMode.BodySendMode.GRPC) {
1806
        return;
×
1807
      }
1808

1809
      HttpBody.Builder bodyBuilder =
1810
          HttpBody.newBuilder();
1✔
1811
      if (bodyByteString != null) {
1✔
1812
        bodyBuilder.setBody(bodyByteString);
1✔
1813
      }
1814
      bodyBuilder.setEndOfStream(endOfStream);
1✔
1815

1816
      ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1✔
1817
          .setResponseBody(bodyBuilder.build());
1✔
1818
      dataPlaneClientCall.mergeAccumulatedWindowUpdates(builder);
1✔
1819
      dataPlaneClientCall.sendToExtProc(builder.build());
1✔
1820
    }
1✔
1821
  }
1822
}
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