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

grpc / grpc-java / #20444

07 Sep 2026 06:51PM UTC coverage: 89.25% (+0.04%) from 89.211%
#20444

push

github

web-flow
feat(xds): Implement  CheckRequestBuilder for external authorization  (#12493)

This PR sits on top of #12492, so only the last commit + any fixups need
to be reviewed.

This commit introduces the `CheckRequestBuilder` library, which is
responsible for constructing the `CheckRequest` message sent to the
external authorization service.

The `CheckRequestBuilder` gathers information from various sources,
including:
- `ServerCall` attributes (local and remote addresses, SSL session).
- `MethodDescriptor` (full method name).
- Request headers.

It uses this information to populate the `AttributeContext` of the
`CheckRequest` message, which provides the authorization service with
the necessary context to make an authorization decision.

This commit also introduces the `ExtAuthzCertificateProvider`, a helper
class for extracting certificate information, such as the principal and
PEM-encoded certificate.

The relevant section of the spec is:
https://github.com/grpc/proposal/pull/481/files#diff-6bb76a24ad2fd8849f164244e68cd54eaR196-R250

Unit tests for the new components are also included.

38780 of 43451 relevant lines covered (89.25%)

0.89 hits per line

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

95.56
/../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("grpc.lb.backend_service"),
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("grpc.lb.backend_service"),
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("grpc.lb.backend_service"),
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("grpc.lb.backend_service"),
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
        (ClientCall<InputStream, InputStream>) (ClientCall<?, ?>)
240
            next.newCall(method, callOptions);
1✔
241

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

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

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

255
  // --- SHARED UTILITY METHODS ---
256

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

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

274
    private final ExternalProcessorGrpc.ExternalProcessorStub stub;
275
    private final ExternalProcessorFilterConfig config;
276
    private final ClientCall<InputStream, InputStream> rawCall;
277
    private final DataPlaneDelayedCall<InputStream, InputStream> delayedCall;
278
    private final ScheduledExecutorService scheduler;
279
    final Object streamLock = new Object();
1✔
280
    @Nullable private volatile EventType expectedRequestResponse;
281
    @Nullable private volatile EventType expectedResponseResponse;
282
    @Nullable private volatile ClientCallStreamObserver<ProcessingRequest>
283
        extProcClientCallRequestObserver;
284
    @GuardedBy("streamLock")
1✔
285
    private final Queue<InputStream> pendingDrainingMessages =
286
        new ConcurrentLinkedQueue<>();
287
    @Nullable private volatile DataPlaneListener wrappedListener;
288
    private final HeaderMutationFilter mutationFilter;
289
    private final HeaderMutator mutator = HeaderMutator.create();
1✔
290
    private final AtomicInteger pendingRequests = new AtomicInteger(0);
1✔
291
    private final ProcessingMode currentProcessingMode;
292

293
    // Default initial window size
294
    private static final long DEFAULT_INITIAL_WINDOW_SIZE = 65536;
295

296
    // Outbound (sending) windows
297
    @GuardedBy("streamLock")
1✔
298
    private long downstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
299
    @GuardedBy("streamLock")
1✔
300
    private long upstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
301

302
    // Inbound (receiving) windows
303
    @GuardedBy("streamLock")
1✔
304
    private long sidestreamToUpstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
305
    @GuardedBy("streamLock")
1✔
306
    private long sidestreamToDownstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
307

308
    // Threshold to trigger standalone client window updates
309
    private static final long WINDOW_UPDATE_THRESHOLD = DEFAULT_INITIAL_WINDOW_SIZE / 2;
310

311
    // Path 1: Pending/buffered request body messages from downstream
312
    @GuardedBy("streamLock")
1✔
313
    private final Queue<ByteString> pendingRequestBodyMessages = new ConcurrentLinkedQueue<>();
314
    // Deferred half-close flag for upstream direction
315
    private final AtomicBoolean pendingUpstreamHalfClose = new AtomicBoolean(false);
1✔
316

317
    // Path 2: Buffered request body messages from ext_proc server to forward upstream
318
    @GuardedBy("streamLock")
1✔
319
    private final Queue<ByteString> pendingUpstreamBodyMessages =
320
        new ConcurrentLinkedQueue<>();
321
    // Path 4: Outstanding requests from downstream for pulling responses
322
    @GuardedBy("streamLock")
1✔
323
    private int downstreamRequestsPending = 0;
324
    // Buffered mutated response bodies from ext_proc server
325
    @GuardedBy("streamLock")
1✔
326
    private final Queue<ByteString> pendingMutatedResponseBodies =
327
        new ConcurrentLinkedQueue<>();
328

329
    // Accumulated client window updates to send to ext_proc
330
    @GuardedBy("streamLock")
1✔
331
    private long accumulatedWindowUpdateSidestreamToUpstream = 0;
332
    @GuardedBy("streamLock")
1✔
333
    private long accumulatedWindowUpdateSidestreamToDownstream = 0;
334

335
    // Flag to track if FlowControlInit was sent in the initial message
336
    @GuardedBy("streamLock")
1✔
337
    private boolean flowControlInitSent = false;
338

339
    private final MethodDescriptor<?, ?> method;
340
    private final Channel channel;
341
    private final MetricRecorder metricsRecorder;
342
    private final String target;
343
    private final String backendService;
344
    private volatile Context callContext = Context.ROOT;
1✔
345

346
    private volatile long clientHeadersStartNanos;
347
    private volatile long clientHalfCloseStartNanos;
348
    private volatile long serverHeadersStartNanos;
349
    private volatile long serverTrailersStartNanos;
350

351
    private boolean protocolConfigSent = false;
1✔
352
    private ImmutableMap<String, Struct> collectedAttributes;
353
    private boolean requestAttributesSent = false;
1✔
354
    @Nullable private volatile Metadata requestHeaders;
355
    final AtomicReference<DataPlaneCallState> dataPlaneCallState =
1✔
356
        new AtomicReference<>(DataPlaneCallState.IDLE);
357
    final AtomicReference<ExtProcStreamState> extProcStreamState =
1✔
358
        new AtomicReference<>(ExtProcStreamState.ACTIVE);
359
    final AtomicBoolean passThroughMode = new AtomicBoolean(false);
1✔
360
    final AtomicBoolean requestSideClosed = new AtomicBoolean(false);
1✔
361
    final AtomicBoolean appHalfClosed = new AtomicBoolean(false);
1✔
362
    final AtomicBoolean isProcessingTrailers = new AtomicBoolean(false);
1✔
363
    final AtomicBoolean pendingHalfClose = new AtomicBoolean(false);
1✔
364
    final AtomicBoolean bodyMessageSentToExtProc = new AtomicBoolean(false);
1✔
365
    private final AtomicBoolean downstreamCancelled = new AtomicBoolean(false);
1✔
366

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

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

416
    private void recordDuration(DoubleHistogramMetricInstrument instrument, long durationNanos) {
417
      if (instrument != null) {
1✔
418
        double durationSecs = (double) durationNanos / 1_000_000_000.0;
1✔
419
        metricsRecorder.recordDoubleHistogram(
1✔
420
            instrument,
421
            durationSecs,
422
            ImmutableList.of(target),
1✔
423
            ImmutableList.of(backendService));
1✔
424
      }
425
    }
1✔
426

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

461
    @Override
462
    public void start(Listener<InputStream> responseListener, Metadata headers) {
463
      this.callContext = Context.current();
1✔
464
      clientHeadersStartNanos = System.nanoTime();
1✔
465
      this.requestHeaders = headers;
1✔
466
      this.wrappedListener = new DataPlaneListener(responseListener, rawCall, this);
1✔
467

468
      // DelayedClientCall.start will buffer the listener and headers until setCall is called.
469
      super.start(wrappedListener, headers);
1✔
470

471
      stub.process(new ClientResponseObserver<ProcessingRequest, ProcessingResponse>() {
1✔
472
        @Override
473
        public void beforeStart(ClientCallStreamObserver<ProcessingRequest> requestStream) {
474
          synchronized (streamLock) {
1✔
475
            extProcClientCallRequestObserver = requestStream;
1✔
476
          }
1✔
477
          requestStream.setOnReadyHandler(DataPlaneClientCall.this::onExtProcStreamReady);
1✔
478
        }
1✔
479

480
        @Override
481
        public void onNext(ProcessingResponse response) {
482
          try {
483
            if (config.getObservabilityMode()) {
1✔
484
              return;
1✔
485
            }
486

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

507
            if (response.hasImmediateResponse()) {
1✔
508
              if (config.getDisableImmediateResponse()) {
1✔
509
                internalOnError(Status.UNAVAILABLE
1✔
510
                    .withDescription(
1✔
511
                        "Immediate response is disabled but received from external processor")
512
                    .asRuntimeException());
1✔
513
                return;
1✔
514
              }
515
              handleImmediateResponse(response.getImmediateResponse(), wrappedListener);
1✔
516
              return;
1✔
517
            }
518

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

569
            if (response.getRequestDrain()) {
1✔
570
              extProcStreamState.set(ExtProcStreamState.DRAINING);
1✔
571
              halfCloseExtProcStream();
1✔
572
              activateCall();
1✔
573
            }
574

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

641
            checkEndOfStream(response);
1✔
642
          } catch (Throwable t) {
1✔
643
            internalOnError(t);
1✔
644
          }
1✔
645
        }
1✔
646

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

664
        @Override
665
        public void onCompleted() {
666
          if (markExtProcStreamCompleted(extProcStreamState)) {
1✔
667
            handleFailOpen(wrappedListener);
1✔
668
          }
669
        }
1✔
670
      });
671

672
      this.collectedAttributes = collectAttributes(
1✔
673
          config.getRequestAttributes(), method, channel.authority(), headers);
1✔
674

675
      boolean sendRequestHeaders =
1✔
676
          currentProcessingMode.getRequestHeaderMode() == ProcessingMode.HeaderSendMode.SEND
1✔
677
          || currentProcessingMode.getRequestHeaderMode()
1✔
678
              == ProcessingMode.HeaderSendMode.DEFAULT;
679

680
      if (sendRequestHeaders) {
1✔
681
        sendToExtProc(ProcessingRequest.newBuilder()
1✔
682
            .setRequestHeaders(HttpHeaders.newBuilder()
1✔
683
                .setHeaders(toHeaderMap(headers, config.getForwardRulesConfig()))
1✔
684
                .setEndOfStream(false)
1✔
685
                .build())
1✔
686
            .build());
1✔
687
      }
688

689
      if (config.getObservabilityMode() || !sendRequestHeaders) {
1✔
690
        activateCall();
1✔
691
      }
692
    }
1✔
693

694
    private void sendToExtProc(ProcessingRequest request) {
695
      synchronized (streamLock) {
1✔
696
        if (extProcStreamState.get().isCompleted()) {
1✔
697
          return;
1✔
698
        }
699
        
700
        if (request.hasRequestHeaders()) {
1✔
701
          expectedRequestResponse = EventType.REQUEST_HEADERS;
1✔
702
        } else if (request.hasResponseHeaders()) {
1✔
703
          expectedResponseResponse = EventType.RESPONSE_HEADERS;
1✔
704
        } else if (request.hasResponseTrailers()) {
1✔
705
          expectedResponseResponse = EventType.RESPONSE_TRAILERS;
1✔
706
        }
707

708
        ProcessingRequest requestToSend = request;
1✔
709
        if (!protocolConfigSent) {
1✔
710
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
711
              .setProtocolConfig(ProtocolConfiguration.newBuilder()
1✔
712
                  .setRequestBodyMode(currentProcessingMode.getRequestBodyMode())
1✔
713
                  .setResponseBodyMode(currentProcessingMode.getResponseBodyMode())
1✔
714
                  .build())
1✔
715
              .build();
1✔
716
          protocolConfigSent = true;
1✔
717
        }
718

719
        boolean isClientServerMessage =
1✔
720
            requestToSend.hasRequestHeaders() || requestToSend.hasRequestBody();
1✔
721
        if (isClientServerMessage
1✔
722
            && !requestAttributesSent
723
            && collectedAttributes != null
724
            && !collectedAttributes.isEmpty()) {
1✔
725
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
726
              .putAllAttributes(collectedAttributes)
1✔
727
              .build();
1✔
728
          requestAttributesSent = true;
1✔
729
        }
730

731
        if (config.getObservabilityMode()) {
1✔
732
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
733
              .setObservabilityMode(true)
1✔
734
              .build();
1✔
735
        } else if (!flowControlInitSent) {
1✔
736
          requestToSend = ProcessingRequest.newBuilder(requestToSend)
1✔
737
              .setFlowControlInit(ProcessingRequest.FlowControlInit.newBuilder()
1✔
738
                  .setInitialWindowDownstreamToSidestream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
739
                  .setInitialWindowSidestreamToUpstream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
740
                  .setInitialWindowUpstreamToSidestream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
741
                  .setInitialWindowSidestreamToDownstream(DEFAULT_INITIAL_WINDOW_SIZE)
1✔
742
                  .build())
1✔
743
              .build();
1✔
744
          flowControlInitSent = true;
1✔
745
        }
746

747
        extProcClientCallRequestObserver.onNext(requestToSend);
1✔
748
      }
1✔
749
    }
1✔
750

751
    // Note: This method not only modifies the builder, but has the side effect of modifying
752
    // the window update bookkeeping.
753
    @GuardedBy("streamLock")
754
    void mergeAccumulatedWindowUpdates(ProcessingRequest.Builder requestBuilder) {
755
      long incrementUpstream = accumulatedWindowUpdateSidestreamToUpstream;
1✔
756
      long incrementDownstream = accumulatedWindowUpdateSidestreamToDownstream;
1✔
757

758
      if (incrementUpstream > 0 || incrementDownstream > 0) {
1✔
759
        requestBuilder.setClientWindowUpdate(
1✔
760
            ProcessingRequest.ClientWindowUpdate.newBuilder()
1✔
761
                .setWindowIncrementSidestreamToUpstream(incrementUpstream)
1✔
762
                .setWindowIncrementSidestreamToDownstream(incrementDownstream)
1✔
763
                .build());
1✔
764
        accumulatedWindowUpdateSidestreamToUpstream -= incrementUpstream;
1✔
765
        accumulatedWindowUpdateSidestreamToDownstream -= incrementDownstream;
1✔
766
        sidestreamToUpstreamWindow += incrementUpstream;
1✔
767
        sidestreamToDownstreamWindow += incrementDownstream;
1✔
768
      }
769
    }
1✔
770

771
    private void trySendAccumulatedWindowUpdates() {
772
      synchronized (streamLock) {
1✔
773
        if (extProcStreamState.get().isCompleted()) {
1✔
774
          return;
1✔
775
        }
776
        long incrementUpstream = accumulatedWindowUpdateSidestreamToUpstream;
1✔
777
        long incrementDownstream = accumulatedWindowUpdateSidestreamToDownstream;
1✔
778

779
        boolean shouldSend = (incrementUpstream > 0 || incrementDownstream > 0) && (
1✔
780
            (incrementUpstream >= WINDOW_UPDATE_THRESHOLD)
781
            || (incrementDownstream >= WINDOW_UPDATE_THRESHOLD)
782
            || (sidestreamToUpstreamWindow <= 0 && accumulatedWindowUpdateSidestreamToUpstream > 0)
783
            || (sidestreamToDownstreamWindow <= 0
784
                && accumulatedWindowUpdateSidestreamToDownstream > 0)
785
        );
786

787
        if (shouldSend) {
1✔
788
          accumulatedWindowUpdateSidestreamToUpstream -= incrementUpstream;
1✔
789
          accumulatedWindowUpdateSidestreamToDownstream -= incrementDownstream;
1✔
790
          sidestreamToUpstreamWindow += incrementUpstream;
1✔
791
          sidestreamToDownstreamWindow += incrementDownstream;
1✔
792

793
          sendToExtProc(ProcessingRequest.newBuilder()
1✔
794
              .setClientWindowUpdate(ProcessingRequest.ClientWindowUpdate.newBuilder()
1✔
795
                  .setWindowIncrementSidestreamToUpstream(incrementUpstream)
1✔
796
                  .setWindowIncrementSidestreamToDownstream(incrementDownstream)
1✔
797
                  .build())
1✔
798
              .build());
1✔
799
        }
800
      }
1✔
801
    }
1✔
802

803
    private void onExtProcStreamReady() {
804
      drainPendingRequests();
1✔
805
      onReadyNotify();
1✔
806
    }
1✔
807

808
    void drainPendingRequests() {
809
      synchronized (streamLock) {
1✔
810
        if (config.getObservabilityMode()
1✔
811
            || currentProcessingMode.getResponseBodyMode() != ProcessingMode.BodySendMode.GRPC
1✔
812
            || extProcStreamState.get().isCompleted()) {
1✔
813
          int toRequest = pendingRequests.getAndSet(0);
1✔
814
          if (toRequest > 0) {
1✔
815
            super.request(toRequest);
1✔
816
          }
817
          return;
1✔
818
        }
819

820
        // Normal mode flow control: pull 1 message at a time
821
        if (isSidecarReady() && upstreamToSidestreamWindow > 0 && pendingRequests.get() > 0) {
1✔
822
          super.request(1);
1✔
823
          pendingRequests.decrementAndGet();
1✔
824
        }
825
      }
1✔
826
    }
1✔
827

828
    private void closeExtProcStream() {
829
      synchronized (streamLock) {
1✔
830
        if (markExtProcStreamCompleted(extProcStreamState)) {
1✔
831
          if (extProcClientCallRequestObserver != null) {
1✔
832
            extProcClientCallRequestObserver.onCompleted();
1✔
833
          }
834
        }
835
      }
1✔
836
    }
1✔
837

838
    private void internalOnError(Throwable t) {
839
      if (markExtProcStreamFailed(extProcStreamState)) {
1✔
840
        synchronized (streamLock) {
1✔
841
          if (extProcClientCallRequestObserver != null) {
1✔
842
            try {
843
              extProcClientCallRequestObserver.onError(t);
1✔
844
            } catch (Throwable ignored) {
×
845
              // Ignore exceptions during cancel/onError propagation
846
            }
1✔
847
            extProcClientCallRequestObserver = null;
1✔
848
          }
849
        }
1✔
850
        if (config.getObservabilityMode()
1✔
851
            || (config.getFailureModeAllow() && !bodyMessageSentToExtProc.get())) {
1✔
852
          handleFailOpen(wrappedListener);
1✔
853
        } else {
854
          String message = "External processor stream failed";
1✔
855
          cancelDownstream(message, t);
1✔
856
          wrappedListener.proceedWithClose();
1✔
857
        }
858
      }
859
    }
1✔
860

861
    private void halfCloseExtProcStream() {
862
      synchronized (streamLock) {
1✔
863
        if (!extProcStreamState.get().isCompleted() && extProcClientCallRequestObserver != null) {
1✔
864
          extProcClientCallRequestObserver.onCompleted();
1✔
865
        }
866
      }
1✔
867
    }
1✔
868

869
    private void onReadyNotify() {
870
      wrappedListener.onReadyNotify();
1✔
871
    }
1✔
872

873
    void onReady() {
874
      boolean isPassThrough;
875
      boolean isCompleted;
876
      boolean isDraining;
877

878
      synchronized (streamLock) {
1✔
879
        isPassThrough = passThroughMode.get();
1✔
880
        ExtProcStreamState state = extProcStreamState.get();
1✔
881
        isCompleted = state.isCompleted();
1✔
882
        isDraining = state.isDraining();
1✔
883
      }
1✔
884

885
      if (isPassThrough) {
1✔
886
        onReadyNotify();
×
887
        return;
×
888
      }
889

890
      if (isCompleted) {
1✔
891
        drainPendingDrainingMessages();
1✔
892
        return;
1✔
893
      }
894

895
      // Normal or Draining operation
896
      drainPendingUpstreamBodyMessages();
1✔
897
      if (!isDraining) {
1✔
898
        trySendAccumulatedWindowUpdates();
1✔
899
      }
900
      drainPendingRequests();
1✔
901
      onReadyNotify();
1✔
902
    }
1✔
903

904
    @GuardedBy("streamLock")
905
    boolean isSidecarReady() {
906
      ExtProcStreamState state = extProcStreamState.get();
1✔
907
      if (state.isCompleted()) {
1✔
908
        return true;
×
909
      }
910
      if (state.isDraining()) {
1✔
911
        return false;
1✔
912
      }
913
      ClientCallStreamObserver<ProcessingRequest> observer = extProcClientCallRequestObserver;
1✔
914
      return observer != null && observer.isReady();
1✔
915
    }
916

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

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

956
        if (normalFlowControl) {
1✔
957
          pendingRequests.addAndGet(numMessages);
1✔
958
          downstreamRequestsPending += numMessages;
1✔
959
          drainPendingMutatedResponseBodies();
1✔
960
          if (isSidecarReady()) {
1✔
961
            drainPendingRequests();
1✔
962
          }
963
        } else {
964
          // Observability mode: gate on readiness but pull all at once
965
          if (isSidecarReady()) {
1✔
966
            super.request(numMessages);
1✔
967
          } else {
968
            pendingRequests.addAndGet(numMessages);
1✔
969
          }
970
        }
971
      }
1✔
972
    }
1✔
973

974
    @Override
975
    public void sendMessage(InputStream message) {
976
      if (requestSideClosed.get()) {
1✔
977
        // External processor already closed the request stream. Discard further messages.
978
        return;
1✔
979
      }
980

981
      if (passThroughMode.get()) {
1✔
982
        super.sendMessage(message);
1✔
983
        return;
1✔
984
      }
985

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

992
        ExtProcStreamState state = extProcStreamState.get();
1✔
993
        if (state.isDraining() || state.isCompleted()) {
1✔
994
          if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1✔
995
            super.sendMessage(message);
1✔
996
            return;
1✔
997
          }
998
          try {
999
            ByteString copiedBody = ByteString.readFrom(message);
1✔
1000
            pendingDrainingMessages.add(new KnownLengthInputStream(copiedBody));
1✔
1001
          } catch (IOException e) {
×
1002
            rawCall.cancel("Failed to copy outbound message for buffering", e);
×
1003
          }
1✔
1004
          return;
1✔
1005
        }
1006

1007
        if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1✔
1008
          super.sendMessage(message);
1✔
1009
          return;
1✔
1010
        }
1011

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

1037
    @GuardedBy("streamLock")
1038
    private void sendRequestBodyToExtProc(ByteString body) {
1039
      downstreamToSidestreamWindow -= body.size();
1✔
1040
      ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1✔
1041
          .setRequestBody(HttpBody.newBuilder()
1✔
1042
              .setBody(body)
1✔
1043
              .setEndOfStream(false)
1✔
1044
              .build());
1✔
1045
      mergeAccumulatedWindowUpdates(builder);
1✔
1046
      sendToExtProc(builder.build());
1✔
1047
      bodyMessageSentToExtProc.set(true);
1✔
1048
    }
1✔
1049

1050
    @GuardedBy("streamLock")
1051
    private void drainPendingRequestBodyMessages() {
1052
      while (downstreamToSidestreamWindow > 0 && !pendingRequestBodyMessages.isEmpty()) {
1✔
1053
        ByteString body = pendingRequestBodyMessages.poll();
1✔
1054
        sendRequestBodyToExtProc(body);
1✔
1055
      }
1✔
1056
      if (pendingRequestBodyMessages.isEmpty() && pendingHalfClose.compareAndSet(true, false)) {
1✔
1057
        halfClose();
1✔
1058
      }
1059
    }
1✔
1060

1061
    private void proceedWithHalfClose() {
1062
      if (clientHalfCloseStartNanos > 0) {
1✔
1063
        long durationNanos = System.nanoTime() - clientHalfCloseStartNanos;
1✔
1064
        recordDuration(clientHalfCloseDuration, durationNanos);
1✔
1065
        clientHalfCloseStartNanos = 0;
1✔
1066
      }
1067
      super.halfClose();
1✔
1068
    }
1✔
1069

1070
    @Override
1071
    public void halfClose() {
1072
      if (appHalfClosed.compareAndSet(false, true)) {
1✔
1073
        clientHalfCloseStartNanos = System.nanoTime();
1✔
1074
      }
1075
      if (passThroughMode.get()) {
1✔
1076
        if (requestSideClosed.compareAndSet(false, true)) {
1✔
1077
          proceedWithHalfClose();
1✔
1078
        }
1079
        return;
1✔
1080
      }
1081

1082
      if (extProcStreamState.get().isCompleted()) {
1✔
1083
        if (passThroughMode.get()) {
1✔
1084
          if (requestSideClosed.compareAndSet(false, true)) {
×
1085
            proceedWithHalfClose();
×
1086
          }
1087
        } else {
1088
          pendingHalfClose.set(true);
1✔
1089
        }
1090
        return;
1✔
1091
      }
1092

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

1111
      if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
1✔
1112
        if (requestSideClosed.compareAndSet(false, true)) {
1✔
1113
          proceedWithHalfClose();
1✔
1114
        }
1115
        return;
1✔
1116
      }
1117

1118
      // Mode is GRPC
1119
      synchronized (streamLock) {
1✔
1120
        if (!pendingRequestBodyMessages.isEmpty()) {
1✔
1121
          pendingHalfClose.set(true);
1✔
1122
          return;
1✔
1123
        }
1124

1125
        ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1✔
1126
            .setRequestBody(HttpBody.newBuilder()
1✔
1127
                .setEndOfStreamWithoutMessage(true)
1✔
1128
                .build());
1✔
1129
        mergeAccumulatedWindowUpdates(builder);
1✔
1130
        sendToExtProc(builder.build());
1✔
1131
      }
1✔
1132
    }
1✔
1133

1134
    private void cancelDownstream(@Nullable String message, @Nullable Throwable cause) {
1135
      if (downstreamCancelled.compareAndSet(false, true)) {
1✔
1136
        delayedCall.cancel(message, cause);
1✔
1137
      }
1138
    }
1✔
1139

1140
    @Override
1141
    public void cancel(@Nullable String message, @Nullable Throwable cause) {
1142
      synchronized (streamLock) {
1✔
1143
        if (!extProcStreamState.get().isCompleted() && extProcClientCallRequestObserver != null) {
1✔
1144
          extProcClientCallRequestObserver.onError(
1✔
1145
              Status.CANCELLED
1146
                  .withDescription(message)
1✔
1147
                  .withCause(cause)
1✔
1148
                  .asRuntimeException());
1✔
1149
        }
1150
      }
1✔
1151
      cancelDownstream(message, cause);
1✔
1152
    }
1✔
1153

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

1191
    private void handleResponseBodyResponse(
1192
        BodyResponse bodyResponse, DataPlaneListener listener) {
1193
      if (bodyResponse.hasResponse() && bodyResponse.getResponse().hasBodyMutation()) {
1✔
1194
        BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
1✔
1195
        if (mutation.hasStreamedResponse()) {
1✔
1196
          StreamedBodyResponse streamed = mutation.getStreamedResponse();
1✔
1197
          ByteString body = streamed.getBody();
1✔
1198
          final int bodySize = body.size();
1✔
1199
          synchronized (streamLock) {
1✔
1200
            sidestreamToDownstreamWindow -= bodySize;
1✔
1201
          }
1✔
1202
          deliverResponseBody(body, listener);
1✔
1203
        }
1204
      }
1205
    }
1✔
1206

1207
    private void deliverResponseBody(ByteString body, DataPlaneListener listener) {
1208
      boolean shouldDeliver = false;
1✔
1209
      synchronized (streamLock) {
1✔
1210
        if (downstreamRequestsPending > 0) {
1✔
1211
          downstreamRequestsPending--;
1✔
1212
          shouldDeliver = true;
1✔
1213
        } else {
1214
          pendingMutatedResponseBodies.add(body);
1✔
1215
        }
1216
      }
1✔
1217
      if (shouldDeliver) {
1✔
1218
        final int bodySize = body.size();
1✔
1219
        callContext.run(() -> {
1✔
1220
          try {
1221
            listener.onExternalBody(body);
1✔
1222
          } finally {
1223
            synchronized (streamLock) {
1✔
1224
              accumulatedWindowUpdateSidestreamToDownstream += bodySize;
1✔
1225
            }
1✔
1226
            trySendAccumulatedWindowUpdates();
1✔
1227
          }
1228
        });
1✔
1229
      }
1230
    }
1✔
1231

1232
    private void drainPendingMutatedResponseBodies() {
1233
      List<ByteString> toDeliver = new ArrayList<>();
1✔
1234
      synchronized (streamLock) {
1✔
1235
        while (downstreamRequestsPending > 0 && !pendingMutatedResponseBodies.isEmpty()) {
1✔
1236
          ByteString body = pendingMutatedResponseBodies.poll();
1✔
1237
          downstreamRequestsPending--;
1✔
1238
          pendingRequests.decrementAndGet();
1✔
1239
          toDeliver.add(body);
1✔
1240
        }
1✔
1241
      }
1✔
1242
      for (ByteString body : toDeliver) {
1✔
1243
        final int bodySize = body.size();
1✔
1244
        callContext.run(() -> {
1✔
1245
          try {
1246
            wrappedListener.onExternalBody(body);
1✔
1247
          } finally {
1248
            synchronized (streamLock) {
1✔
1249
              accumulatedWindowUpdateSidestreamToDownstream += bodySize;
1✔
1250
            }
1✔
1251
            trySendAccumulatedWindowUpdates();
1✔
1252
          }
1253
        });
1✔
1254
      }
1✔
1255
    }
1✔
1256

1257
    // Used to immediately flush any mutated response chunks that we already received and buffered
1258
    // before the stream failed, ensuring the application receives them in the correct order
1259
    void drainPendingMutatedResponseBodiesDirect(DataPlaneListener listener) {
1260
      List<ByteString> toDeliver = new ArrayList<>();
1✔
1261
      synchronized (streamLock) {
1✔
1262
        ByteString body;
1263
        while ((body = pendingMutatedResponseBodies.poll()) != null) {
1✔
1264
          toDeliver.add(body);
×
1265
        }
1266
      }
1✔
1267
      for (ByteString body : toDeliver) {
1✔
1268
        listener.onExternalBody(body);
×
1269
      }
×
1270
    }
1✔
1271

1272
    void drainPendingUpstreamBodyMessages() {
1273
      while (true) {
1274
        ByteString body = null;
1✔
1275
        boolean triggerHalfClose = false;
1✔
1276
        synchronized (streamLock) {
1✔
1277
          if (!pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
1✔
1278
            body = pendingUpstreamBodyMessages.poll();
1✔
1279
            accumulatedWindowUpdateSidestreamToUpstream += body.size();
1✔
1280
            if (pendingUpstreamBodyMessages.isEmpty()
1✔
1281
                && pendingUpstreamHalfClose.compareAndSet(true, false)) {
1✔
1282
              triggerHalfClose = true;
1✔
1283
            }
1284
          }
1285
        }
1✔
1286
        if (body == null) {
1✔
1287
          break;
1✔
1288
        }
1289
        super.sendMessage(new KnownLengthInputStream(body));
1✔
1290
        trySendAccumulatedWindowUpdates();
1✔
1291
        if (triggerHalfClose) {
1✔
1292
          if (requestSideClosed.compareAndSet(false, true)) {
1✔
1293
            proceedWithHalfClose();
1✔
1294
          }
1295
        }
1296
      }
1✔
1297
    }
1✔
1298

1299
    private void handleImmediateResponse(ImmediateResponse immediate, DataPlaneListener listener)
1300
        throws HeaderMutationDisallowedException {
1301
      Status status = Status.fromCodeValue(immediate.getGrpcStatus().getStatus());
1✔
1302
      if (!immediate.getDetails().isEmpty()) {
1✔
1303
        status = status.withDescription(immediate.getDetails());
1✔
1304
      }
1305

1306
      Metadata trailers = new Metadata();
1✔
1307
      if (immediate.hasHeaders()) {
1✔
1308
        applyHeaderMutations(trailers, immediate.getHeaders(), mutationFilter, mutator);
1✔
1309
      }
1310

1311
      listener.setImmediateResponse(status, trailers);
1✔
1312

1313
      if (isProcessingTrailers.get()) {
1✔
1314
        // If sent in response to a server trailers event, sets the status and optionally
1315
        // headers to be included in the trailers.
1316
        listener.unblockAfterStreamComplete();
1✔
1317
      } else {
1318
        // If sent in response to any other event, it will cause the data plane RPC to
1319
        // immediately fail with the specified status as if it were an out-of-band
1320
        // cancellation.
1321
        rawCall.cancel(status.getDescription(), null);
1✔
1322
        listener.unblockAfterStreamComplete();
1✔
1323
      }
1324
      closeExtProcStream();
1✔
1325
    }
1✔
1326

1327
    private void drainPendingDrainingMessages() {
1328
      while (true) {
1329
        Object msg = null; // Can be ByteString or InputStream
1✔
1330
        boolean isMutated = false;
1✔
1331
        boolean triggerHalfClose = false;
1✔
1332

1333
        synchronized (streamLock) {
1✔
1334
          if (!pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
1✔
1335
            msg = pendingUpstreamBodyMessages.poll();
1✔
1336
            isMutated = true;
1✔
1337
          } else if (pendingUpstreamBodyMessages.isEmpty()
1✔
1338
              && !pendingRequestBodyMessages.isEmpty() && super.isReady()) {
1✔
1339
            msg = pendingRequestBodyMessages.poll();
1✔
1340
            isMutated = true;
1✔
1341
          } else if (pendingUpstreamBodyMessages.isEmpty()
1✔
1342
              && pendingRequestBodyMessages.isEmpty()
1✔
1343
              && !pendingDrainingMessages.isEmpty() && super.isReady()) {
1✔
1344
            msg = pendingDrainingMessages.poll();
1✔
1345
            isMutated = false;
1✔
1346
          }
1347

1348
          if (msg == null) {
1✔
1349
            if (pendingUpstreamBodyMessages.isEmpty()
1✔
1350
                && pendingRequestBodyMessages.isEmpty()
1✔
1351
                && pendingDrainingMessages.isEmpty()) {
1✔
1352
              passThroughMode.set(true);
1✔
1353
              if (appHalfClosed.get()) {
1✔
1354
                triggerHalfClose = true;
1✔
1355
              }
1356
            }
1357
          }
1358
        }
1✔
1359

1360
        if (msg == null) {
1✔
1361
          if (triggerHalfClose) {
1✔
1362
            if (requestSideClosed.compareAndSet(false, true)) {
1✔
1363
              proceedWithHalfClose();
1✔
1364
            }
1365
          }
1366
          break;
1367
        }
1368

1369
        if (isMutated) {
1✔
1370
          super.sendMessage(new KnownLengthInputStream((ByteString) msg));
1✔
1371
        } else {
1372
          super.sendMessage((InputStream) msg);
1✔
1373
        }
1374
      }
1✔
1375
    }
1✔
1376

1377
    private void handleFailOpen(DataPlaneListener listener) {
1378
      if (!activateCall()) {
1✔
1379
        drainPendingRequests();
1✔
1380
      }
1381
      listener.unblockAfterStreamComplete();
1✔
1382
      closeExtProcStream();
1✔
1383
    }
1✔
1384

1385
    private void checkEndOfStream(ProcessingResponse response) {
1386
      boolean terminal = false;
1✔
1387
      if (response.hasResponseTrailers()) {
1✔
1388
        terminal = true;
1✔
1389
      } else if (response.hasResponseHeaders() && wrappedListener.isTrailersOnly()) {
1✔
1390
        terminal = true;
1✔
1391
      }
1392

1393
      if (terminal) {
1✔
1394
        wrappedListener.unblockAfterStreamComplete();
1✔
1395
        closeExtProcStream();
1✔
1396
      }
1397
    }
1✔
1398

1399
    long getServerHeadersStartNanos() {
1400
      return serverHeadersStartNanos;
1✔
1401
    }
1402

1403
    void setServerHeadersStartNanos(long serverHeadersStartNanos) {
1404
      this.serverHeadersStartNanos = serverHeadersStartNanos;
1✔
1405
    }
1✔
1406

1407
    long getServerTrailersStartNanos() {
1408
      return serverTrailersStartNanos;
1✔
1409
    }
1410

1411
    void setServerTrailersStartNanos(long serverTrailersStartNanos) {
1412
      this.serverTrailersStartNanos = serverTrailersStartNanos;
1✔
1413
    }
1✔
1414

1415
    AtomicReference<ExtProcStreamState> getExtProcStreamState() {
1416
      return extProcStreamState;
1✔
1417
    }
1418

1419
    ProcessingMode getCurrentProcessingMode() {
1420
      return currentProcessingMode;
1✔
1421
    }
1422

1423
    AtomicBoolean getPassThroughMode() {
1424
      return passThroughMode;
1✔
1425
    }
1426

1427
    ExternalProcessorFilterConfig getConfig() {
1428
      return config;
1✔
1429
    }
1430

1431
    Context getCallContext() {
1432
      return callContext;
1✔
1433
    }
1434

1435
    ScheduledExecutorService getScheduler() {
1436
      return scheduler;
1✔
1437
    }
1438

1439
    AtomicBoolean getIsProcessingTrailers() {
1440
      return isProcessingTrailers;
1✔
1441
    }
1442
  }
1443

1444
  private static class DataPlaneListener extends SimpleForwardingClientCallListener<InputStream> {
1445
    private final ClientCall<?, ?> rawCall;
1446
    private final DataPlaneClientCall dataPlaneClientCall;
1447
    // Path 3: Upstream response bodies queued because upstream to sidestream window not available,
1448
    // response headers not cleared by ext_proc or ext_proc stream draining
1449
    private final Queue<InputStream> savedMessages = new ConcurrentLinkedQueue<>();
1✔
1450
    private boolean inboundPassThrough = false;
1✔
1451
    @Nullable private volatile Metadata savedHeaders;
1452
    @Nullable private volatile Metadata savedTrailers;
1453
    @Nullable private volatile Status savedStatus;
1454
    private final AtomicBoolean terminationTriggered = new AtomicBoolean(false);
1✔
1455
    private final AtomicBoolean responseHeadersSent = new AtomicBoolean(false);
1✔
1456
    private final AtomicBoolean trailersOnly = new AtomicBoolean(false);
1✔
1457

1458
    protected DataPlaneListener(
1459
        ClientCall.Listener<InputStream> delegate,
1460
        ClientCall<?, ?> rawCall,
1461
        DataPlaneClientCall dataPlaneClientCall) {
1462
      super(delegate);
1✔
1463
      this.rawCall = rawCall;
1✔
1464
      this.dataPlaneClientCall = dataPlaneClientCall;
1✔
1465
    }
1✔
1466

1467
    boolean isTrailersOnly() {
1468
      return trailersOnly.get();
1✔
1469
    }
1470

1471
    Metadata getSavedHeaders() {
1472
      return savedHeaders;
1✔
1473
    }
1474

1475
    Metadata getSavedTrailers() {
1476
      return savedTrailers;
1✔
1477
    }
1478

1479
    void setImmediateResponse(Status status, Metadata trailers) {
1480
      this.savedStatus = status;
1✔
1481
      this.savedTrailers = trailers;
1✔
1482
    }
1✔
1483

1484
    @Override
1485
    public void onReady() {
1486
      dataPlaneClientCall.onReady();
1✔
1487
    }
1✔
1488

1489
    @Override
1490
    public void onHeaders(Metadata headers) {
1491
      dataPlaneClientCall.setServerHeadersStartNanos(System.nanoTime());
1✔
1492
      responseHeadersSent.set(true);
1✔
1493
      boolean sendResponseHeaders =
1✔
1494
          dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1495
              == ProcessingMode.HeaderSendMode.SEND
1496
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1497
              == ProcessingMode.HeaderSendMode.DEFAULT;
1498

1499
      if (dataPlaneClientCall.getExtProcStreamState().get().isDraining() && sendResponseHeaders) {
1✔
1500
        this.savedHeaders = headers;
1✔
1501
        return;
1✔
1502
      }
1503

1504
      if (dataPlaneClientCall.getPassThroughMode().get()
1✔
1505
          || dataPlaneClientCall.getExtProcStreamState().get().isCompleted() 
1✔
1506
          || !sendResponseHeaders) {
1507
        proceedWithHeaders(headers);
1✔
1508
        return;
1✔
1509
      }
1510

1511
      this.savedHeaders = headers;
1✔
1512
      dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1✔
1513
          .setResponseHeaders(HttpHeaders.newBuilder()
1✔
1514
              .setHeaders(
1✔
1515
                  toHeaderMap(headers, dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1✔
1516
              .build())
1✔
1517
          .build());
1✔
1518

1519
      if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1520
        proceedWithHeaders();
1✔
1521
      }
1522
    }
1✔
1523

1524
    @Override
1525
    public void onMessage(InputStream message) {
1526
      synchronized (dataPlaneClientCall.streamLock) {
1✔
1527
        if (inboundPassThrough) {
1✔
1528
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
1✔
1529
          return;
1✔
1530
        }
1531

1532
        boolean checkDrain = dataPlaneClientCall.getExtProcStreamState().get().isDraining()
1✔
1533
            && dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1✔
1534
                == ProcessingMode.BodySendMode.GRPC;
1535

1536
        if (savedHeaders != null || checkDrain) {
1✔
1537
          try {
1538
            ByteString copiedBody = ByteString.readFrom(message);
1✔
1539
            savedMessages.add(new KnownLengthInputStream(copiedBody));
1✔
1540
          } catch (IOException e) {
×
1541
            rawCall.cancel("Failed to copy inbound message for buffering", e);
×
1542
          }
1✔
1543
          return;
1✔
1544
        }
1545

1546
        if (dataPlaneClientCall.getPassThroughMode().get()) {
1✔
1547
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
×
1548
          return;
×
1549
        }
1550

1551
        if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1✔
1552
            || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1✔
1553
                != ProcessingMode.BodySendMode.GRPC) {
1554
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
1✔
1555
          return;
1✔
1556
        }
1557

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

1582
    void drainSavedMessages() {
1583
      synchronized (dataPlaneClientCall.streamLock) {
1✔
1584
        while (dataPlaneClientCall.isSidecarReady()
1✔
1585
            && dataPlaneClientCall.upstreamToSidestreamWindow > 0
1✔
1586
            && !savedMessages.isEmpty()) {
1✔
1587
          InputStream msg = savedMessages.poll();
1✔
1588
          if (msg != null) {
1✔
1589
            try {
1590
              ByteString bodyByteString = ByteString.readFrom(msg);
1✔
1591
              dataPlaneClientCall.upstreamToSidestreamWindow -= bodyByteString.size();
1✔
1592
              sendResponseBodyToExtProc(bodyByteString, false);
1✔
1593
              dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
1✔
1594
            } catch (IOException e) {
×
1595
              rawCall.cancel("Failed to read buffered response body", e);
×
1596
            }
1✔
1597
          }
1598
        }
1✔
1599
        dataPlaneClientCall.drainPendingRequests();
1✔
1600
      }
1✔
1601
    }
1✔
1602

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

1625
      if (this.savedStatus == null) {
1✔
1626
        this.savedStatus = status;
1✔
1627
        this.savedTrailers = trailers;
1✔
1628
      }
1629

1630
      // If we are still waiting for the external processor to validate response headers,
1631
      // buffer the close status/trailers and defer the close until headers are processed.
1632
      if (savedHeaders != null) {
1✔
1633
        return;
1✔
1634
      }
1635

1636
      boolean sendResponseTrailers =
1✔
1637
          dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1✔
1638
              == ProcessingMode.HeaderSendMode.SEND;
1639

1640
      if (dataPlaneClientCall.getExtProcStreamState().get().isDraining() && sendResponseTrailers) {
1✔
1641
        return;
×
1642
      }
1643

1644
      if (!responseHeadersSent.get()) {
1✔
1645
        trailersOnly.set(true);
1✔
1646
      }
1647

1648
      triggerCloseHandshake();
1✔
1649
    }
1✔
1650

1651
    void onReadyNotify() {
1652
      dataPlaneClientCall.getCallContext().run(() -> delegate().onReady());
1✔
1653
    }
1✔
1654

1655
    void proceedWithHeaders() {
1656
      if (savedHeaders != null) {
1✔
1657
        proceedWithHeaders(savedHeaders);
1✔
1658
        synchronized (dataPlaneClientCall.streamLock) {
1✔
1659
          savedHeaders = null;
1✔
1660
          if (!dataPlaneClientCall.getExtProcStreamState().get().isDraining()) {
1✔
1661
            InputStream msg;
1662
            while ((msg = savedMessages.poll()) != null) {
1✔
1663
              onMessage(msg);
1✔
1664
            }
1665
          }
1666
        }
1✔
1667
        onReadyNotify();
1✔
1668
        if (savedStatus != null) {
1✔
1669
          triggerCloseHandshake();
1✔
1670
        }
1671
      }
1672
    }
1✔
1673

1674
    private void proceedWithHeaders(Metadata headers) {
1675
      if (dataPlaneClientCall.getServerHeadersStartNanos() > 0) {
1✔
1676
        long durationNanos = System.nanoTime() - dataPlaneClientCall.getServerHeadersStartNanos();
1✔
1677
        dataPlaneClientCall.recordDuration(serverHeadersDuration, durationNanos);
1✔
1678
        dataPlaneClientCall.setServerHeadersStartNanos(0);
1✔
1679
      }
1680
      dataPlaneClientCall.getCallContext().run(() -> delegate().onHeaders(headers));
1✔
1681
    }
1✔
1682

1683
    void proceedWithClose() {
1684
      if (savedStatus != null) {
1✔
1685
        if (markDataPlaneCallClosed(dataPlaneClientCall.dataPlaneCallState)) {
1✔
1686
          proceedWithClose(savedStatus, savedTrailers);
1✔
1687
        }
1688
        savedStatus = null;
1✔
1689
        savedTrailers = null;
1✔
1690
      }
1691
    }
1✔
1692

1693
    private void proceedWithClose(Status status, Metadata trailers) {
1694
      if (dataPlaneClientCall.getServerTrailersStartNanos() > 0) {
1✔
1695
        long durationNanos = System.nanoTime() - dataPlaneClientCall.getServerTrailersStartNanos();
1✔
1696
        dataPlaneClientCall.recordDuration(serverTrailersDuration, durationNanos);
1✔
1697
        dataPlaneClientCall.setServerTrailersStartNanos(0);
1✔
1698
      }
1699
      dataPlaneClientCall.getCallContext().run(() -> delegate().onClose(status, trailers));
1✔
1700
    }
1✔
1701

1702
    void onExternalBody(ByteString body) {
1703
      // If needed, downstream reading can be made more optimal by creating a wrapped
1704
      // Inputstream wraps the underlying bytestring and that implements HasByteBuffer,
1705
      // Detachable, KnownLength
1706
      dataPlaneClientCall.getCallContext().run(
1✔
1707
          () -> delegate().onMessage(body.newInput()));
1✔
1708
    }
1✔
1709

1710
    void unblockAfterStreamComplete() {
1711
      proceedWithHeaders();
1✔
1712
      // 1. Drain mutated responses first
1713
      dataPlaneClientCall.drainPendingMutatedResponseBodiesDirect(this);
1✔
1714
      // 2. Drain raw responses
1715
      proceedWithSavedMessages();
1✔
1716
      // 3. Drain outbound requests
1717
      dataPlaneClientCall.drainPendingDrainingMessages();
1✔
1718
      proceedWithClose();
1✔
1719
    }
1✔
1720

1721
    private void proceedWithSavedMessages() {
1722
      synchronized (dataPlaneClientCall.streamLock) {
1✔
1723
        InputStream msg;
1724
        while ((msg = savedMessages.poll()) != null) {
1✔
1725
          final InputStream finalMsg = msg;
×
1726
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(finalMsg));
×
1727
        }
×
1728
        inboundPassThrough = true;
1✔
1729
      }
1✔
1730
    }
1✔
1731

1732
    private void triggerCloseHandshake() {
1733
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1✔
1734
          || !terminationTriggered.compareAndSet(false, true)) {
1✔
1735
        return;
1✔
1736
      }
1737

1738
      boolean sendResponseHeaders =
1✔
1739
          dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1740
              == ProcessingMode.HeaderSendMode.SEND
1741
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1✔
1742
              == ProcessingMode.HeaderSendMode.DEFAULT;
1743

1744
      boolean sendResponseTrailers =
1✔
1745
          dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1✔
1746
              == ProcessingMode.HeaderSendMode.SEND;
1747

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

1782
      if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1✔
1783
        proceedWithClose();
1✔
1784
        @SuppressWarnings("unused")
1785
        ScheduledFuture<?> unused = dataPlaneClientCall.getScheduler().schedule(
1✔
1786
            dataPlaneClientCall::closeExtProcStream,
1✔
1787
            dataPlaneClientCall.getConfig().getDeferredCloseTimeoutNanos(),
1✔
1788
            TimeUnit.NANOSECONDS);
1789
      }
1790
    }
1✔
1791

1792
    @GuardedBy("dataPlaneClientCall.streamLock")
1793
    private void sendResponseBodyToExtProc(
1794
        @Nullable ByteString bodyByteString, boolean endOfStream) {
1795
      if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
1✔
1796
          || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
1✔
1797
              != ProcessingMode.BodySendMode.GRPC) {
1798
        return;
×
1799
      }
1800

1801
      HttpBody.Builder bodyBuilder =
1802
          HttpBody.newBuilder();
1✔
1803
      if (bodyByteString != null) {
1✔
1804
        bodyBuilder.setBody(bodyByteString);
1✔
1805
      }
1806
      bodyBuilder.setEndOfStream(endOfStream);
1✔
1807

1808
      ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1✔
1809
          .setResponseBody(bodyBuilder.build());
1✔
1810
      dataPlaneClientCall.mergeAccumulatedWindowUpdates(builder);
1✔
1811
      dataPlaneClientCall.sendToExtProc(builder.build());
1✔
1812
    }
1✔
1813
  }
1814
}
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