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

grpc / grpc-java / #20490

24 Sep 2026 01:27PM UTC coverage: 89.306% (+0.002%) from 89.304%
#20490

push

github

web-flow
api: move ATTR_ADDRESS_NAME to EquivalentAddressGroup (#13072)

Correction to the implementation of [gRFC
A81](https://github.com/grpc/proposal/blob/master/A81-xds-authority-rewriting.md)-xds-authority-rewriting.
The proposal says:
> Note that the resolver attribute used here should be a general-purpose
one, not something specific to EDS;

The endpoint hostname attribute from gRFC A81 is currently defined in
`XdsInternalAttributes`, which makes it unreachable from other modules.

Move the key to `EquivalentAddressGroup` alongside the other endpoint
attributes and re-export it via `InternalEquivalentAddressGroup`.
`XdsInternalAttributes` held nothing else, so it is removed and its
callers now reference the key directly.

This is needed by the autosharding LB policy [gRFC
A119](https://github.com/grpc/proposal/pull/551), which keys its
endpoint map on the A81 hostname. `autosharding` cannot depend on `xds`
because `xds` will depend on autosharding.

39143 of 43830 relevant lines covered (89.31%)

0.89 hits per line

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

95.79
/../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;
1 ✔
915
      }
916
      if (state.isDraining()) {
1 ✔
917
        return false;
1 ✔
918
      }
919
      ClientCallStreamObserver<ProcessingRequest> observer = extProcClientCallRequestObserver;
1 ✔
920
      return observer != null && observer.isReady();
1 ✔
921
    }
922

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1816
      ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
1 ✔
1817
          .setResponseBody(bodyBuilder.build());
1 ✔
1818
      dataPlaneClientCall.mergeAccumulatedWindowUpdates(builder);
1 ✔
1819
      dataPlaneClientCall.sendToExtProc(builder.build());
1 ✔
1820
    }
1 ✔
1821
  }
1822
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc