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

grpc / grpc-java / #20486

17 Sep 2026 08:05AM UTC coverage: 89.329% (+0.02%) from 89.313%
#20486

push

github

web-flow
xds: fix ext_authz unit test issues creating import issues (#13061)

The fixes were created by patching the currently failing CL , making
fixes and running blaze to fix tests and then moving the public changes
out of it to this PR.

- Remove anonymous AutoValue subclassing in test which triggers Error
Prone. While this decreases statistical test coverage, the branch is
defensive programming that cannot occur in production.
- Fix non-deterministic HashSet iteration order in
CheckRequestBuilderTest.
- Clean up and add explicit imports in test files.

39153 of 43830 relevant lines covered (89.33%)

0.89 hits per line

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

95.99
/../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();
×
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;
1✔
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;
1✔
1734
          dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(finalMsg));
1✔
1735
        }
1✔
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