diff --git a/xds/src/generated/thirdparty/grpc/io/envoyproxy/envoy/service/ext_proc/v3/ExternalProcessorGrpc.java b/xds/src/generated/thirdparty/grpc/io/envoyproxy/envoy/service/ext_proc/v3/ExternalProcessorGrpc.java
index fc3ce3a2723..20064af7844 100644
--- a/xds/src/generated/thirdparty/grpc/io/envoyproxy/envoy/service/ext_proc/v3/ExternalProcessorGrpc.java
+++ b/xds/src/generated/thirdparty/grpc/io/envoyproxy/envoy/service/ext_proc/v3/ExternalProcessorGrpc.java
@@ -4,31 +4,25 @@
/**
*
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
@io.grpc.stub.annotations.GrpcGenerated
@@ -131,31 +125,25 @@ public ExternalProcessorFutureStub newStub(io.grpc.Channel channel, io.grpc.Call
/**
*
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
public interface AsyncService {
@@ -164,7 +152,7 @@ public interface AsyncService {
*
* This begins the bidirectional stream that the data plane will use to
* give the server control over what the filter does. The actual
- * protocol is described by the ProcessingRequest and ProcessingResponse
+ * protocol is described by the ``ProcessingRequest`` and ``ProcessingResponse``
* messages below.
*
*/
@@ -177,31 +165,25 @@ default io.grpc.stub.StreamObserver
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
public static abstract class ExternalProcessorImplBase
@@ -215,31 +197,25 @@ public static abstract class ExternalProcessorImplBase
/**
* A stub to allow clients to do asynchronous rpc calls to service ExternalProcessor.
*
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
public static final class ExternalProcessorStub
@@ -259,7 +235,7 @@ protected ExternalProcessorStub build(
*
* This begins the bidirectional stream that the data plane will use to
* give the server control over what the filter does. The actual
- * protocol is described by the ProcessingRequest and ProcessingResponse
+ * protocol is described by the ``ProcessingRequest`` and ``ProcessingResponse``
* messages below.
*
*/
@@ -273,31 +249,25 @@ public io.grpc.stub.StreamObserver
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
public static final class ExternalProcessorBlockingV2Stub
@@ -317,7 +287,7 @@ protected ExternalProcessorBlockingV2Stub build(
*
* This begins the bidirectional stream that the data plane will use to
* give the server control over what the filter does. The actual
- * protocol is described by the ProcessingRequest and ProcessingResponse
+ * protocol is described by the ``ProcessingRequest`` and ``ProcessingResponse``
* messages below.
*
*/
@@ -332,31 +302,25 @@ protected ExternalProcessorBlockingV2Stub build(
/**
* A stub to allow clients to do limited synchronous rpc calls to service ExternalProcessor.
*
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
public static final class ExternalProcessorBlockingStub
@@ -376,31 +340,25 @@ protected ExternalProcessorBlockingStub build(
/**
* A stub to allow clients to do ListenableFuture-style rpc calls to service ExternalProcessor.
*
- * A service that can access and modify HTTP requests and responses
- * as part of a filter chain.
+ * A service that can access and modify HTTP requests and responses as part of a filter chain.
* The overall external processing protocol works like this:
* 1. The data plane sends to the service information about the HTTP request.
- * 2. The service sends back a ProcessingResponse message that directs
- * the data plane to either stop processing, continue without it, or send
- * it the next chunk of the message body.
- * 3. If so requested, the data plane sends the server the message body in
- * chunks, or the entire body at once. In either case, the server may send
- * back a ProcessingResponse for each message it receives, or wait for
- * a certain amount of body chunks received before streaming back the
- * ProcessingResponse messages.
- * 4. If so requested, the data plane sends the server the HTTP trailers,
- * and the server sends back a ProcessingResponse.
- * 5. At this point, request processing is done, and we pick up again
- * at step 1 when the data plane receives a response from the upstream
- * server.
- * 6. At any point above, if the server closes the gRPC stream cleanly,
- * then the data plane proceeds without consulting the server.
- * 7. At any point above, if the server closes the gRPC stream with an error,
- * then the data plane returns a 500 error to the client, unless the filter
- * was configured to ignore errors.
- * In other words, the process is a request/response conversation, but
- * using a gRPC stream to make it easier for the server to
- * maintain state.
+ * 2. The service sends back a ``ProcessingResponse`` message that directs the data plane to either
+ * stop processing, continue without it, or send it the next chunk of the message body.
+ * 3. If so requested, the data plane sends the server the message body in chunks, or the entire
+ * body at once. In either case, the server may send back a ``ProcessingResponse`` for each
+ * message it receives, or wait for a certain amount of body chunks to be received before
+ * streaming back the ``ProcessingResponse`` messages.
+ * 4. If so requested, the data plane sends the server the HTTP trailers, and the server sends back
+ * a ``ProcessingResponse``.
+ * 5. At this point, request processing is done, and we pick up again at step 1 when the data plane
+ * receives a response from the upstream server.
+ * 6. At any point above, if the server closes the gRPC stream cleanly, then the data plane
+ * proceeds without consulting the server.
+ * 7. At any point above, if the server closes the gRPC stream with an error, then the data plane
+ * returns a ``500`` error to the client, unless the filter was configured to ignore errors.
+ * In other words, the process is a request/response conversation, but using a gRPC stream to make
+ * it easier for the server to maintain state.
*
*/
public static final class ExternalProcessorFutureStub
diff --git a/xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java b/xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java
index 7dd33f3f483..3a0a1c9fbc8 100644
--- a/xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java
+++ b/xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java
@@ -186,11 +186,6 @@ ExternalProcessorFilterConfig getFilterConfig() {
return filterConfig;
}
- @VisibleForTesting
- ManagedChannel getExtProcChannel() {
- return extProcChannel;
- }
-
@Override
@SuppressWarnings("unchecked")
public ClientCall interceptCall(
@@ -232,7 +227,6 @@ public ClientCall interceptCall(
io.grpc.stub.MetadataUtils.newAttachHeadersInterceptor(extraHeaders));
}
-
// The filter chain is preceded by RawMessageClientInterceptor, so ReqT and RespT are
// InputStream.
MethodDescriptor rawMethod =
@@ -290,6 +284,42 @@ private static class DataPlaneClientCall
private final HeaderMutator mutator = HeaderMutator.create();
private final AtomicInteger pendingRequests = new AtomicInteger(0);
private final ProcessingMode currentProcessingMode;
+
+ // Default initial window size
+ private static final long DEFAULT_INITIAL_WINDOW_SIZE = 65536;
+
+ // Outbound (sending) windows
+ private long downstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+ private long upstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+
+ // Inbound (receiving) windows
+ private long sidestreamToUpstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+ private long sidestreamToDownstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+
+ // Threshold to trigger standalone client window updates
+ private static final long WINDOW_UPDATE_THRESHOLD = DEFAULT_INITIAL_WINDOW_SIZE / 2;
+
+ // Path 1: Pending/buffered request body messages from downstream
+ private final Queue pendingRequestBodyMessages = new ConcurrentLinkedQueue<>();
+
+ // Path 2: Buffered request body messages from ext_proc server to forward upstream
+ private final Queue pendingUpstreamBodyMessages =
+ new java.util.concurrent.ConcurrentLinkedQueue<>();
+ // Path 4: Outstanding requests from downstream for pulling responses
+ private int downstreamRequestsPending = 0;
+ // Buffered mutated response bodies from ext_proc server
+ private final Queue pendingMutatedResponseBodies =
+ new java.util.concurrent.ConcurrentLinkedQueue<>();
+ // Deferred half-close flag for upstream direction
+ private final AtomicBoolean pendingUpstreamHalfClose = new AtomicBoolean(false);
+
+ // Accumulated client window updates to send to ext_proc
+ private long accumulatedWindowUpdateSidestreamToUpstream = 0;
+ private long accumulatedWindowUpdateSidestreamToDownstream = 0;
+
+ // Flag to track if FlowControlInit was sent in the initial message
+ private boolean flowControlInitSent = false;
+
private final MethodDescriptor, ?> method;
private final Channel channel;
private final MetricRecorder metricsRecorder;
@@ -316,6 +346,10 @@ private static class DataPlaneClientCall
final AtomicBoolean pendingHalfClose = new AtomicBoolean(false);
final AtomicBoolean bodyMessageSentToExtProc = new AtomicBoolean(false);
+ Object getStreamLock() {
+ return streamLock;
+ }
+
protected DataPlaneClientCall(
DataPlaneDelayedCall delayedCall,
ClientCall rawCall,
@@ -343,8 +377,6 @@ protected DataPlaneClientCall(
this.backendService = checkNotNull(backendService, "backendService");
}
-
-
private void activateCall() {
if ((extProcStreamState.get() == ExtProcStreamState.FAILED
&& !config.getFailureModeAllow()
@@ -411,8 +443,6 @@ private boolean validateCompressionSupport(BodyResponse bodyResponse) {
return true;
}
-
-
@Override
public void start(Listener responseListener, Metadata headers) {
this.callContext = Context.current();
@@ -439,6 +469,25 @@ public void onNext(ProcessingResponse response) {
return;
}
+ if (response.hasServerWindowUpdate()) {
+ ProcessingResponse.ServerWindowUpdate update = response.getServerWindowUpdate();
+ boolean wasReady = isReady();
+ synchronized (streamLock) {
+ downstreamToSidestreamWindow += update.getWindowIncrementDownstreamToSidestream();
+ upstreamToSidestreamWindow += update.getWindowIncrementUpstreamToSidestream();
+ drainPendingRequestBodyMessages();
+ drainPendingRequests();
+ if (wrappedListener != null) {
+ wrappedListener.drainSavedMessages();
+ }
+ }
+ // If isReady() becomes true (depends on updated downstreamToSidestreamWindow),
+ // notify the client application via onReadyNotify() (runs unlocked).
+ if (!wasReady && isReady()) {
+ onReadyNotify();
+ }
+ }
+
if (response.hasImmediateResponse()) {
if (config.getDisableImmediateResponse()) {
internalOnError(Status.UNAVAILABLE
@@ -667,21 +716,94 @@ private void sendToExtProc(ProcessingRequest request) {
requestToSend = ProcessingRequest.newBuilder(requestToSend)
.setObservabilityMode(true)
.build();
+ } else if (!flowControlInitSent) {
+ requestToSend = ProcessingRequest.newBuilder(requestToSend)
+ .setFlowControlInit(ProcessingRequest.FlowControlInit.newBuilder()
+ .setInitialWindowDownstreamToSidestream(DEFAULT_INITIAL_WINDOW_SIZE)
+ .setInitialWindowSidestreamToUpstream(DEFAULT_INITIAL_WINDOW_SIZE)
+ .setInitialWindowUpstreamToSidestreama(DEFAULT_INITIAL_WINDOW_SIZE)
+ .setInitialWindowSidestreamToDownstream(DEFAULT_INITIAL_WINDOW_SIZE)
+ .build())
+ .build();
+ flowControlInitSent = true;
}
extProcClientCallRequestObserver.onNext(requestToSend);
}
}
+ void mergeAccumulatedWindowUpdates(ProcessingRequest.Builder requestBuilder) {
+ synchronized (streamLock) {
+ long incrementUpstream = accumulatedWindowUpdateSidestreamToUpstream;
+ long incrementDownstream = accumulatedWindowUpdateSidestreamToDownstream;
+
+ if (incrementUpstream > 0 || incrementDownstream > 0) {
+ requestBuilder.setClientWindowUpdate(
+ ProcessingRequest.ClientWindowUpdate.newBuilder()
+ .setWindowIncrementSidestreamToUpstream(incrementUpstream)
+ .setWindowIncrementSidestreamToDownstream(incrementDownstream)
+ .build());
+ accumulatedWindowUpdateSidestreamToUpstream -= incrementUpstream;
+ accumulatedWindowUpdateSidestreamToDownstream -= incrementDownstream;
+ sidestreamToUpstreamWindow += incrementUpstream;
+ sidestreamToDownstreamWindow += incrementDownstream;
+ }
+ }
+ }
+
+ private void trySendAccumulatedWindowUpdates() {
+ synchronized (streamLock) {
+ if (extProcStreamState.get().isCompleted()) {
+ return;
+ }
+ long incrementUpstream = accumulatedWindowUpdateSidestreamToUpstream;
+ long incrementDownstream = accumulatedWindowUpdateSidestreamToDownstream;
+
+ boolean shouldSend = (incrementUpstream > 0 || incrementDownstream > 0) && (
+ (incrementUpstream >= WINDOW_UPDATE_THRESHOLD)
+ || (incrementDownstream >= WINDOW_UPDATE_THRESHOLD)
+ || (sidestreamToUpstreamWindow <= 0 && accumulatedWindowUpdateSidestreamToUpstream > 0)
+ || (sidestreamToDownstreamWindow <= 0
+ && accumulatedWindowUpdateSidestreamToDownstream > 0)
+ );
+
+ if (shouldSend) {
+ accumulatedWindowUpdateSidestreamToUpstream -= incrementUpstream;
+ accumulatedWindowUpdateSidestreamToDownstream -= incrementDownstream;
+ sidestreamToUpstreamWindow += incrementUpstream;
+ sidestreamToDownstreamWindow += incrementDownstream;
+
+ sendToExtProc(ProcessingRequest.newBuilder()
+ .setClientWindowUpdate(ProcessingRequest.ClientWindowUpdate.newBuilder()
+ .setWindowIncrementSidestreamToUpstream(incrementUpstream)
+ .setWindowIncrementSidestreamToDownstream(incrementDownstream)
+ .build())
+ .build());
+ }
+ }
+ }
+
private void onExtProcStreamReady() {
drainPendingRequests();
onReadyNotify();
}
- private void drainPendingRequests() {
- int toRequest = pendingRequests.getAndSet(0);
- if (toRequest > 0) {
- super.request(toRequest);
+ void drainPendingRequests() {
+ synchronized (streamLock) {
+ if (config.getObservabilityMode()
+ || currentProcessingMode.getResponseBodyMode() != ProcessingMode.BodySendMode.GRPC) {
+ int toRequest = pendingRequests.getAndSet(0);
+ if (toRequest > 0) {
+ super.request(toRequest);
+ }
+ return;
+ }
+
+ // Normal mode flow control: pull 1 message at a time
+ if (isSidecarReady() && upstreamToSidestreamWindow > 0 && pendingRequests.get() > 0) {
+ super.request(1);
+ pendingRequests.decrementAndGet();
+ }
}
}
@@ -730,7 +852,7 @@ private void onReadyNotify() {
wrappedListener.onReadyNotify();
}
- private boolean isSidecarReady() {
+ boolean isSidecarReady() {
ExtProcStreamState state = extProcStreamState.get();
if (state.isCompleted()) {
return true;
@@ -755,11 +877,14 @@ public boolean isReady() {
if (dataPlaneCallState.get() == DataPlaneCallState.IDLE && !config.getObservabilityMode()) {
return false;
}
- boolean sidecarReady = isSidecarReady();
- if (config.getObservabilityMode()) {
- return super.isReady() && sidecarReady;
+ synchronized (streamLock) {
+ boolean sidecarReady = isSidecarReady();
+ if (config.getObservabilityMode()) {
+ return super.isReady() && sidecarReady;
+ }
+ return downstreamToSidestreamWindow > 0 && sidecarReady
+ && pendingRequestBodyMessages.isEmpty();
}
- return sidecarReady;
}
@Override
@@ -768,16 +893,43 @@ public void request(int numMessages) {
super.request(numMessages);
return;
}
- if (!config.getObservabilityMode()
+ if (!config.getObservabilityMode()
&& currentProcessingMode.getResponseBodyMode() != ProcessingMode.BodySendMode.GRPC) {
super.request(numMessages);
return;
}
- if (!isSidecarReady()) {
- pendingRequests.addAndGet(numMessages);
- return;
+ synchronized (streamLock) {
+ boolean sendResponseBodiesToExtProc = config.getObservabilityMode()
+ || currentProcessingMode.getResponseBodyMode() == ProcessingMode.BodySendMode.GRPC;
+
+ if (!sendResponseBodiesToExtProc) {
+ // We do not send response bodies to ext_proc server at all. Bypassed.
+ super.request(numMessages);
+ return;
+ }
+
+ // We send response bodies to ext_proc server (either in normal GRPC mode or
+ // observability mode).
+ // Gated by ext_proc server readiness.
+ // i.e. normal GRPC response body mode
+ boolean normalFlowControl = !config.getObservabilityMode();
+
+ if (normalFlowControl) {
+ pendingRequests.addAndGet(numMessages);
+ downstreamRequestsPending += numMessages;
+ drainPendingMutatedResponseBodies();
+ if (isSidecarReady()) {
+ drainPendingRequests();
+ }
+ } else {
+ // Observability mode: gate on readiness but pull all at once
+ if (isSidecarReady()) {
+ super.request(numMessages);
+ } else {
+ pendingRequests.addAndGet(numMessages);
+ }
+ }
}
- super.request(numMessages);
}
@Override
@@ -812,29 +964,60 @@ public void sendMessage(InputStream message) {
}
return;
}
- }
- if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
- super.sendMessage(message);
- return;
+ if (currentProcessingMode.getRequestBodyMode() == ProcessingMode.BodySendMode.NONE) {
+ super.sendMessage(message);
+ return;
+ }
+
+ // Mode is GRPC
+ try {
+ ByteString bodyByteString = outboundStreamToByteString(message);
+ if (config.getObservabilityMode()) {
+ sendToExtProc(ProcessingRequest.newBuilder()
+ .setRequestBody(HttpBody.newBuilder()
+ .setBody(bodyByteString)
+ .setEndOfStream(false)
+ .build())
+ .build());
+ bodyMessageSentToExtProc.set(true);
+ super.sendMessage(new KnownLengthInputStream(bodyByteString));
+ } else {
+ if (downstreamToSidestreamWindow <= 0 || !pendingRequestBodyMessages.isEmpty()) {
+ pendingRequestBodyMessages.add(bodyByteString);
+ } else {
+ sendRequestBodyToExtProc(bodyByteString);
+ }
+ }
+ } catch (IOException e) {
+ rawCall.cancel("Failed to serialize message for External Processor", e);
+ }
}
+ }
- // Mode is GRPC
- try {
- ByteString bodyByteString = outboundStreamToByteString(message);
- sendToExtProc(ProcessingRequest.newBuilder()
+ private void sendRequestBodyToExtProc(ByteString body) {
+ synchronized (streamLock) {
+ downstreamToSidestreamWindow -= body.size();
+ ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
.setRequestBody(HttpBody.newBuilder()
- .setBody(bodyByteString)
+ .setBody(body)
.setEndOfStream(false)
- .build())
- .build());
+ .build());
+ mergeAccumulatedWindowUpdates(builder);
+ sendToExtProc(builder.build());
bodyMessageSentToExtProc.set(true);
+ }
+ }
- if (config.getObservabilityMode()) {
- super.sendMessage(new KnownLengthInputStream(bodyByteString));
+ private void drainPendingRequestBodyMessages() {
+ synchronized (streamLock) {
+ while (downstreamToSidestreamWindow > 0 && !pendingRequestBodyMessages.isEmpty()) {
+ ByteString body = pendingRequestBodyMessages.poll();
+ sendRequestBodyToExtProc(body);
+ }
+ if (pendingRequestBodyMessages.isEmpty() && pendingHalfClose.get()) {
+ halfClose();
}
- } catch (IOException e) {
- rawCall.cancel("Failed to serialize message for External Processor", e);
}
}
@@ -887,11 +1070,18 @@ public void halfClose() {
}
// Mode is GRPC
- sendToExtProc(ProcessingRequest.newBuilder()
- .setRequestBody(HttpBody.newBuilder()
- .setEndOfStreamWithoutMessage(true)
- .build())
- .build());
+ synchronized (streamLock) {
+ if (!pendingRequestBodyMessages.isEmpty()) {
+ return;
+ }
+
+ ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
+ .setRequestBody(HttpBody.newBuilder()
+ .setEndOfStreamWithoutMessage(true)
+ .build());
+ mergeAccumulatedWindowUpdates(builder);
+ sendToExtProc(builder.build());
+ }
}
@Override
@@ -914,11 +1104,39 @@ private void handleRequestBodyResponse(BodyResponse bodyResponse) {
if (mutation.hasStreamedResponse()) {
StreamedBodyResponse streamed = mutation.getStreamedResponse();
if (!streamed.getEndOfStreamWithoutMessage()) {
- super.sendMessage(new KnownLengthInputStream(streamed.getBody()));
+ com.google.protobuf.ByteString body = streamed.getBody();
+ boolean sendImmediately = false;
+ synchronized (streamLock) {
+ if (sidestreamToUpstreamWindow <= 0) {
+ internalOnError(Status.INTERNAL
+ .withDescription(
+ "Flow control violation: received client body from ext_proc "
+ + "when window is closed")
+ .asRuntimeException());
+ return;
+ }
+ sidestreamToUpstreamWindow -= body.size();
+ if (super.isReady() && pendingUpstreamBodyMessages.isEmpty()) {
+ sendImmediately = true;
+ accumulatedWindowUpdateSidestreamToUpstream += body.size();
+ } else {
+ pendingUpstreamBodyMessages.add(body);
+ }
+ }
+ if (sendImmediately) {
+ super.sendMessage(new KnownLengthInputStream(body));
+ trySendAccumulatedWindowUpdates();
+ }
}
if (streamed.getEndOfStream() || streamed.getEndOfStreamWithoutMessage()) {
- if (requestSideClosed.compareAndSet(false, true)) {
- proceedWithHalfClose();
+ synchronized (streamLock) {
+ if (pendingUpstreamBodyMessages.isEmpty()) {
+ if (requestSideClosed.compareAndSet(false, true)) {
+ proceedWithHalfClose();
+ }
+ } else {
+ pendingUpstreamHalfClose.set(true);
+ }
}
}
}
@@ -931,7 +1149,89 @@ private void handleResponseBodyResponse(
BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
if (mutation.hasStreamedResponse()) {
StreamedBodyResponse streamed = mutation.getStreamedResponse();
- listener.onExternalBody(streamed.getBody());
+ com.google.protobuf.ByteString body = streamed.getBody();
+ final int bodySize = body.size();
+ synchronized (streamLock) {
+ if (sidestreamToDownstreamWindow <= 0) {
+ internalOnError(Status.INTERNAL
+ .withDescription(
+ "Flow control violation: received server body from ext_proc "
+ + "when window is closed")
+ .asRuntimeException());
+ return;
+ }
+ sidestreamToDownstreamWindow -= bodySize;
+ }
+ deliverResponseBody(body, listener);
+ }
+ }
+ }
+
+ private void deliverResponseBody(ByteString body, DataPlaneListener listener) {
+ synchronized (streamLock) {
+ if (downstreamRequestsPending > 0) {
+ downstreamRequestsPending--;
+ final int bodySize = body.size();
+ callContext.run(() -> {
+ try {
+ listener.onExternalBody(body);
+ } finally {
+ synchronized (streamLock) {
+ accumulatedWindowUpdateSidestreamToDownstream += bodySize;
+ }
+ trySendAccumulatedWindowUpdates();
+ }
+ });
+ } else {
+ pendingMutatedResponseBodies.add(body);
+ }
+ }
+ }
+
+ private void drainPendingMutatedResponseBodies() {
+ synchronized (streamLock) {
+ while (downstreamRequestsPending > 0 && !pendingMutatedResponseBodies.isEmpty()) {
+ ByteString body = pendingMutatedResponseBodies.poll();
+ downstreamRequestsPending--;
+ pendingRequests.decrementAndGet();
+ final int bodySize = body.size();
+ callContext.run(() -> {
+ try {
+ wrappedListener.onExternalBody(body);
+ } finally {
+ synchronized (streamLock) {
+ accumulatedWindowUpdateSidestreamToDownstream += bodySize;
+ }
+ trySendAccumulatedWindowUpdates();
+ }
+ });
+ }
+ }
+ }
+
+ void drainPendingUpstreamBodyMessages() {
+ while (true) {
+ ByteString body = null;
+ boolean triggerHalfClose = false;
+ synchronized (streamLock) {
+ if (super.isReady() && !pendingUpstreamBodyMessages.isEmpty()) {
+ body = pendingUpstreamBodyMessages.poll();
+ accumulatedWindowUpdateSidestreamToUpstream += body.size();
+ if (pendingUpstreamBodyMessages.isEmpty()
+ && pendingUpstreamHalfClose.compareAndSet(true, false)) {
+ triggerHalfClose = true;
+ }
+ }
+ }
+ if (body == null) {
+ break;
+ }
+ super.sendMessage(new KnownLengthInputStream(body));
+ trySendAccumulatedWindowUpdates();
+ if (triggerHalfClose) {
+ if (requestSideClosed.compareAndSet(false, true)) {
+ proceedWithHalfClose();
+ }
}
}
}
@@ -1048,6 +1348,8 @@ AtomicBoolean getIsProcessingTrailers() {
private static class DataPlaneListener extends SimpleForwardingClientCallListener {
private final ClientCall, ?> rawCall;
private final DataPlaneClientCall dataPlaneClientCall;
+ // Path 3: Upstream response bodies queued because upstream to sidestream window not available,
+ // response headers not cleared by ext_proc or ext_proc stream draining
private final Queue savedMessages = new ConcurrentLinkedQueue<>();
private boolean inboundPassThrough = false;
@Nullable private volatile Metadata savedHeaders;
@@ -1085,6 +1387,8 @@ void setImmediateResponse(Status status, Metadata trailers) {
@Override
public void onReady() {
+ dataPlaneClientCall.drainPendingUpstreamBodyMessages();
+ dataPlaneClientCall.trySendAccumulatedWindowUpdates();
dataPlaneClientCall.drainPendingRequests();
onReadyNotify();
}
@@ -1104,7 +1408,7 @@ public void onHeaders(Metadata headers) {
return;
}
- if (dataPlaneClientCall.getPassThroughMode().get()
+ if (dataPlaneClientCall.getPassThroughMode().get()
|| dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
|| !sendResponseHeaders) {
proceedWithHeaders(headers);
@@ -1126,7 +1430,7 @@ public void onHeaders(Metadata headers) {
@Override
public void onMessage(InputStream message) {
- synchronized (savedMessages) {
+ synchronized (dataPlaneClientCall.getStreamLock()) {
if (inboundPassThrough) {
dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
return;
@@ -1145,34 +1449,60 @@ public void onMessage(InputStream message) {
}
return;
}
- }
- if (dataPlaneClientCall.getPassThroughMode().get()) {
- dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
- return;
- }
+ if (dataPlaneClientCall.getPassThroughMode().get()) {
+ dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
+ return;
+ }
- if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
- || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
- != ProcessingMode.BodySendMode.GRPC) {
- dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
- return;
- }
+ if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
+ || dataPlaneClientCall.getCurrentProcessingMode().getResponseBodyMode()
+ != ProcessingMode.BodySendMode.GRPC) {
+ dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
+ return;
+ }
- try {
- ByteString bodyByteString = ByteString.readFrom(message);
- sendResponseBodyToExtProc(bodyByteString, false);
- dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
+ try {
+ ByteString bodyByteString = ByteString.readFrom(message);
+ if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
+ sendResponseBodyToExtProc(bodyByteString, false);
+ dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
+ dataPlaneClientCall.getCallContext().run(
+ () -> delegate().onMessage(bodyByteString.newInput()));
+ } else {
+ if (dataPlaneClientCall.upstreamToSidestreamWindow <= 0 || !savedMessages.isEmpty()) {
+ savedMessages.add(new KnownLengthInputStream(bodyByteString));
+ } else {
+ dataPlaneClientCall.upstreamToSidestreamWindow -= bodyByteString.size();
+ sendResponseBodyToExtProc(bodyByteString, false);
+ dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
+ }
+ dataPlaneClientCall.drainPendingRequests();
+ }
+ } catch (IOException e) {
+ rawCall.cancel("Failed to read server response", e);
+ }
+ }
+ }
- if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
- // If needed, downstream reading can be made more optimal by creating a wrapped
- // Inputstream wraps the underlying bytestring and that implements HasByteBuffer,
- // Detachable, KnownLength
- dataPlaneClientCall.getCallContext().run(
- () -> delegate().onMessage(bodyByteString.newInput()));
+ void drainSavedMessages() {
+ synchronized (dataPlaneClientCall.getStreamLock()) {
+ while (dataPlaneClientCall.isSidecarReady()
+ && dataPlaneClientCall.upstreamToSidestreamWindow > 0
+ && !savedMessages.isEmpty()) {
+ InputStream msg = savedMessages.poll();
+ if (msg != null) {
+ try {
+ ByteString bodyByteString = ByteString.readFrom(msg);
+ dataPlaneClientCall.upstreamToSidestreamWindow -= bodyByteString.size();
+ sendResponseBodyToExtProc(bodyByteString, false);
+ dataPlaneClientCall.bodyMessageSentToExtProc.set(true);
+ } catch (IOException e) {
+ rawCall.cancel("Failed to read buffered response body", e);
+ }
+ }
}
- } catch (IOException e) {
- rawCall.cancel("Failed to read server response", e);
+ dataPlaneClientCall.drainPendingRequests();
}
}
@@ -1231,7 +1561,7 @@ void onReadyNotify() {
void proceedWithHeaders() {
if (savedHeaders != null) {
proceedWithHeaders(savedHeaders);
- synchronized (savedMessages) {
+ synchronized (dataPlaneClientCall.getStreamLock()) {
savedHeaders = null;
if (!dataPlaneClientCall.getExtProcStreamState().get().isDraining()) {
InputStream msg;
@@ -1291,7 +1621,7 @@ void unblockAfterStreamComplete() {
}
private void proceedWithSavedMessages() {
- synchronized (savedMessages) {
+ synchronized (dataPlaneClientCall.getStreamLock()) {
InputStream msg;
while ((msg = savedMessages.poll()) != null) {
final InputStream finalMsg = msg;
@@ -1376,9 +1706,10 @@ private void sendResponseBodyToExtProc(
}
bodyBuilder.setEndOfStream(endOfStream);
- dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
- .setResponseBody(bodyBuilder.build())
- .build());
+ ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
+ .setResponseBody(bodyBuilder.build());
+ dataPlaneClientCall.mergeAccumulatedWindowUpdates(builder);
+ dataPlaneClientCall.sendToExtProc(builder.build());
}
}
}
diff --git a/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java b/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java
index 2f7387c1f12..de507889ed5 100644
--- a/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java
+++ b/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java
@@ -71,11 +71,8 @@
import io.grpc.stub.StreamObserver;
import io.grpc.testing.GrpcCleanupRule;
import io.grpc.util.MutableHandlerRegistry;
-import io.grpc.xds.ConfigOrError;
import io.grpc.xds.ExternalProcessorFilter.ExternalProcessorFilterConfig;
import io.grpc.xds.ExternalProcessorFilter.ExternalProcessorFilterOverrideConfig;
-import io.grpc.xds.Filter;
-import io.grpc.xds.XdsNameResolver;
import io.grpc.xds.client.Bootstrapper;
import io.grpc.xds.client.EnvoyProtoData.Node;
import io.grpc.xds.internal.grpcservice.CachedChannelManager;
@@ -91,6 +88,7 @@
import java.util.Collection;
import java.util.Collections;
import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
@@ -112,6 +110,10 @@
*/
@RunWith(JUnit4.class)
public class ExternalProcessorClientInterceptorTest {
+ private static final String INSECURE_CREDENTIALS_TYPE_URL =
+ "type.googleapis.com/envoy.extensions.grpc_service."
+ + "channel_credentials.insecure.v3.InsecureCredentials";
+
static {
System.setProperty("GRPC_EXPERIMENTAL_XDS_EXT_PROC_ON_CLIENT", "true");
}
@@ -3461,7 +3463,6 @@ public void onClose(Status status, Metadata trailers) {
long startTime = System.currentTimeMillis();
while (sidecarBodyLatch.getCount() > 0 && System.currentTimeMillis() - startTime < 5000) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(10);
}
assertThat(capturedRequest.get().getResponseBody().getBody().toStringUtf8())
.isEqualTo("Server Message");
@@ -3469,7 +3470,6 @@ public void onClose(Status status, Metadata trailers) {
while ((appMessageLatch.getCount() > 0 || appCloseLatch.getCount() > 0)
&& System.currentTimeMillis() - startTime < 5000) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(10);
}
proxyCall.cancel("Cleanup", null);
@@ -3619,16 +3619,13 @@ public void onClose(Status status, Metadata trailers) {
long startTime = System.currentTimeMillis();
while (sidecarBodyLatch.getCount() > 0 && System.currentTimeMillis() - startTime < 5000) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(10);
}
while (appMessageLatch.getCount() > 0 && System.currentTimeMillis() - startTime < 5000) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(10);
}
assertThat(capturedMessage.get()).isEqualTo("Mutated Server");
while (appCloseLatch.getCount() > 0 && System.currentTimeMillis() - startTime < 5000) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(10);
}
proxyCall.cancel("Cleanup", null);
@@ -5946,7 +5943,6 @@ public void onReady() {
// Wait for sidecar call to start and listener to be captured
long startTime = System.currentTimeMillis();
while (sidecarListenerRef.get() == null && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
}
assertThat(sidecarListenerRef.get()).isNotNull();
@@ -6270,7 +6266,6 @@ public void onClose(Status status, Metadata trailers) {
assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
// Wait for the drain signal to be received and processed by client call
- Thread.sleep(100);
// Call is now in DRAINING state.
// Send a message. Since request_body_mode is NONE, it should go directly to data plane.
@@ -6769,13 +6764,12 @@ public void onMessage(String message) {
assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
// Wait for the drain signal to be received and processed by client call
- Thread.sleep(100);
// Send response headers first (they bypass ext_proc because send mode is default SKIP, so
// they proceed immediately)
StreamObserver upstreamResponseObserver = dataPlaneResponseObserverRef.get();
upstreamResponseObserver.onNext("Dummy for headers");
-
+
// Now call is in DRAINING state, and savedHeaders is null.
// Send response body message. Since response_body_mode is NONE, it should go directly
// downstream.
@@ -6894,7 +6888,6 @@ public void onHeaders(Metadata headers) {
assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
// Wait for the drain signal to be received and processed by client call
- Thread.sleep(100);
// Call is in DRAINING state.
// Send response headers from server. Since response_header_mode is SKIP, they should go
@@ -7015,7 +7008,6 @@ public void onClose(Status status, Metadata trailers) {
assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
// Wait for the drain signal to be received and processed by client call
- Thread.sleep(100);
// Call is in DRAINING state.
// Complete the server call. Since response_trailer_mode is SKIP, onClose should trigger
@@ -7120,7 +7112,6 @@ public void onCompleted() {
// Use a small loop because of SerializingExecutor delay even with directExecutor.
long start = System.currentTimeMillis();
while (proxyCall.isReady() && System.currentTimeMillis() - start < 2000) {
- Thread.sleep(10);
}
assertThat(proxyCall.isReady()).isFalse();
@@ -7175,10 +7166,10 @@ public void onNext(ProcessingRequest request) {
sidecarOnNextLatch.countDown();
try {
if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) {
- sidecarOnCompletedLatch.countDown();
synchronized (responseObserver) {
responseObserver.onCompleted();
}
+ sidecarOnCompletedLatch.countDown();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -7280,10 +7271,6 @@ public void onReady() {
// After sidecar stream completes, it should trigger onReady and become ready
assertThat(onReadyLatch.await(5, TimeUnit.SECONDS)).isTrue();
- for (int i = 0; i < 50 && !proxyCall.isReady(); i++) {
- fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
- Thread.sleep(10);
- }
assertThat(proxyCall.isReady()).isTrue();
proxyCall.cancel("Cleanup", null);
@@ -7519,7 +7506,6 @@ public void onMessage(String message) {
// Wait for drain to be processed
long startTime = System.currentTimeMillis();
while (proxyCall.isReady() && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
}
assertThat(proxyCall.isReady()).isFalse();
@@ -7532,7 +7518,6 @@ public void onMessage(String message) {
// Wait for it to become ready again
startTime = System.currentTimeMillis();
while (!proxyCall.isReady() && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
}
assertThat(proxyCall.isReady()).isTrue();
@@ -8512,7 +8497,6 @@ public void request(int numMessages) {
// Wait for sidecar call to start
long startTime = System.currentTimeMillis();
while (sidecarListenerRef.get() == null && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
}
assertThat(sidecarListenerRef.get()).isNotNull();
@@ -8638,7 +8622,6 @@ public void request(int numMessages) {
// Wait for drain to be processed
long startTime = System.currentTimeMillis();
while (proxyCall.isReady() && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
}
assertThat(proxyCall.isReady()).isFalse();
@@ -8776,7 +8759,6 @@ public void request(int numMessages) {
// Wait for sidecar call to start
long startTime = System.currentTimeMillis();
while (sidecarListenerRef.get() == null && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
}
assertThat(sidecarListenerRef.get()).isNotNull();
@@ -9761,7 +9743,6 @@ public void onClose(Status status, Metadata trailers) {
// Wait for sidecar to receive headers and filter to activate call
for (int i = 0; i < 5000 && closedLatch.getCount() > 0; i++) {
fakeClock.forwardTime(10, TimeUnit.MILLISECONDS);
- Thread.sleep(1);
}
// Trigger request body processing to hit the unsupported compression check
@@ -9840,7 +9821,8 @@ public void onNext(ProcessingRequest request) {
.build())
.build());
} else if (request.hasResponseBody()) {
- // Simulate sidecar sending compressed body mutation (unsupported) for response body
+ // Simulate sidecar sending compressed body mutation (unsupported) for
+ // response body
responseObserver.onNext(ProcessingResponse.newBuilder()
.setResponseBody(BodyResponse.newBuilder()
.setResponse(CommonResponse.newBuilder()
@@ -10048,7 +10030,6 @@ public void onCompleted() {
for (int i = 0; i < 1000 && finishLatch.getCount() > 0; i++) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(1);
}
assertThat(finishLatch.await(5, TimeUnit.SECONDS)).isTrue();
@@ -10069,12 +10050,138 @@ public void onCompleted() {
|| !extProcServer.isTerminated());
i++) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(1);
}
channelManager.close();
}
}
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testObservabilityMode_ProceedsWithoutBlockingOnExtProcResponseHeaders()
+ throws Exception {
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + uniqueExtProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl("type.googleapis.com/envoy.extensions.grpc_service."
+ + "channel_credentials.insecure.v3.InsecureCredentials")
+ .build())
+ .build())
+ .build())
+ .setObservabilityMode(true)
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final CountDownLatch extProcReceivedHeadersLatch = new CountDownLatch(1);
+ final AtomicReference extProcReceivedRequest = new AtomicReference<>();
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ if (request.hasResponseHeaders()) {
+ extProcReceivedRequest.set(request);
+ extProcReceivedHeadersLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ };
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final CountDownLatch dataPlaneLatch = new CountDownLatch(1);
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
+ (request, responseObserver) -> {
+ responseObserver.onNext("Hello " + request);
+ responseObserver.onCompleted();
+ dataPlaneLatch.countDown();
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ final List appReceivedMessages = new java.util.concurrent.CopyOnWriteArrayList<>();
+ final CountDownLatch appMessageLatch = new CountDownLatch(1);
+ final CountDownLatch appCloseLatch = new CountDownLatch(1);
+ final AtomicReference appReceivedHeaders = new AtomicReference<>();
+ final AtomicReference appReceivedStatus = new AtomicReference<>();
+
+ ClientCall.Listener appListener = new ClientCall.Listener() {
+ @Override
+ public void onHeaders(Metadata headers) {
+ appReceivedHeaders.set(headers);
+ }
+
+ @Override
+ public void onMessage(String message) {
+ appReceivedMessages.add(message);
+ appMessageLatch.countDown();
+ }
+
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ appReceivedStatus.set(status);
+ appCloseLatch.countDown();
+ }
+ };
+
+ CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
+ proxyCall.start(appListener, new Metadata());
+ proxyCall.request(1);
+
+ proxyCall.sendMessage("test");
+ proxyCall.halfClose();
+
+ // Verify data plane server received the request and processed it
+ assertThat(dataPlaneLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // In observability mode, the app should receive response headers and messages immediately
+ // without waiting for the external processor stream to complete.
+ assertThat(appMessageLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(appReceivedHeaders.get()).isNotNull();
+ assertThat(appReceivedMessages).containsExactly("Hello test");
+
+ assertThat(appCloseLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(appReceivedStatus.get().isOk()).isTrue();
+
+ // Also verify that the external processor received the response headers in the background
+ assertThat(extProcReceivedHeadersLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(extProcReceivedRequest.get().hasResponseHeaders()).isTrue();
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
// --- Category 17: Immediate Response Handling ---
@Test
@@ -10390,7 +10497,6 @@ public void onCompleted() {
for (int i = 0; i < 1000 && closedLatch.getCount() > 0; i++) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(1);
}
// Verify app listener notified with an error (not the sidecar's UNAUTHENTICATED)
assertThat(closedLatch.await(5, TimeUnit.SECONDS)).isTrue();
@@ -10404,7 +10510,6 @@ public void onCompleted() {
i < 100 && (!dataPlaneChannel.isTerminated() || !extProcServer.isTerminated());
i++) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(1);
}
channelManager.close();
}
@@ -10514,7 +10619,6 @@ public void onCompleted() {
for (int i = 0; i < 1000 && closedLatch.getCount() > 0; i++) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(1);
}
// Verify app listener notified with UNIMPLEMENTED because data plane connection succeeded
// but the method was not registered, and it failed before the ext-proc stream failed
@@ -10529,7 +10633,6 @@ public void onCompleted() {
i < 100 && (!dataPlaneChannel.isTerminated() || !extProcServer.isTerminated());
i++) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(1);
}
channelManager.close();
}
@@ -10786,7 +10889,6 @@ public void onCompleted() {
// Wait for activation
for (int i = 0; i < 50 && !proxyCall.isReady(); i++) {
fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
- Thread.sleep(10);
}
assertThat(proxyCall.isReady()).isTrue();
@@ -10826,6 +10928,7 @@ public void givenObservabilityModeFalse_whenExtProcBusy_thenIsReadyReturnsFalse(
new java.util.concurrent.CopyOnWriteArrayList<>();
// Sidecar server
final CountDownLatch sidecarActionLatch = new CountDownLatch(1);
+ final CountDownLatch responseSentLatch = new CountDownLatch(1);
ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl;
extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
@@ -10845,6 +10948,7 @@ public void onNext(ProcessingRequest request) {
.setRequestHeaders(HeadersResponse.newBuilder().build())
.build());
}
+ responseSentLatch.countDown();
}
}).start();
}
@@ -10928,10 +11032,7 @@ public boolean isReady() {
// Wait for activation
assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
- for (int i = 0; i < 50 && !proxyCall.isReady(); i++) {
- fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
- Thread.sleep(10);
- }
+ assertThat(responseSentLatch.await(5, TimeUnit.SECONDS)).isTrue();
assertThat(proxyCall.isReady()).isTrue();
// Sidecar becomes busy -> proxyCall becomes busy
@@ -10967,7 +11068,9 @@ public void givenObservabilityModeFalse_whenExtProcBusy_thenAppRequestsAreBuffer
.build())
.build())
.setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
.setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
.setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
.build())
.setObservabilityMode(false)
@@ -10979,6 +11082,7 @@ public void givenObservabilityModeFalse_whenExtProcBusy_thenAppRequestsAreBuffer
// Sidecar server
final CountDownLatch sidecarActionLatch = new CountDownLatch(1);
+ final CountDownLatch responseSentLatch = new CountDownLatch(1);
ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl;
extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
@@ -10997,6 +11101,27 @@ public void onNext(ProcessingRequest request) {
.setRequestHeaders(HeadersResponse.newBuilder().build())
.build());
}
+ responseSentLatch.countDown();
+ } else if (request.hasResponseHeaders()) {
+ synchronized (responseObserver) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ }
+ } else if (request.hasResponseBody()) {
+ synchronized (responseObserver) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(request.getResponseBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ }
}
}).start();
}
@@ -11053,13 +11178,31 @@ public boolean isReady() {
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
- dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello");
- responseObserver.onCompleted();
- }))
- .build());
+ final AtomicReference> dataPlaneResponseObserverRef =
+ new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_BIDI_STREAMING,
+ ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(
+ StreamObserver responseObserver) {
+ dataPlaneResponseObserverRef.set(responseObserver);
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ }))
+ .build());
final AtomicInteger dataPlaneRequestCount = new AtomicInteger(0);
ManagedChannel dataPlaneChannel = grpcCleanup.register(
@@ -11083,116 +11226,4783 @@ public void request(int numMessages) {
CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
- interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
+ interceptCall(interceptor, METHOD_BIDI_STREAMING, callOptions, dataPlaneChannel);
proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1); // Bootstrap request for headers
// Wait for activation
assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
- for (int i = 0; i < 50 && !proxyCall.isReady(); i++) {
- fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
- Thread.sleep(10);
- }
+ assertThat(responseSentLatch.await(5, TimeUnit.SECONDS)).isTrue();
assertThat(proxyCall.isReady()).isTrue();
// Sidecar busy -> request(5) should be buffered
sidecarReady.set(false);
proxyCall.request(5);
- assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(1);
+ // (Only the initial bootstrap request went through)
// Sidecar becomes ready -> buffered requests should be drained
sidecarReady.set(true);
sidecarListenerRef.get().onReady();
long startTime2 = System.currentTimeMillis();
+ while (dataPlaneRequestCount.get() < 2 && System.currentTimeMillis() - startTime2 < 5000) {
+ fakeClock.forwardTime(1, TimeUnit.SECONDS);
+ }
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(2);
+
+ StreamObserver upstreamResponseObserver = dataPlaneResponseObserverRef.get();
+
+ // Server sends response headers (Dummy)
+ upstreamResponseObserver.onNext("Dummy for headers");
+
+ startTime2 = System.currentTimeMillis();
+ while (dataPlaneRequestCount.get() < 3 && System.currentTimeMillis() - startTime2 < 5000) {
+ fakeClock.forwardTime(1, TimeUnit.SECONDS);
+ }
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(3);
+
+ // Server sends first data message -> pulls next
+ upstreamResponseObserver.onNext("Msg 1");
+
+ startTime2 = System.currentTimeMillis();
+ while (dataPlaneRequestCount.get() < 4 && System.currentTimeMillis() - startTime2 < 5000) {
+ fakeClock.forwardTime(1, TimeUnit.SECONDS);
+ }
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(4);
+
+ // Server sends second data message -> pulls next
+ upstreamResponseObserver.onNext("Msg 2");
+
+ startTime2 = System.currentTimeMillis();
while (dataPlaneRequestCount.get() < 5 && System.currentTimeMillis() - startTime2 < 5000) {
fakeClock.forwardTime(1, TimeUnit.SECONDS);
- Thread.sleep(10);
}
assertThat(dataPlaneRequestCount.get()).isEqualTo(5);
+ // Server sends third data message -> pulls next (which drains the final pending request)
+ upstreamResponseObserver.onNext("Msg 3");
+
+ startTime2 = System.currentTimeMillis();
+ while (dataPlaneRequestCount.get() < 6 && System.currentTimeMillis() - startTime2 < 5000) {
+ fakeClock.forwardTime(1, TimeUnit.SECONDS);
+ }
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(6);
+
proxyCall.cancel("Cleanup", null);
channelManager.close();
}
- // --- Category 21: Streaming Completeness (Client & Bi-Di) ---
-
@Test
- @SuppressWarnings({"unchecked", "FutureReturnValueIgnored"})
- public void givenClientStreamingRpc_whenExtProcMutatesAll_thenAllTargetsReceiveMutatedData()
+ @SuppressWarnings("unchecked")
+ public void givenResponseBodyModeNone_whenExtProcBusy_thenAppRequestsAreNotBuffered()
throws Exception {
- String uniqueExtProcServerName =
- "extProc-client-stream-" + InProcessServerBuilder.generateName();
- String uniqueDataPlaneServerName =
- "dataPlane-client-stream-" + InProcessServerBuilder.generateName();
- ExternalProcessor proto = createBaseProto(uniqueExtProcServerName)
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl("type.googleapis.com/envoy.extensions.grpc_service."
+ + "channel_credentials.insecure.v3.InsecureCredentials")
+ .build())
+ .build())
+ .build())
.setProcessingMode(ProcessingMode.newBuilder()
- .setRequestHeaderMode(ProcessingMode.HeaderSendMode.SEND)
.setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
.setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
- .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
.setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
.build())
+ .setObservabilityMode(false)
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
assertThat(configOrError.errorDetail).isNull();
ExternalProcessorFilterConfig filterConfig = configOrError.config;
- final Metadata.Key reqKey =
- Metadata.Key.of("req-mutated", Metadata.ASCII_STRING_MARSHALLER);
-
- final List receivedPhases = Collections.synchronizedList(new ArrayList<>());
- final CountDownLatch sidecarActionLatch = new CountDownLatch(6);
- final ExecutorService sidecarResponseExecutor = Executors.newSingleThreadExecutor();
- // External Processor Server
+ // Sidecar server
+ final CountDownLatch sidecarActionLatch = new CountDownLatch(1);
+ final CountDownLatch responseSentLatch = new CountDownLatch(1);
ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl;
extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
+ @SuppressWarnings("unchecked")
public StreamObserver process(
final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
return new StreamObserver() {
@Override
public void onNext(ProcessingRequest request) {
- sidecarResponseExecutor.submit(() -> {
- synchronized (responseObserver) {
- ProcessingResponse.Builder resp = ProcessingResponse.newBuilder();
- if (request.hasRequestHeaders()) {
- receivedPhases.add("REQ_HEADERS");
- resp.setRequestHeaders(
- HeadersResponse.newBuilder()
- .setResponse(
- CommonResponse.newBuilder()
- .setHeaderMutation(
- HeaderMutation.newBuilder()
- .addSetHeaders(
- io.envoyproxy.envoy.config.core.v3.HeaderValueOption
- .newBuilder()
- .setHeader(
- io.envoyproxy.envoy.config.core.v3.HeaderValue
- .newBuilder()
- .setKey("req-mutated")
- .setValue("true")
- .build())
- .build())
- .build())
- .build())
- .build());
- } else if (request.hasRequestBody()) {
- if (request.getRequestBody().getEndOfStream()
- || request.getRequestBody().getEndOfStreamWithoutMessage()) {
- receivedPhases.add("REQ_BODY_EOS");
- resp.setRequestBody(
- BodyResponse.newBuilder()
- .setResponse(
- CommonResponse.newBuilder()
- .setBodyMutation(
+ new Thread(() -> {
+ if (request.hasRequestHeaders()) {
+ sidecarActionLatch.countDown();
+ synchronized (responseObserver) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ }
+ responseSentLatch.countDown();
+ }
+ }).start();
+ }
+
+ @Override
+ public void onError(Throwable t) {
+ }
+
+ @Override
+ public void onCompleted() {
+ new Thread(() -> responseObserver.onCompleted()).start();
+ }
+ };
+ }
+ };
+ grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ final AtomicBoolean sidecarReady = new AtomicBoolean(true);
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(extProcServerName)
+ .directExecutor()
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new io.grpc.ForwardingClientCall.SimpleForwardingClientCall(
+ next.newCall(method, callOptions)) {
+ @Override
+ public void start(Listener responseListener, Metadata headers) {
+ super.start(responseListener, headers);
+ }
+
+ @Override
+ public boolean isReady() {
+ return sidecarReady.get();
+ }
+ };
+ }
+ })
+ .build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_BIDI_STREAMING, ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ }))
+ .build());
+
+ final AtomicInteger dataPlaneRequestCount = new AtomicInteger(0);
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName)
+ .directExecutor()
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new io.grpc.ForwardingClientCall.SimpleForwardingClientCall(
+ next.newCall(method, callOptions)) {
+ @Override
+ public void request(int numMessages) {
+ dataPlaneRequestCount.addAndGet(numMessages);
+ super.request(numMessages);
+ }
+ };
+ }
+ })
+ .build());
+
+ CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING, callOptions, dataPlaneChannel);
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1); // Bootstrap request
+
+ // Wait for activation
+ assertThat(sidecarActionLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(responseSentLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(proxyCall.isReady()).isTrue();
+
+ // Sidecar server busy
+ sidecarReady.set(false);
+
+ // Since responseBodyMode is NONE and not in observabilityMode, request(5) should
+ // be passed upstream immediately
+ proxyCall.request(5);
+
+ long startTime = System.currentTimeMillis();
+ while (dataPlaneRequestCount.get() < 6 && System.currentTimeMillis() - startTime < 5000) {
+ fakeClock.forwardTime(1, TimeUnit.SECONDS);
+ }
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(6); // 1 bootstrap + 5 requested
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testFlowControlStateInitialization() throws Exception {
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
+ .build())
+ .build())
+ .build())
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests = new CopyOnWriteArrayList<>();
+ final CountDownLatch sidecarLatch = new CountDownLatch(2);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ sidecarLatch.countDown();
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(request.getRequestBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
+ (request, responseObserver) -> {
+ responseObserver.onNext("Hello " + request);
+ responseObserver.onCompleted();
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_SAY_HELLO,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.sendMessage("Message 1");
+
+ assertThat(sidecarLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ assertThat(receivedRequests).hasSize(2);
+ ProcessingRequest firstRequest = receivedRequests.get(0);
+ ProcessingRequest secondRequest = receivedRequests.get(1);
+
+ assertThat(firstRequest.hasRequestHeaders()).isTrue();
+ assertThat(firstRequest.hasFlowControlInit()).isTrue();
+ assertThat(firstRequest.getFlowControlInit().getInitialWindowDownstreamToSidestream())
+ .isEqualTo(65536);
+ assertThat(firstRequest.getFlowControlInit().getInitialWindowSidestreamToUpstream())
+ .isEqualTo(65536);
+
+ assertThat(secondRequest.hasRequestBody()).isTrue();
+ assertThat(secondRequest.hasFlowControlInit()).isFalse();
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testDownstreamToSidestreamFlowControl_EnforcesWindow() throws Exception {
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
+ .build())
+ .build())
+ .build())
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests = new CopyOnWriteArrayList<>();
+ final CountDownLatch firstBodyLatch = new CountDownLatch(2); // Headers + First Body
+ final CountDownLatch secondBodyLatch = new CountDownLatch(1);
+ final AtomicReference>
+ responseObserverRef = new AtomicReference<>();
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ responseObserverRef.set(responseObserver);
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ firstBodyLatch.countDown();
+ } else if (request.hasRequestBody()) {
+ if (request.getRequestBody().getEndOfStreamWithoutMessage()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setEndOfStreamWithoutMessage(true)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ return;
+ }
+ if (firstBodyLatch.getCount() > 0) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(request.getRequestBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ firstBodyLatch.countDown();
+ } else {
+ // This is the second body (30000 bytes)
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(request.getRequestBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ secondBodyLatch.countDown();
+ }
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final List dataPlaneReceivedMessages = new CopyOnWriteArrayList<>();
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_CLIENT_STREAMING, ServerCalls.asyncClientStreamingCall(
+ new ServerCalls.ClientStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {
+ dataPlaneReceivedMessages.add(value);
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onNext("Response");
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ }))
+ .build());
+
+ final List dataPlaneResponseMessages = new CopyOnWriteArrayList<>();
+ final CountDownLatch callClosedLatch = new CountDownLatch(1);
+ final AtomicReference callClosedStatus = new AtomicReference<>();
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_CLIENT_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ dataPlaneResponseMessages.add(message);
+ }
+
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ callClosedStatus.set(status);
+ callClosedLatch.countDown();
+ }
+ }, new Metadata());
+ proxyCall.request(1);
+
+ // Generate large messages
+ String largeMessage70k = new String(new char[70000]).replace('\0', 'a');
+ String largeMessage30k = new String(new char[30000]).replace('\0', 'b');
+
+ // Send first message (70000 bytes) - fits in 65536 window
+ proxyCall.sendMessage(largeMessage70k);
+ assertThat(firstBodyLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(proxyCall.isReady()).isFalse();
+
+ // Send second message (30000 bytes) - total 100000 > 65536, should buffer
+ proxyCall.sendMessage(largeMessage30k);
+
+ // Call halfClose() while the second message is still buffered.
+ // This should NOT trigger immediate half-close, but mark pendingHalfClose = true.
+ proxyCall.halfClose();
+
+ // Assert that it is NOT delivered to ext_proc (delivery is synchronous on
+ // directExecutor, so we can check immediately)
+ assertThat(receivedRequests).hasSize(3);
+ // (Headers + First Body + Client Window Update (Path 2 replenishment))
+ assertThat(proxyCall.isReady()).isFalse();
+
+ // Now send ServerWindowUpdate from ext_proc to interceptor to increment window by 40000
+ responseObserverRef.get().onNext(ProcessingResponse.newBuilder()
+ .setServerWindowUpdate(ProcessingResponse.ServerWindowUpdate.newBuilder()
+ .setWindowIncrementDownstreamToSidestream(40000)
+ .build())
+ .build());
+
+ // The second body should now be flushed and received by ext_proc, and then half-closed
+ assertThat(secondBodyLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // Verify both messages reached the backend service
+ assertThat(dataPlaneReceivedMessages)
+ .containsExactly(largeMessage70k, largeMessage30k).inOrder();
+
+ // Wait for the call to close successfully.
+ assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(callClosedStatus.get().isOk()).isTrue();
+ assertThat(dataPlaneResponseMessages).containsExactly("Response");
+
+ // The mock ext_proc should have received 5 requests:
+ // 1. Headers
+ // 2. First Body (70000)
+ // 3. Client Window Update
+ // 4. Second Body (30000)
+ // 5. EndOfStreamWithoutMessage (half-close)
+ assertThat(receivedRequests.size()).isEqualTo(5);
+ assertThat(receivedRequests.get(0).hasRequestHeaders()).isTrue();
+ assertThat(receivedRequests.get(1).hasRequestBody()).isTrue();
+ assertThat(receivedRequests.get(2).hasClientWindowUpdate()).isTrue();
+ assertThat(receivedRequests.get(3).hasRequestBody()).isTrue();
+ assertThat(receivedRequests.get(3).getRequestBody().getBody().size()).isEqualTo(30000);
+ assertThat(receivedRequests.get(4).hasRequestBody()).isTrue();
+ assertThat(receivedRequests.get(4).getRequestBody().getEndOfStreamWithoutMessage()).isTrue();
+
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testUpstreamToSidestreamFlowControl_EnforcesWindow() throws Exception {
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
+ .build())
+ .build())
+ .build())
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests = new CopyOnWriteArrayList<>();
+ final CountDownLatch sidecarLatch = new CountDownLatch(4);
+ // (Headers, Request Body, Response Headers, Response Body 1)
+ final CountDownLatch secondResponseBodyLatch = new CountDownLatch(1);
+ final AtomicReference>
+ responseObserverRef = new AtomicReference<>();
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ responseObserverRef.set(responseObserver);
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ sidecarLatch.countDown();
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(request.getRequestBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ } else if (request.hasResponseHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasResponseBody()) {
+ com.google.protobuf.ByteString originalBody = request.getResponseBody().getBody();
+ com.google.protobuf.ByteString bodyToSend = originalBody;
+ // Return the original 70,000 bytes as-is
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(bodyToSend)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ if (originalBody.size() == 30000) {
+ secondResponseBodyLatch.countDown();
+ }
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final AtomicReference> dataPlaneResponseObserverRef =
+ new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_BIDI_STREAMING,
+ ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(
+ StreamObserver responseObserver) {
+ dataPlaneResponseObserverRef.set(responseObserver);
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ final List appReceivedMessages = new CopyOnWriteArrayList<>();
+ final CountDownLatch messagesLatch2 = new CountDownLatch(2);
+ final CountDownLatch messagesLatch3 = new CountDownLatch(3);
+ ClientCall.Listener appListener = new ClientCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ appReceivedMessages.add(message);
+ messagesLatch2.countDown();
+ messagesLatch3.countDown();
+ }
+ };
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
+
+ proxyCall.start(appListener, new Metadata());
+ proxyCall.request(10);
+
+ // Send first dummy message to initialize headers and stream
+ proxyCall.sendMessage("Client Msg");
+
+ StreamObserver upstreamResponseObserver = dataPlaneResponseObserverRef.get();
+ upstreamResponseObserver.onNext("Dummy for headers");
+
+ String largeMessage70k = new String(new char[70000]).replace('\0', 'a');
+ String largeMessage30k = new String(new char[30000]).replace('\0', 'b');
+
+ // Wait for the initialization (headers, request body, response headers) to reach the ext_proc
+ // server
+ assertThat(sidecarLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // Upstream sends 70k response chunk. Since window is 65,536, this drives the window
+ // negative (-4,464).
+ upstreamResponseObserver.onNext(largeMessage70k);
+
+ // Verify Chunk 1 is successfully delivered (2 messages total in app: dummy and chunk 1)
+ assertThat(messagesLatch2.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(appReceivedMessages).hasSize(2);
+
+ // Upstream sends 30k response chunk. Since the window is negative, the filter
+ // must block/buffer this chunk.
+ upstreamResponseObserver.onNext(largeMessage30k);
+
+ // Wait a brief period and verify that the 30k chunk has NOT been sent to the ext_proc server
+ assertThat(secondResponseBodyLatch.getCount()).isEqualTo(1);
+ assertThat(appReceivedMessages).hasSize(2);
+
+ // Sidecar server sends a ServerWindowUpdate of 40k to the filter, unblocking the window.
+ responseObserverRef.get().onNext(ProcessingResponse.newBuilder()
+ .setServerWindowUpdate(ProcessingResponse.ServerWindowUpdate.newBuilder()
+ .setWindowIncrementUpstreamToSidestream(40000)
+ .build())
+ .build());
+
+ // Once the window is unblocked, the filter immediately forwards the 30k chunk
+ // to the ext_proc server, which processes it.
+ assertThat(secondResponseBodyLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // Verify that Chunk 2 is now successfully delivered to the client application
+ assertThat(messagesLatch3.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(appReceivedMessages).hasSize(3);
+ assertThat(appReceivedMessages.get(2)).isEqualTo(largeMessage30k);
+ assertThat(receivedRequests).isNotEmpty();
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToDownstreamFlowControl_Violations() throws Exception {
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
+ .build())
+ .build())
+ .build())
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final String mutatedMessageTooLarge = new String(new char[70000]).replace('\0', 'd');
+ final CountDownLatch callClosedLatch = new CountDownLatch(1);
+ final AtomicReference capturedStatus = new AtomicReference<>();
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(request.getRequestBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ } else if (request.hasResponseHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasResponseBody()) {
+ // Send Chunk 1 (70k) -> Consumed by client app's pending request count of 1.
+ // Replenishes window.
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(ByteString.copyFromUtf8(mutatedMessageTooLarge))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+
+ // Send Chunk 2 (70k) -> Excess chunk. Stored in queue. Window stays negative.
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(ByteString.copyFromUtf8(mutatedMessageTooLarge))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+
+ // Send Chunk 3 (1 byte) -> Arrives when window is negative, triggering violation.
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setResponseBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(ByteString.copyFromUtf8("a"))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final AtomicReference> dataPlaneResponseObserverRef =
+ new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_BIDI_STREAMING, ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ dataPlaneResponseObserverRef.set(responseObserver);
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ capturedStatus.set(status);
+ callClosedLatch.countDown();
+ }
+ }, new Metadata());
+ proxyCall.request(1);
+
+ proxyCall.sendMessage("Client Msg");
+
+ // Send a response from upstream to trigger headers and then the body response
+ StreamObserver upstreamResponseObserver = dataPlaneResponseObserverRef.get();
+ upstreamResponseObserver.onNext("Response Msg");
+
+ assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // The call should fail immediately with INTERNAL error code due to flow control violation
+ assertThat(capturedStatus.get().getCode()).isEqualTo(Status.Code.INTERNAL);
+ assertThat(capturedStatus.get().getDescription())
+
+ .isEqualTo("External processor stream failed");
+ assertThat(capturedStatus.get().getCause()).isInstanceOf(io.grpc.StatusRuntimeException.class);
+ assertThat(capturedStatus.get().getCause().getMessage())
+ .contains("Flow control violation: received server body from ext_proc "
+ + "when window is closed");
+
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToUpstreamFlowControl_Violations() throws Exception {
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
+ .build())
+ .build())
+ .build())
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final String mutatedMessageTooLarge = new String(new char[70000]).replace('\0', 'c');
+ final CountDownLatch callClosedLatch = new CountDownLatch(1);
+ final AtomicReference capturedStatus = new AtomicReference<>();
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ ByteString original = request.getRequestBody().getBody();
+ if (original.toStringUtf8().equals("Message 1")) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(ByteString.copyFromUtf8(mutatedMessageTooLarge))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ } else {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder()
+ .setBodyMutation(BodyMutation.newBuilder()
+ .setStreamedResponse(StreamedBodyResponse.newBuilder()
+ .setBody(ByteString.copyFromUtf8("a"))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ }
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_CLIENT_STREAMING, ServerCalls.asyncClientStreamingCall(
+ new ServerCalls.ClientStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onNext("Response");
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ }))
+ .build());
+
+ // Build the data plane channel with a custom ClientInterceptor.
+ // This interceptor overrides isReady() to always return false.
+ // This simulates a blocked backend server
+ // (transport flow control buffer is full) and blocks the filter's upstream
+ // window replenishment.
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName)
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new io.grpc.ForwardingClientCall.SimpleForwardingClientCall(
+ next.newCall(method, callOptions)) {
+ @Override
+ public boolean isReady() {
+ return false;
+ }
+ };
+ }
+ })
+ .directExecutor()
+ .build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_CLIENT_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()), dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ capturedStatus.set(status);
+ callClosedLatch.countDown();
+ }
+ }, new Metadata());
+
+ // Send first message. Window becomes negative.
+ proxyCall.sendMessage("Message 1");
+
+ // Send second message. Window is still negative, triggering flow control violation.
+ proxyCall.sendMessage("Message 2");
+
+ assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // The call should fail immediately with INTERNAL error code due to flow control violation
+ assertThat(capturedStatus.get().getCode()).isEqualTo(Status.Code.INTERNAL);
+ assertThat(capturedStatus.get().getDescription())
+
+ .isEqualTo("External processor stream failed");
+ assertThat(capturedStatus.get().getCause()).isInstanceOf(io.grpc.StatusRuntimeException.class);
+ assertThat(capturedStatus.get().getCause().getMessage())
+ .contains("Flow control violation: received client body from ext_proc "
+ + "when window is closed");
+
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToUpstreamFlowControl_QueuingAndDraining() throws Exception {
+ ExternalProcessor proto =
+ ExternalProcessor.newBuilder()
+ .setGrpcService(
+ GrpcService.newBuilder()
+ .setGoogleGrpc(
+ GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(
+ Any.newBuilder().setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL).build())
+ .build())
+ .build())
+ .setProcessingMode(
+ ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final CountDownLatch finishLatch = new CountDownLatch(1);
+ final List serverReceivedBodies = new CopyOnWriteArrayList<>();
+ final CountDownLatch serverReceivedLatch = new CountDownLatch(2);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ ByteString original = request.getRequestBody().getBody();
+ boolean eos =
+ request.getRequestBody().getEndOfStream()
+ || request.getRequestBody().getEndOfStreamWithoutMessage();
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(
+ ByteString.copyFromUtf8(
+ eos
+ ? ""
+ : "Mutated"
+ + original
+ .toStringUtf8()))
+ .setEndOfStream(eos)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_CLIENT_STREAMING,
+ ServerCalls.asyncClientStreamingCall(
+ new ServerCalls.ClientStreamingMethod() {
+ @Override
+ public StreamObserver invoke(
+ StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {
+ serverReceivedBodies.add(value);
+ serverReceivedLatch.countDown();
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onNext("Response");
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ }))
+ .build());
+
+ final AtomicBoolean transportReady = new AtomicBoolean(false);
+ final AtomicReference> capturedListenerRef = new AtomicReference<>();
+
+ class TriggerableForwardingCall
+ extends io.grpc.ForwardingClientCall.SimpleForwardingClientCall {
+ TriggerableForwardingCall(ClientCall delegate) {
+ super(delegate);
+ }
+
+ @Override
+ public void start(Listener responseListener, Metadata headers) {
+ capturedListenerRef.set(responseListener);
+ super.start(responseListener, headers);
+ }
+
+ @Override
+ public boolean isReady() {
+ return transportReady.get();
+ }
+ }
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName)
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new TriggerableForwardingCall<>(next.newCall(method, callOptions));
+ }
+ })
+ .directExecutor()
+ .build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_CLIENT_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()), dataPlaneChannel);
+
+ proxyCall.start(
+ new ClientCall.Listener() {
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ finishLatch.countDown();
+ }
+ },
+ new Metadata());
+ proxyCall.request(1);
+
+ // Send first message. This gets mutated to "MutatedOriginalRequest 1" by ext_proc.
+ proxyCall.sendMessage("OriginalRequest 1");
+
+ // Give some time to process and ensure the message is NOT received on the server side because
+ // transport is not ready
+ assertThat(serverReceivedLatch.await(500, TimeUnit.MILLISECONDS)).isFalse();
+ assertThat(serverReceivedBodies).isEmpty();
+
+ // Now make the transport ready (super.isReady() returns true), but do NOT trigger onReady drain
+ // yet.
+ transportReady.set(true);
+
+ // Send second message. This gets mutated to "MutatedOriginalRequest 2" by ext_proc.
+ // Since transportReady is true but there's still a pending message in the queue,
+ // the second message should also be queued (to preserve order).
+ proxyCall.sendMessage("OriginalRequest 2");
+
+ // Ensure still no message is received on the server side (since we haven't triggered drain via
+ // onReady)
+ assertThat(serverReceivedLatch.await(500, TimeUnit.MILLISECONDS)).isFalse();
+ assertThat(serverReceivedBodies).isEmpty();
+
+ // Now trigger onReady callback to drain the queue.
+ ClientCall.Listener> listener = capturedListenerRef.get();
+ assertThat(listener).isNotNull();
+ listener.onReady();
+
+ // Both messages should be drained and forwarded to the backend server in order.
+ assertThat(serverReceivedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(serverReceivedBodies)
+ .containsExactly("MutatedOriginalRequest 1", "MutatedOriginalRequest 2")
+ .inOrder();
+
+ proxyCall.halfClose();
+ assertThat(finishLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToUpstreamFlowControl_DelayedHalfClose() throws Exception {
+ ExternalProcessor proto =
+ ExternalProcessor.newBuilder()
+ .setGrpcService(
+ GrpcService.newBuilder()
+ .setGoogleGrpc(
+ GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(
+ Any.newBuilder().setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL).build())
+ .build())
+ .build())
+ .setProcessingMode(
+ ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final CountDownLatch sidecarLatch = new CountDownLatch(1);
+ final List serverReceivedBodies = new CopyOnWriteArrayList<>();
+ final CountDownLatch serverReceivedLatch = new CountDownLatch(1);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(
+ ByteString.copyFromUtf8("Mutated1"))
+ .setEndOfStream(true)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ sidecarLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_CLIENT_STREAMING,
+ ServerCalls.asyncClientStreamingCall(
+ new ServerCalls.ClientStreamingMethod() {
+ @Override
+ public StreamObserver invoke(
+ StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {
+ serverReceivedBodies.add(value);
+ serverReceivedLatch.countDown();
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onNext("Response");
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ }))
+ .build());
+
+ final AtomicBoolean transportReady = new AtomicBoolean(false);
+ final AtomicReference> capturedListenerRef = new AtomicReference<>();
+ final AtomicInteger halfCloseCallCount = new AtomicInteger(0);
+
+ class DelayedHalfCloseForwardingCall
+ extends io.grpc.ForwardingClientCall.SimpleForwardingClientCall {
+ DelayedHalfCloseForwardingCall(ClientCall delegate) {
+ super(delegate);
+ }
+
+ @Override
+ public void start(Listener responseListener, Metadata headers) {
+ capturedListenerRef.set(responseListener);
+ super.start(responseListener, headers);
+ }
+
+ @Override
+ public boolean isReady() {
+ return transportReady.get();
+ }
+
+ @Override
+ public void halfClose() {
+ halfCloseCallCount.incrementAndGet();
+ super.halfClose();
+ }
+ }
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName)
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new DelayedHalfCloseForwardingCall<>(next.newCall(method, callOptions));
+ }
+ })
+ .directExecutor()
+ .build());
+
+ ClientCall proxyCall =
+ interceptCall(
+ interceptor,
+ METHOD_CLIENT_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1);
+
+ // Send the first client message. This gets mutated to "Mutated1" by ext_proc.
+ proxyCall.sendMessage("OriginalRequest 1");
+
+ assertThat(sidecarLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // Since transportReady is false, the mutated body is queued in pendingUpstreamBodyMessages.
+ // And since it was unilateral half-close, pendingUpstreamHalfClose is set to true.
+ // Verify that the call is NOT half-closed on transport yet.
+ assertThat(halfCloseCallCount.get()).isEqualTo(0);
+ assertThat(serverReceivedBodies).isEmpty();
+
+ // Now make the transport ready and trigger onReady callback
+ transportReady.set(true);
+ ClientCall.Listener> listener = capturedListenerRef.get();
+ assertThat(listener).isNotNull();
+ listener.onReady();
+
+ // The queued message should be drained, forwarded to backend server,
+ // and the delayed half-close should be triggered.
+ assertThat(serverReceivedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(serverReceivedBodies).containsExactly("Mutated1");
+ assertThat(halfCloseCallCount.get()).isEqualTo(1);
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToUpstreamFlowControl_FailOpenDuringDelayedHalfClose()
+ throws Exception {
+ ExternalProcessor proto =
+ ExternalProcessor.newBuilder()
+ .setGrpcService(
+ GrpcService.newBuilder()
+ .setGoogleGrpc(
+ GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(
+ Any.newBuilder().setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL).build())
+ .build())
+ .build())
+ .setProcessingMode(
+ ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
+ .build())
+ .setFailureModeAllow(true)
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final CountDownLatch sidecarLatch = new CountDownLatch(1);
+ final AtomicReference> responseObserverRef =
+ new AtomicReference<>();
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ responseObserverRef.set(responseObserver);
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ // Unilaterally send a request body response containing mutated body and
+ // endOfStream = true
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(
+ ByteString.copyFromUtf8("Mutated1"))
+ .setEndOfStream(true)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ sidecarLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(
+ InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build()
+ .start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor =
+ new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_CLIENT_STREAMING,
+ ServerCalls.asyncClientStreamingCall(
+ new ServerCalls.ClientStreamingMethod() {
+ @Override
+ public StreamObserver invoke(
+ StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onNext("Response");
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ }))
+ .build());
+
+ final AtomicBoolean transportReady = new AtomicBoolean(false);
+ final AtomicReference> capturedListenerRef = new AtomicReference<>();
+ final AtomicInteger halfCloseCallCount = new AtomicInteger(0);
+ final AtomicInteger sendMessageCount = new AtomicInteger(0);
+
+ class FailOpenDelayedHalfCloseForwardingCall
+ extends io.grpc.ForwardingClientCall.SimpleForwardingClientCall {
+ FailOpenDelayedHalfCloseForwardingCall(ClientCall delegate) {
+ super(delegate);
+ }
+
+ @Override
+ public void start(Listener responseListener, Metadata headers) {
+ capturedListenerRef.set(responseListener);
+ super.start(responseListener, headers);
+ }
+
+ @Override
+ public boolean isReady() {
+ return transportReady.get();
+ }
+
+ @Override
+ public void sendMessage(ReqT message) {
+ sendMessageCount.incrementAndGet();
+ super.sendMessage(message);
+ }
+
+ @Override
+ public void halfClose() {
+ halfCloseCallCount.incrementAndGet();
+ super.halfClose();
+ }
+ }
+
+ ManagedChannel dataPlaneChannel =
+ grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName)
+ .intercept(
+ new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method,
+ CallOptions callOptions,
+ Channel next) {
+ return new FailOpenDelayedHalfCloseForwardingCall<>(
+ next.newCall(method, callOptions));
+ }
+ })
+ .directExecutor()
+ .build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_CLIENT_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()), dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1);
+
+ // Call halfClose immediately. This sets pendingHalfClose = true.
+ proxyCall.halfClose();
+
+ assertThat(sidecarLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // Since transportReady is false, the mutated body is queued.
+ // And since it was unilateral half-close, pendingUpstreamHalfClose is set to true.
+ // Verify that the call is NOT half-closed on transport yet.
+ assertThat(halfCloseCallCount.get()).isEqualTo(0);
+
+ // Fail the ext_proc stream to trigger fail-open.
+ responseObserverRef.get().onError(Status.INTERNAL.asRuntimeException());
+
+ // Fail-open will see pendingHalfClose = true, set requestSideClosed = true,
+ // and immediately invoke transport.halfClose()
+ assertThat(halfCloseCallCount.get()).isEqualTo(1);
+
+ // Now make transport ready and trigger onReady callback
+ transportReady.set(true);
+ ClientCall.Listener> listener = capturedListenerRef.get();
+ assertThat(listener).isNotNull();
+ listener.onReady();
+
+ // The queued message should be drained, forwarded to backend server call.
+ // The delayed half-close is triggered, but since requestSideClosed was already true (via
+ // fail-open),
+ // it will evaluate to false in compareAndSet and NOT invoke transport.halfClose() again.
+ assertThat(sendMessageCount.get()).isEqualTo(1);
+ // Verify that transport.halfClose() was still called exactly once
+ assertThat(halfCloseCallCount.get()).isEqualTo(1);
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToDownstreamFlowControl_QueuingAndWithholdingWindowUpdates()
+ throws Exception {
+ ExternalProcessor proto =
+ ExternalProcessor.newBuilder()
+ .setGrpcService(
+ GrpcService.newBuilder()
+ .setGoogleGrpc(
+ GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(
+ Any.newBuilder().setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL).build())
+ .build())
+ .build())
+ .setProcessingMode(
+ ProcessingMode.newBuilder()
+ .setRequestHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setRequestBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests =
+ Collections.synchronizedList(new ArrayList<>());
+ final CountDownLatch headersLatch = new CountDownLatch(1);
+ final CountDownLatch firstBodyResponseLatch = new CountDownLatch(1);
+ final CountDownLatch secondBodyResponseLatch = new CountDownLatch(1);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ headersLatch.countDown();
+ } else if (request.hasResponseBody()) {
+ ByteString body = request.getResponseBody().getBody();
+ boolean eos = request.getResponseBody().getEndOfStream();
+ if (body.size() == 40001) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setResponseBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(body)
+ .setEndOfStream(eos)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ firstBodyResponseLatch.countDown();
+
+ // Send second body (40002) - spoofed
+ ByteString body2 =
+ ByteString.copyFromUtf8(new String(new char[40002]).replace('\0', 'y'));
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setResponseBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(body2)
+ .setEndOfStream(eos)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ secondBodyResponseLatch.countDown();
+ }
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager =
+ new CachedChannelManager(
+ config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName)
+ .directExecutor()
+ .build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor =
+ new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final AtomicReference> dataPlaneResponseObserverRef =
+ new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_BIDI_STREAMING,
+ ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(
+ StreamObserver responseObserver) {
+ dataPlaneResponseObserverRef.set(responseObserver);
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()), dataPlaneChannel);
+
+ final List receivedResponses = Collections.synchronizedList(new ArrayList<>());
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ receivedResponses.add(message);
+ }
+ }, new Metadata());
+
+ // Wait for the headers handshake to complete and activate the call
+ assertThat(headersLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ proxyCall.request(1);
+
+ String resp1 = new String(new char[40001]).replace('\0', 'x');
+ dataPlaneResponseObserverRef.get().onNext(resp1);
+
+ assertThat(firstBodyResponseLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(secondBodyResponseLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ assertThat(receivedResponses).containsExactly(resp1);
+
+ List windowUpdates = new ArrayList<>();
+ for (ProcessingRequest req : receivedRequests) {
+ if (req.hasClientWindowUpdate()) {
+ windowUpdates.add(req);
+ }
+ }
+ assertThat(windowUpdates).hasSize(1);
+ assertThat(
+ windowUpdates.get(0).getClientWindowUpdate().getWindowIncrementSidestreamToDownstream())
+ .isEqualTo(40001);
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testSidestreamToDownstreamFlowControl_DrainingAndSendingWindowUpdates()
+ throws Exception {
+ ExternalProcessor proto =
+ ExternalProcessor.newBuilder()
+ .setGrpcService(
+ GrpcService.newBuilder()
+ .setGoogleGrpc(
+ GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(
+ Any.newBuilder().setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL).build())
+ .build())
+ .build())
+ .setProcessingMode(
+ ProcessingMode.newBuilder()
+ .setRequestHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setRequestBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests =
+ Collections.synchronizedList(new ArrayList<>());
+ final CountDownLatch headersLatch = new CountDownLatch(1);
+ final CountDownLatch firstBodyResponseLatch = new CountDownLatch(1);
+ final CountDownLatch secondBodyResponseLatch = new CountDownLatch(1);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ headersLatch.countDown();
+ } else if (request.hasResponseBody()) {
+ ByteString body = request.getResponseBody().getBody();
+ boolean eos = request.getResponseBody().getEndOfStream();
+ if (body.size() == 40001) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setResponseBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(body)
+ .setEndOfStream(eos)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ firstBodyResponseLatch.countDown();
+
+ // Send second body (40002) - spoofed
+ ByteString body2 =
+ ByteString.copyFromUtf8(new String(new char[40002]).replace('\0', 'y'));
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setResponseBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(body2)
+ .setEndOfStream(eos)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ secondBodyResponseLatch.countDown();
+ }
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(
+ InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build()
+ .start());
+
+ CachedChannelManager channelManager =
+ new CachedChannelManager(
+ config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName)
+ .directExecutor()
+ .build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor =
+ new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final AtomicReference> dataPlaneResponseObserverRef =
+ new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_BIDI_STREAMING, ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ dataPlaneResponseObserverRef.set(responseObserver);
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ })).build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()), dataPlaneChannel);
+
+ final List receivedResponses = Collections.synchronizedList(new ArrayList<>());
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ receivedResponses.add(message);
+ }
+ }, new Metadata());
+
+ // Wait for the headers handshake to complete and activate the call
+ assertThat(headersLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ proxyCall.request(1);
+
+ String resp1 = new String(new char[40001]).replace('\0', 'x');
+ String resp2 = new String(new char[40002]).replace('\0', 'y');
+ dataPlaneResponseObserverRef.get().onNext(resp1);
+
+ assertThat(firstBodyResponseLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(secondBodyResponseLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ assertThat(receivedResponses).containsExactly(resp1);
+
+ List windowUpdates = new ArrayList<>();
+ for (ProcessingRequest req : receivedRequests) {
+ if (req.hasClientWindowUpdate()) {
+ windowUpdates.add(req);
+ }
+ }
+ assertThat(windowUpdates).hasSize(1);
+ assertThat(
+ windowUpdates.get(0).getClientWindowUpdate().getWindowIncrementSidestreamToDownstream())
+ .isEqualTo(40001);
+
+ proxyCall.request(1);
+ assertThat(receivedResponses).containsExactly(resp1, resp2);
+
+ windowUpdates.clear();
+ for (ProcessingRequest req : receivedRequests) {
+ if (req.hasClientWindowUpdate()) {
+ windowUpdates.add(req);
+ }
+ }
+ assertThat(windowUpdates).hasSize(2);
+ assertThat(
+ windowUpdates.get(1).getClientWindowUpdate().getWindowIncrementSidestreamToDownstream())
+ .isEqualTo(40002);
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testThresholdBasedWindowUpdates() throws Exception {
+ ExternalProcessor proto = ExternalProcessor.newBuilder()
+ .setGrpcService(GrpcService.newBuilder()
+ .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
+ .setTargetUri("in-process:///" + extProcServerName)
+ .addChannelCredentialsPlugin(Any.newBuilder()
+ .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
+ .build())
+ .build())
+ .build())
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests = new CopyOnWriteArrayList<>();
+ final CountDownLatch sidecarLatch = new CountDownLatch(1);
+ final List> observers = new ArrayList<>();
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ observers.add(responseObserver);
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ sidecarLatch.countDown();
+ } else if (request.hasRequestBody()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(
+ request.getRequestBody().getBody())
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_BIDI_STREAMING, ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
+ }
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING, DEFAULT_CALL_OPTIONS, dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+
+ assertThat(sidecarLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(filterClientRequests(receivedRequests)).hasSize(1);
+ assertThat(filterClientRequests(receivedRequests).get(0).hasRequestHeaders()).isTrue();
+
+ // 1. Send body message. It should be sent immediately.
+ proxyCall.sendMessage("Msg 1"); // size = 5 bytes
+
+ assertThat(filterClientRequests(receivedRequests)).hasSize(2);
+ assertThat(filterClientRequests(receivedRequests).get(1).hasRequestBody()).isTrue();
+ assertThat(filterClientRequests(receivedRequests).get(1)
+ .getRequestBody().getBody().toStringUtf8())
+
+ .isEqualTo("Msg 1");
+ // No window updates were accumulated yet.
+ assertThat(filterClientRequests(receivedRequests).get(1).hasClientWindowUpdate()).isFalse();
+
+ // 2. Trigger window replenishment below threshold (e.g. 5 bytes from Msg 1 response).
+ // The interceptor processes the response, forwards it upstream, and increments
+ // accumulatedWindowUpdateSidestreamToUpstream. Since 5 < 32768, it won't send
+ // standalone updates.
+ // We send another message "Msg 2" to trigger piggybacking.
+ proxyCall.sendMessage("Msg 2");
+
+ assertThat(filterClientRequests(receivedRequests)).hasSize(3);
+ assertThat(filterClientRequests(receivedRequests).get(2).hasRequestBody()).isTrue();
+ assertThat(filterClientRequests(receivedRequests).get(2)
+ .getRequestBody().getBody().toStringUtf8())
+
+ .isEqualTo("Msg 2");
+ // Verify accumulated 5 bytes update is piggybacked.
+ assertThat(filterClientRequests(receivedRequests).get(2).hasClientWindowUpdate()).isTrue();
+ assertThat(filterClientRequests(receivedRequests).get(2)
+ .getClientWindowUpdate().getWindowIncrementSidestreamToUpstream())
+
+ .isEqualTo(5);
+
+ // 3. Accumulate past threshold (e.g. 35,000 bytes) without sending body messages.
+ // This should trigger an immediate standalone ClientWindowUpdate.
+ StreamObserver responseObserver = observers.get(0);
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(ByteString.copyFrom(new byte[35000]))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+
+ // standalone client window update received.
+ assertThat(filterClientRequests(receivedRequests)).hasSize(4);
+ assertThat(filterClientRequests(receivedRequests).get(3).hasClientWindowUpdate()).isTrue();
+ assertThat(
+ filterClientRequests(receivedRequests)
+ .get(3)
+ .getClientWindowUpdate()
+ .getWindowIncrementSidestreamToUpstream())
+ .isEqualTo(35005);
+ assertThat(filterClientRequests(receivedRequests).get(3).hasRequestBody()).isFalse();
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testWindowUpdateWithheldWhenUpstreamCapacityExistsAndBelowThreshold()
+ throws Exception {
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ ExternalProcessor proto = createBaseProto(uniqueExtProcServerName)
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests = new CopyOnWriteArrayList<>();
+ final CountDownLatch sidecarLatch = new CountDownLatch(2);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ sidecarLatch.countDown();
+ } else if (request.hasRequestBody()) {
+ // Mutate request body and send back 10000 bytes (below threshold 32768)
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(
+ ByteString.copyFrom(new byte[10000]))
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ sidecarLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
+
+ CachedChannelManager channelManager = new CachedChannelManager(config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_CLIENT_STREAMING, ServerCalls.asyncClientStreamingCall(
+ new ServerCalls.ClientStreamingMethod() {
+ @Override
+ public StreamObserver invoke(StreamObserver responseObserver) {
+ return new StreamObserver() {
+ @Override
+ public void onNext(String value) {}
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onNext("Response");
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ }))
+ .build());
+
+ ManagedChannel dataPlaneChannel =
+ grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientCall proxyCall =
+ interceptCall(
+ interceptor,
+ METHOD_CLIENT_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
+
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+
+ // Send 10k message to ext_proc.
+ String body10k = new String(new char[10000]).replace('\0', 'a');
+ proxyCall.sendMessage(body10k);
+
+ assertThat(sidecarLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // Since the window has capacity (65536 - 10000 = 55536 > 0) and the increment (10000)
+ // is below the threshold, NO window update should be sent.
+ // receivedRequests should only contain Headers and RequestBody (size = 2).
+ assertThat(receivedRequests).hasSize(2);
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testWindowUpdateWithheldWhenDownstreamCapacityExistsAndBelowThreshold()
+ throws Exception {
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ ExternalProcessor proto =
+ createBaseProto(uniqueExtProcServerName)
+ .setProcessingMode(
+ ProcessingMode.newBuilder()
+ .setRequestHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setRequestBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final List receivedRequests =
+ Collections.synchronizedList(new ArrayList<>());
+ final CountDownLatch headersLatch = new CountDownLatch(1);
+ final CountDownLatch bodyResponseLatch = new CountDownLatch(1);
+
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ @Override
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ receivedRequests.add(request);
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ headersLatch.countDown();
+ } else if (request.hasResponseBody()) {
+ boolean eos = request.getResponseBody().getEndOfStream();
+ // Mutate response body and send back 10000 bytes (below threshold 32768)
+ responseObserver.onNext(
+ ProcessingResponse.newBuilder()
+ .setResponseBody(
+ BodyResponse.newBuilder()
+ .setResponse(
+ CommonResponse.newBuilder()
+ .setBodyMutation(
+ BodyMutation.newBuilder()
+ .setStreamedResponse(
+ StreamedBodyResponse.newBuilder()
+ .setBody(
+ ByteString.copyFrom(new byte[10000]))
+ .setEndOfStream(eos)
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ bodyResponseLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
+ };
+
+ grpcCleanup.register(
+ InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build()
+ .start());
+
+ CachedChannelManager channelManager =
+ new CachedChannelManager(
+ config -> {
+ return grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueExtProcServerName)
+ .directExecutor()
+ .build());
+ });
+
+ ExternalProcessorClientInterceptor interceptor =
+ new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+
+ final AtomicReference> dataPlaneResponseObserverRef =
+ new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(
+ METHOD_BIDI_STREAMING,
+ ServerCalls.asyncBidiStreamingCall(
+ new ServerCalls.BidiStreamingMethod() {
+ @Override
+ public StreamObserver