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..4dc49ccf1bf 100644
--- a/xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java
+++ b/xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java
@@ -77,6 +77,7 @@
import io.grpc.xds.internal.headermutations.HeaderMutator;
import java.io.IOException;
import java.io.InputStream;
+import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import java.util.Queue;
@@ -89,6 +90,7 @@
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import javax.annotation.Nullable;
+import javax.annotation.concurrent.GuardedBy;
/**
* Client-side interceptor for external processing filter.
@@ -186,11 +188,6 @@ ExternalProcessorFilterConfig getFilterConfig() {
return filterConfig;
}
- @VisibleForTesting
- ManagedChannel getExtProcChannel() {
- return extProcChannel;
- }
-
@Override
@SuppressWarnings("unchecked")
public ClientCall interceptCall(
@@ -232,7 +229,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 =
@@ -278,11 +274,12 @@ private static class DataPlaneClientCall
private final ClientCall rawCall;
private final DataPlaneDelayedCall delayedCall;
private final ScheduledExecutorService scheduler;
- private final Object streamLock = new Object();
+ final Object streamLock = new Object();
@Nullable private volatile EventType expectedRequestResponse;
@Nullable private volatile EventType expectedResponseResponse;
@Nullable private volatile ClientCallStreamObserver
extProcClientCallRequestObserver;
+ @GuardedBy("streamLock")
private final Queue pendingDrainingMessages =
new ConcurrentLinkedQueue<>();
@Nullable private volatile DataPlaneListener wrappedListener;
@@ -290,6 +287,53 @@ 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
+ @GuardedBy("streamLock")
+ private long downstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+ @GuardedBy("streamLock")
+ private long upstreamToSidestreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+
+ // Inbound (receiving) windows
+ @GuardedBy("streamLock")
+ private long sidestreamToUpstreamWindow = DEFAULT_INITIAL_WINDOW_SIZE;
+ @GuardedBy("streamLock")
+ 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
+ @GuardedBy("streamLock")
+ private final Queue pendingRequestBodyMessages = new ConcurrentLinkedQueue<>();
+ // Deferred half-close flag for upstream direction
+ private final AtomicBoolean pendingUpstreamHalfClose = new AtomicBoolean(false);
+
+ // Path 2: Buffered request body messages from ext_proc server to forward upstream
+ @GuardedBy("streamLock")
+ private final Queue pendingUpstreamBodyMessages =
+ new java.util.concurrent.ConcurrentLinkedQueue<>();
+ // Path 4: Outstanding requests from downstream for pulling responses
+ @GuardedBy("streamLock")
+ private int downstreamRequestsPending = 0;
+ // Buffered mutated response bodies from ext_proc server
+ @GuardedBy("streamLock")
+ private final Queue pendingMutatedResponseBodies =
+ new java.util.concurrent.ConcurrentLinkedQueue<>();
+
+ // Accumulated client window updates to send to ext_proc
+ @GuardedBy("streamLock")
+ private long accumulatedWindowUpdateSidestreamToUpstream = 0;
+ @GuardedBy("streamLock")
+ private long accumulatedWindowUpdateSidestreamToDownstream = 0;
+
+ // Flag to track if FlowControlInit was sent in the initial message
+ @GuardedBy("streamLock")
+ private boolean flowControlInitSent = false;
+
private final MethodDescriptor, ?> method;
private final Channel channel;
private final MetricRecorder metricsRecorder;
@@ -343,8 +387,6 @@ protected DataPlaneClientCall(
this.backendService = checkNotNull(backendService, "backendService");
}
-
-
private void activateCall() {
if ((extProcStreamState.get() == ExtProcStreamState.FAILED
&& !config.getFailureModeAllow()
@@ -411,8 +453,6 @@ private boolean validateCompressionSupport(BodyResponse bodyResponse) {
return true;
}
-
-
@Override
public void start(Listener responseListener, Metadata headers) {
this.callContext = Context.current();
@@ -439,6 +479,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 +726,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);
}
}
+ @GuardedBy("streamLock")
+ void mergeAccumulatedWindowUpdates(ProcessingRequest.Builder requestBuilder) {
+ 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
+ || extProcStreamState.get().isCompleted()) {
+ 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 +862,39 @@ private void onReadyNotify() {
wrappedListener.onReadyNotify();
}
- private boolean isSidecarReady() {
+ void onReady() {
+ boolean isPassThrough;
+ boolean isCompleted;
+ boolean isDraining;
+
+ synchronized (streamLock) {
+ isPassThrough = passThroughMode.get();
+ ExtProcStreamState state = extProcStreamState.get();
+ isCompleted = state.isCompleted();
+ isDraining = state.isDraining();
+ }
+
+ if (isPassThrough) {
+ onReadyNotify();
+ return;
+ }
+
+ if (isCompleted) {
+ drainPendingDrainingMessages();
+ return;
+ }
+
+ // Normal or Draining operation
+ drainPendingUpstreamBodyMessages();
+ if (!isDraining) {
+ trySendAccumulatedWindowUpdates();
+ }
+ drainPendingRequests();
+ onReadyNotify();
+ }
+
+ @GuardedBy("streamLock")
+ boolean isSidecarReady() {
ExtProcStreamState state = extProcStreamState.get();
if (state.isCompleted()) {
return true;
@@ -738,10 +902,8 @@ private boolean isSidecarReady() {
if (state.isDraining()) {
return false;
}
- synchronized (streamLock) {
- ClientCallStreamObserver observer = extProcClientCallRequestObserver;
- return observer != null && observer.isReady();
- }
+ ClientCallStreamObserver observer = extProcClientCallRequestObserver;
+ return observer != null && observer.isReady();
}
@Override
@@ -755,11 +917,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 +933,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 +1004,58 @@ 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()
- .setRequestBody(HttpBody.newBuilder()
- .setBody(bodyByteString)
- .setEndOfStream(false)
- .build())
- .build());
- bodyMessageSentToExtProc.set(true);
+ @GuardedBy("streamLock")
+ private void sendRequestBodyToExtProc(ByteString body) {
+ downstreamToSidestreamWindow -= body.size();
+ ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()
+ .setRequestBody(HttpBody.newBuilder()
+ .setBody(body)
+ .setEndOfStream(false)
+ .build());
+ mergeAccumulatedWindowUpdates(builder);
+ sendToExtProc(builder.build());
+ bodyMessageSentToExtProc.set(true);
+ }
- if (config.getObservabilityMode()) {
- super.sendMessage(new KnownLengthInputStream(bodyByteString));
- }
- } catch (IOException e) {
- rawCall.cancel("Failed to serialize message for External Processor", e);
+ @GuardedBy("streamLock")
+ private void drainPendingRequestBodyMessages() {
+ while (downstreamToSidestreamWindow > 0 && !pendingRequestBodyMessages.isEmpty()) {
+ ByteString body = pendingRequestBodyMessages.poll();
+ sendRequestBodyToExtProc(body);
+ }
+ if (pendingRequestBodyMessages.isEmpty() && pendingHalfClose.get()) {
+ halfClose();
}
}
@@ -887,11 +1108,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 +1142,31 @@ 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) {
+ sidestreamToUpstreamWindow -= body.size();
+ if (pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
+ 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 +1179,102 @@ 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) {
+ sidestreamToDownstreamWindow -= bodySize;
+ }
+ deliverResponseBody(body, listener);
+ }
+ }
+ }
+
+ private void deliverResponseBody(ByteString body, DataPlaneListener listener) {
+ boolean shouldDeliver = false;
+ synchronized (streamLock) {
+ if (downstreamRequestsPending > 0) {
+ downstreamRequestsPending--;
+ shouldDeliver = true;
+ } else {
+ pendingMutatedResponseBodies.add(body);
+ }
+ }
+ if (shouldDeliver) {
+ final int bodySize = body.size();
+ callContext.run(() -> {
+ try {
+ listener.onExternalBody(body);
+ } finally {
+ synchronized (streamLock) {
+ accumulatedWindowUpdateSidestreamToDownstream += bodySize;
+ }
+ trySendAccumulatedWindowUpdates();
+ }
+ });
+ }
+ }
+
+ private void drainPendingMutatedResponseBodies() {
+ List toDeliver = new ArrayList<>();
+ synchronized (streamLock) {
+ while (downstreamRequestsPending > 0 && !pendingMutatedResponseBodies.isEmpty()) {
+ ByteString body = pendingMutatedResponseBodies.poll();
+ downstreamRequestsPending--;
+ pendingRequests.decrementAndGet();
+ toDeliver.add(body);
+ }
+ }
+ for (ByteString body : toDeliver) {
+ final int bodySize = body.size();
+ callContext.run(() -> {
+ try {
+ wrappedListener.onExternalBody(body);
+ } finally {
+ synchronized (streamLock) {
+ accumulatedWindowUpdateSidestreamToDownstream += bodySize;
+ }
+ trySendAccumulatedWindowUpdates();
+ }
+ });
+ }
+ }
+
+ void drainPendingMutatedResponseBodiesDirect(DataPlaneListener listener) {
+ List toDeliver = new ArrayList<>();
+ synchronized (streamLock) {
+ ByteString body;
+ while ((body = pendingMutatedResponseBodies.poll()) != null) {
+ toDeliver.add(body);
+ }
+ }
+ for (ByteString body : toDeliver) {
+ listener.onExternalBody(body);
+ }
+ }
+
+ void drainPendingUpstreamBodyMessages() {
+ while (true) {
+ ByteString body = null;
+ boolean triggerHalfClose = false;
+ synchronized (streamLock) {
+ if (!pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
+ 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();
+ }
}
}
}
@@ -965,16 +1308,51 @@ private void handleImmediateResponse(ImmediateResponse immediate, DataPlaneListe
}
private void drainPendingDrainingMessages() {
- synchronized (streamLock) {
- InputStream msg;
- while ((msg = pendingDrainingMessages.poll()) != null) {
- super.sendMessage(msg);
+ while (true) {
+ Object msg = null; // Can be ByteString or InputStream
+ boolean isMutated = false;
+ boolean triggerHalfClose = false;
+
+ synchronized (streamLock) {
+ if (!pendingUpstreamBodyMessages.isEmpty() && super.isReady()) {
+ msg = pendingUpstreamBodyMessages.poll();
+ isMutated = true;
+ } else if (pendingUpstreamBodyMessages.isEmpty()
+ && !pendingRequestBodyMessages.isEmpty() && super.isReady()) {
+ msg = pendingRequestBodyMessages.poll();
+ isMutated = true;
+ } else if (pendingUpstreamBodyMessages.isEmpty()
+ && pendingRequestBodyMessages.isEmpty()
+ && !pendingDrainingMessages.isEmpty() && super.isReady()) {
+ msg = pendingDrainingMessages.poll();
+ isMutated = false;
+ }
+
+ if (msg == null) {
+ if (pendingUpstreamBodyMessages.isEmpty()
+ && pendingRequestBodyMessages.isEmpty()
+ && pendingDrainingMessages.isEmpty()) {
+ passThroughMode.set(true);
+ if (pendingHalfClose.get()) {
+ triggerHalfClose = true;
+ }
+ }
+ }
}
- passThroughMode.set(true);
- if (pendingHalfClose.get()) {
- if (requestSideClosed.compareAndSet(false, true)) {
- proceedWithHalfClose();
+
+ if (msg == null) {
+ if (triggerHalfClose) {
+ if (requestSideClosed.compareAndSet(false, true)) {
+ proceedWithHalfClose();
+ }
}
+ break;
+ }
+
+ if (isMutated) {
+ super.sendMessage(new KnownLengthInputStream((ByteString) msg));
+ } else {
+ super.sendMessage((InputStream) msg);
}
}
}
@@ -1048,6 +1426,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,8 +1465,7 @@ void setImmediateResponse(Status status, Metadata trailers) {
@Override
public void onReady() {
- dataPlaneClientCall.drainPendingRequests();
- onReadyNotify();
+ dataPlaneClientCall.onReady();
}
@Override
@@ -1104,7 +1483,7 @@ public void onHeaders(Metadata headers) {
return;
}
- if (dataPlaneClientCall.getPassThroughMode().get()
+ if (dataPlaneClientCall.getPassThroughMode().get()
|| dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
|| !sendResponseHeaders) {
proceedWithHeaders(headers);
@@ -1126,7 +1505,7 @@ public void onHeaders(Metadata headers) {
@Override
public void onMessage(InputStream message) {
- synchronized (savedMessages) {
+ synchronized (dataPlaneClientCall.streamLock) {
if (inboundPassThrough) {
dataPlaneClientCall.getCallContext().run(() -> delegate().onMessage(message));
return;
@@ -1145,34 +1524,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.streamLock) {
+ 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 +1636,7 @@ void onReadyNotify() {
void proceedWithHeaders() {
if (savedHeaders != null) {
proceedWithHeaders(savedHeaders);
- synchronized (savedMessages) {
+ synchronized (dataPlaneClientCall.streamLock) {
savedHeaders = null;
if (!dataPlaneClientCall.getExtProcStreamState().get().isDraining()) {
InputStream msg;
@@ -1285,13 +1690,17 @@ void onExternalBody(ByteString body) {
void unblockAfterStreamComplete() {
proceedWithHeaders();
+ // 1. Drain mutated responses first
+ dataPlaneClientCall.drainPendingMutatedResponseBodiesDirect(this);
+ // 2. Drain raw responses
proceedWithSavedMessages();
+ // 3. Drain outbound requests
dataPlaneClientCall.drainPendingDrainingMessages();
proceedWithClose();
}
private void proceedWithSavedMessages() {
- synchronized (savedMessages) {
+ synchronized (dataPlaneClientCall.streamLock) {
InputStream msg;
while ((msg = savedMessages.poll()) != null) {
final InputStream finalMsg = msg;
@@ -1361,6 +1770,7 @@ private void triggerCloseHandshake() {
}
}
+ @GuardedBy("dataPlaneClientCall.streamLock")
private void sendResponseBodyToExtProc(
@Nullable ByteString bodyByteString, boolean endOfStream) {
if (dataPlaneClientCall.getExtProcStreamState().get().isCompleted()
@@ -1376,9 +1786,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..096dc636140 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");
}
@@ -274,7 +276,6 @@ private ExternalProcessor.Builder createBaseProto(String targetName) {
.build());
}
-
// --- Category 1: Configuration Override ---
@Test
@@ -3461,7 +3462,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 +3469,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 +3618,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 +5942,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 +6265,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 +6763,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 +6887,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 +7007,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 +7111,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 +7165,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 +7270,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 +7505,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 +7517,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();
@@ -8378,182 +8362,158 @@ public void onClose(Status status, Metadata trailers) {
channelManager.close();
}
- // --- Category 15: Inbound Backpressure (request(n) / pendingRequests) ---
+ // --- Category 15: Ext-proc fail-open draining of flow-control queues ---
@Test
@SuppressWarnings("unchecked")
- public void givenObservabilityTrue_whenExtProcBusy_thenAppRequestsBuffered()
- throws Exception {
- 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())
+ public void testFailOpen_DrainsInboundQueuesInOrder() throws Exception {
+ ExternalProcessor proto = createBaseProto(extProcServerName)
+ .setFailureModeAllow(true)
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
.build())
- .setObservabilityMode(true)
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
assertThat(configOrError.errorDetail).isNull();
ExternalProcessorFilterConfig filterConfig = configOrError.config;
- // External Processor Server
- ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl;
- extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() {
- @Override
- @SuppressWarnings("unchecked")
- public StreamObserver process(
- StreamObserver responseObserver) {
- ((ServerCallStreamObserver) responseObserver).request(100);
- return new StreamObserver() {
- @Override
- public void onNext(ProcessingRequest request) {
- }
+ final CountDownLatch extProcReceivedHeadersLatch = new CountDownLatch(1);
+ final AtomicReference> responseObserverRef =
+ new AtomicReference<>();
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
- public void onError(Throwable t) {
- }
+ 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());
+ extProcReceivedHeadersLatch.countDown();
+ }
+ }
- @Override
- public void onCompleted() {
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
}
};
- }
- };
- grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
.addService(extProcImpl)
.directExecutor()
.build().start());
- final AtomicBoolean sidecarReady = new AtomicBoolean(true);
- final AtomicReference> sidecarListenerRef =
- new AtomicReference<>();
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<
- ReqT, RespT>(next.newCall(method, callOptions)) {
- @Override
- public void start(Listener responseListener, Metadata headers) {
- sidecarListenerRef.set((Listener) responseListener);
- super.start(responseListener, headers);
- }
-
- @Override
- public boolean isReady() {
- return sidecarReady.get();
- }
- };
- }
- })
- .build());
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
});
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
- final AtomicInteger dataPlaneRequestCount = new AtomicInteger(0);
+ final CountDownLatch backendSentMessage2Latch = new CountDownLatch(1);
dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, 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) {
- }
+ .addMethod(METHOD_BIDI_STREAMING, (call, headers) -> {
+ call.sendHeaders(new Metadata());
+ // Send message 1 (70k to close window)
+ String largeMessage70k = new String(new char[70000]).replace('\0', 'a');
+ call.sendMessage(largeMessage70k);
- @Override
- public void onCompleted() {
- responseObserver.onCompleted();
- }
- };
+ new Thread(() -> {
+ try {
+ if (extProcReceivedHeadersLatch.await(5, TimeUnit.SECONDS)) {
+ // Send message 2 (unsolicited, will be buffered in savedMessages)
+ call.sendMessage("backend-msg-2");
+ backendSentMessage2Latch.countDown();
}
- }))
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }).start();
+
+ return new ServerCall.Listener() {
+ @Override
+ public void onMessage(String message) {}
+
+ @Override
+ public void onHalfClose() {}
+
+ @Override
+ public void onCancel() {}
+ };
+ })
.build());
+ final List appReceivedMessages = new CopyOnWriteArrayList<>();
+ final CountDownLatch callClosedLatch = new CountDownLatch(1);
+
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());
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
- CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
- interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
- proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ dataPlaneChannel);
- // 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();
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ appReceivedMessages.add(message);
+ }
- // Sidecar is busy
- sidecarReady.set(false);
- assertThat(proxyCall.isReady()).isFalse();
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ callClosedLatch.countDown();
+ }
+ }, new Metadata());
- proxyCall.request(5);
+ proxyCall.request(10);
- // Verify data plane call NOT requested yet (due to observability mode and sidecar busy)
- assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
+ // Wait for backend to send message 2
+ assertThat(backendSentMessage2Latch.await(5, TimeUnit.SECONDS)).isTrue();
- // Sidecar becomes ready
- sidecarReady.set(true);
- sidecarListenerRef.get().onReady();
+ // Verify app received nothing yet (buffered in savedMessages)
+ assertThat(appReceivedMessages).isEmpty();
+
+ // Trigger fail-open by error on ext_proc stream
+ responseObserverRef.get().onError(Status.UNAVAILABLE.asException());
+
+ // Verify call is NOT closed
+ assertThat(callClosedLatch.getCount()).isEqualTo(1);
+
+ // Verify all buffered messages are drained in order: largeMessage70k then backend-msg-2
+ String largeMessage70k = new String(new char[70000]).replace('\0', 'a');
+ assertThat(appReceivedMessages).containsExactly(largeMessage70k, "backend-msg-2").inOrder();
- // After sidecar becomes ready, pending requests should be drained to data plane.
- assertThat(dataPlaneRequestCount.get()).isEqualTo(5);
- assertThat(proxyCall.isReady()).isTrue();
-
proxyCall.cancel("Cleanup", null);
channelManager.close();
}
@Test
@SuppressWarnings("unchecked")
- public void givenRequestDrainActive_whenAppRequestsMessages_thenRequestsBuffered()
- throws Exception {
- 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())
+ public void testFailOpen_DrainsBlockedRequests() throws Exception {
+ ExternalProcessor proto = createBaseProto(extProcServerName)
+ .setFailureModeAllow(true)
.setProcessingMode(ProcessingMode.newBuilder()
- .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
- .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
.build())
.build();
ConfigOrError configOrError =
@@ -8561,258 +8521,438 @@ public void givenRequestDrainActive_whenAppRequestsMessages_thenRequestsBuffered
assertThat(configOrError.errorDetail).isNull();
ExternalProcessorFilterConfig filterConfig = configOrError.config;
- // External Processor Server
- 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) {
- if (request.hasRequestHeaders()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestDrain(true)
- .build());
- }
- }
+ final AtomicReference> responseObserverRef =
+ new AtomicReference<>();
+ final CountDownLatch extProcReceivedHeadersLatch = new CountDownLatch(1);
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
- public void onError(Throwable t) {
+ 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()) {
+ // Respond with headers AND negative window update to block outbound body
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .setServerWindowUpdate(ProcessingResponse.ServerWindowUpdate.newBuilder()
+ .setWindowIncrementDownstreamToSidestream(-65536) // Reduce window to 0
+ .build())
+ .build());
+ extProcReceivedHeadersLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
}
+ };
- @Override
- public void onCompleted() {
- }
- };
- }
- };
- grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
.addService(extProcImpl)
.directExecutor()
.build().start());
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(extProcServerName).directExecutor().build());
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
});
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+ final List sentToBackend = new CopyOnWriteArrayList<>();
dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello " + request);
- responseObserver.onCompleted();
- }))
+ .addMethod(METHOD_BIDI_STREAMING, (call, headers) -> {
+ call.sendHeaders(new Metadata());
+ call.request(100);
+ return new ServerCall.Listener() {
+ @Override
+ public void onMessage(String message) {}
+
+ @Override
+ public void onHalfClose() {}
+
+ @Override
+ public void onCancel() {}
+ };
+ })
.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());
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientInterceptor backendInterceptor = new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ ClientCall delegateCall = next.newCall(method, callOptions);
+ return new SimpleForwardingClientCall(delegateCall) {
+ @Override
+ public void sendMessage(ReqT message) {
+ try {
+ InputStream is = (InputStream) message;
+ byte[] bytes = com.google.common.io.ByteStreams.toByteArray(is);
+ String str = new String(bytes, StandardCharsets.UTF_8);
+ sentToBackend.add(str);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ super.sendMessage(message);
+ }
+ };
+ }
+ };
+ Channel interceptedChannel =
+ ClientInterceptors.intercept(dataPlaneChannel, backendInterceptor);
- CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
- interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ interceptedChannel);
+
proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1);
- // Wait for drain to be processed
- long startTime = System.currentTimeMillis();
- while (proxyCall.isReady() && System.currentTimeMillis() - startTime < 5000) {
- Thread.sleep(10);
- }
- assertThat(proxyCall.isReady()).isFalse();
+ assertThat(extProcReceivedHeadersLatch.await(5, TimeUnit.SECONDS)).isTrue();
- // App requests more messages
- proxyCall.request(3);
+ // These should now be buffered because window is 0
+ proxyCall.sendMessage("msg-1");
+ proxyCall.sendMessage("msg-2");
+
+ assertThat(sentToBackend).isEmpty();
+
+ // Trigger fail-open
+ responseObserverRef.get().onError(Status.UNAVAILABLE.asException());
+
+ // Verify messages are drained
+ assertThat(sentToBackend).containsExactly("msg-1", "msg-2").inOrder();
- // Verify requests are buffered and not sent to data plane
- assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
- // proxyCall.isReady() should remain false during drain
- assertThat(proxyCall.isReady()).isFalse();
-
proxyCall.cancel("Cleanup", null);
channelManager.close();
}
@Test
@SuppressWarnings("unchecked")
- public void givenBufferedRequests_whenExtProcStreamBecomesReady_thenDataPlaneDrained()
- throws Exception {
- 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())
+ public void testFailOpen_DrainsDrainingRequests() throws Exception {
+ ExternalProcessor proto = createBaseProto(extProcServerName)
+ .setFailureModeAllow(true)
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
.build())
- .setObservabilityMode(true)
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
assertThat(configOrError.errorDetail).isNull();
ExternalProcessorFilterConfig filterConfig = configOrError.config;
- // External Processor Server
- 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) {
- if (request.hasRequestHeaders()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestHeaders(HeadersResponse.newBuilder().build())
- .build());
- }
- }
+ final AtomicReference> responseObserverRef =
+ new AtomicReference<>();
+ final CountDownLatch extProcReceivedHeadersLatch = new CountDownLatch(1);
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
- public void onError(Throwable t) {
- }
+ 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()) {
+ // Respond with headers AND request_drain = true to trigger DRAINING state
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .setRequestDrain(true)
+ .build());
+ extProcReceivedHeadersLatch.countDown();
+ }
+ }
- @Override
- public void onCompleted() {
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void onCompleted() {}
+ };
}
};
- }
- };
- grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
+
+ String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
.addService(extProcImpl)
.directExecutor()
.build().start());
- final AtomicBoolean sidecarReady = new AtomicBoolean(true);
- final AtomicReference> sidecarListenerRef =
- new AtomicReference<>();
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<
- ReqT, RespT>(next.newCall(method, callOptions)) {
- @Override
- public void start(Listener responseListener, Metadata headers) {
- sidecarListenerRef.set((Listener) responseListener);
- super.start(responseListener, headers);
- }
-
- @Override
- public boolean isReady() {
- return sidecarReady.get();
- }
- };
- }
- })
- .build());
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
});
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+ final List sentToBackend = new CopyOnWriteArrayList<>();
dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello " + request);
- responseObserver.onCompleted();
- }))
+ .addMethod(METHOD_BIDI_STREAMING, (call, headers) -> {
+ call.sendHeaders(new Metadata());
+ call.request(100);
+ return new ServerCall.Listener() {
+ @Override
+ public void onMessage(String message) {}
+
+ @Override
+ public void onHalfClose() {}
+
+ @Override
+ public void onCancel() {}
+ };
+ })
.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());
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientInterceptor backendInterceptor = new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ ClientCall delegateCall = next.newCall(method, callOptions);
+ return new SimpleForwardingClientCall(delegateCall) {
+ @Override
+ public void sendMessage(ReqT message) {
+ try {
+ InputStream is = (InputStream) message;
+ byte[] bytes = com.google.common.io.ByteStreams.toByteArray(is);
+ String str = new String(bytes, StandardCharsets.UTF_8);
+ sentToBackend.add(str);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ super.sendMessage(message);
+ }
+ };
+ }
+ };
+ Channel interceptedChannel =
+ ClientInterceptors.intercept(dataPlaneChannel, backendInterceptor);
- CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
- interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ interceptedChannel);
+
proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1);
- // 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();
+ assertThat(extProcReceivedHeadersLatch.await(5, TimeUnit.SECONDS)).isTrue();
- // Sidecar is busy initially
- sidecarReady.set(false);
-
- // Request from application
- proxyCall.request(10);
- assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
+ // Send msg-1. Since state is DRAINING, it should be buffered in pendingDrainingMessages.
+ proxyCall.sendMessage("msg-1");
- // Sidecar becomes ready
- sidecarReady.set(true);
- sidecarListenerRef.get().onReady();
+ assertThat(sentToBackend).isEmpty();
+
+ // Trigger fail-open
+ responseObserverRef.get().onError(Status.UNAVAILABLE.asException());
+
+ // Verify message is drained
+ assertThat(sentToBackend).containsExactly("msg-1");
- // Verify buffered request drained
- assertThat(dataPlaneRequestCount.get()).isEqualTo(10);
- assertThat(proxyCall.isReady()).isTrue();
-
proxyCall.cancel("Cleanup", null);
channelManager.close();
}
@Test
@SuppressWarnings("unchecked")
- public void givenExtProcStreamCompleted_whenAppRequestsMessages_thenRequestsForwarded()
- throws Exception {
- 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())
+ public void testFailOpen_ResumesDrainingOnReady() throws Exception {
+ ExternalProcessor proto = createBaseProto(extProcServerName)
+ .setFailureModeAllow(true)
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+
+ final AtomicReference> responseObserverRef =
+ new AtomicReference<>();
+ final CountDownLatch extProcReceivedHeadersLatch = new CountDownLatch(1);
+
+ 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()) {
+ // Respond with headers AND negative window update to block outbound body
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .setServerWindowUpdate(ProcessingResponse.ServerWindowUpdate.newBuilder()
+ .setWindowIncrementDownstreamToSidestream(-65536) // Reduce window to 0
+ .build())
+ .build());
+ extProcReceivedHeadersLatch.countDown();
+ }
+ }
+
+ @Override
+ public void onError(Throwable t) {}
+
+ @Override
+ public void 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 sentToBackend = new CopyOnWriteArrayList<>();
+ final AtomicBoolean backendReady = new AtomicBoolean(true);
+ final AtomicReference> backendListenerRef = new AtomicReference<>();
+
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_BIDI_STREAMING, (call, headers) -> {
+ call.sendHeaders(new Metadata());
+ call.request(100);
+ return new ServerCall.Listener() {
+ @Override
+ public void onMessage(String message) {}
+
+ @Override
+ public void onHalfClose() {}
+
+ @Override
+ public void onCancel() {}
+ };
+ })
+ .build());
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ ClientInterceptor backendInterceptor = new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ ClientCall delegateCall = next.newCall(method, callOptions);
+ return new SimpleForwardingClientCall(delegateCall) {
+ @Override
+ public void start(ClientCall.Listener responseListener, Metadata headers) {
+ backendListenerRef.set((ClientCall.Listener) responseListener);
+ super.start(responseListener, headers);
+ }
+
+ @Override
+ public void sendMessage(ReqT message) {
+ try {
+ InputStream is = (InputStream) message;
+ byte[] bytes = com.google.common.io.ByteStreams.toByteArray(is);
+ String str = new String(bytes, StandardCharsets.UTF_8);
+ sentToBackend.add(str);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ super.sendMessage(message);
+ }
+
+ @Override
+ public boolean isReady() {
+ return backendReady.get();
+ }
+ };
+ }
+ };
+ Channel interceptedChannel =
+ ClientInterceptors.intercept(dataPlaneChannel, backendInterceptor);
+
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
+ interceptedChannel);
+
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
+ proxyCall.request(1);
+
+ assertThat(extProcReceivedHeadersLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ // 1. Backend not ready
+ backendReady.set(false);
+
+ // 2. Send msg-1, msg-2 (buffered in pendingRequestBodyMessages)
+ proxyCall.sendMessage("msg-1");
+ proxyCall.sendMessage("msg-2");
+
+ assertThat(sentToBackend).isEmpty();
+
+ // 3. Trigger fail-open while backend is NOT ready
+ responseObserverRef.get().onError(Status.UNAVAILABLE.asException());
+
+ // Verify still nothing sent
+ assertThat(sentToBackend).isEmpty();
+
+ // 4. Make backend ready and trigger onReady
+ backendReady.set(true);
+ backendListenerRef.get().onReady();
+
+ // Verify messages are drained
+ assertThat(sentToBackend).containsExactly("msg-1", "msg-2").inOrder();
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
+ }
+
+ // --- Category 16: Inbound Backpressure (request(n) / pendingRequests) ---
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void givenObservabilityTrue_whenExtProcBusy_thenAppRequestsBuffered()
+ throws Exception {
+ 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())
+ .setObservabilityMode(true)
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
@@ -8825,15 +8965,11 @@ public void givenExtProcStreamCompleted_whenAppRequestsMessages_thenRequestsForw
@Override
@SuppressWarnings("unchecked")
public StreamObserver process(
- final StreamObserver responseObserver) {
+ StreamObserver responseObserver) {
((ServerCallStreamObserver) responseObserver).request(100);
return new StreamObserver() {
@Override
public void onNext(ProcessingRequest request) {
- if (request.hasRequestHeaders()) {
- // Immediately complete the stream from server side
- responseObserver.onCompleted();
- }
}
@Override
@@ -8851,23 +8987,62 @@ public void onCompleted() {
.directExecutor()
.build().start());
+ final AtomicBoolean sidecarReady = new AtomicBoolean(true);
+ final AtomicReference> sidecarListenerRef =
+ new AtomicReference<>();
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(extProcServerName).directExecutor().build());
+ InProcessChannelBuilder.forName(extProcServerName)
+ .directExecutor()
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new io.grpc.ForwardingClientCall.SimpleForwardingClientCall<
+ ReqT, RespT>(next.newCall(method, callOptions)) {
+ @Override
+ public void start(Listener responseListener, Metadata headers) {
+ sidecarListenerRef.set((Listener) responseListener);
+ super.start(responseListener, headers);
+ }
+
+ @Override
+ public boolean isReady() {
+ return sidecarReady.get();
+ }
+ };
+ }
+ })
+ .build());
});
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+ final AtomicInteger dataPlaneRequestCount = new AtomicInteger(0);
dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello " + request);
- responseObserver.onCompleted();
+ .addMethod(METHOD_SAY_HELLO, 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() {
+ responseObserver.onCompleted();
+ }
+ };
+ }
}))
.build());
- final AtomicInteger dataPlaneRequestCount = new AtomicInteger(0);
ManagedChannel dataPlaneChannel = grpcCleanup.register(
InProcessChannelBuilder.forName(dataPlaneServerName)
.directExecutor()
@@ -8887,97 +9062,78 @@ public void request(int numMessages) {
})
.build());
- final CountDownLatch readyLatch = new CountDownLatch(1);
CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
- proxyCall.start(new ClientCall.Listener() {
- @Override
- public void onReady() {
- readyLatch.countDown();
- }
- }, new Metadata());
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
- // Wait for sidecar stream completion
- assertThat(readyLatch.await(5, TimeUnit.SECONDS)).isTrue();
- assertThat(proxyCall.isReady()).isTrue();
+ // Wait for sidecar call to start
+ long startTime = System.currentTimeMillis();
+ while (sidecarListenerRef.get() == null && System.currentTimeMillis() - startTime < 5000) {
+ }
+ assertThat(sidecarListenerRef.get()).isNotNull();
- proxyCall.request(7);
+ // Sidecar is busy
+ sidecarReady.set(false);
+ assertThat(proxyCall.isReady()).isFalse();
- // Verify request forwarded immediately
- assertThat(dataPlaneRequestCount.get()).isEqualTo(7);
- // proxyCall.isReady() should remain true as sidecar is gone
+ proxyCall.request(5);
+
+ // Verify data plane call NOT requested yet (due to observability mode and sidecar busy)
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
+
+ // Sidecar becomes ready
+ sidecarReady.set(true);
+ sidecarListenerRef.get().onReady();
+
+ // After sidecar becomes ready, pending requests should be drained to data plane.
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(5);
assertThat(proxyCall.isReady()).isTrue();
proxyCall.cancel("Cleanup", null);
channelManager.close();
}
- // --- Category 16: Error Handling & Security ---
-
@Test
- @SuppressWarnings("FutureReturnValueIgnored")
- public void givenPendingData_whenImmediateResponseReceived_thenDeliversDataBeforeStatus()
+ @SuppressWarnings("unchecked")
+ public void givenRequestDrainActive_whenAppRequestsMessages_thenRequestsBuffered()
throws Exception {
- final String uniqueExtProcServerName = InProcessServerBuilder.generateName();
- final String uniqueDataPlaneServerName = InProcessServerBuilder.generateName();
- final List appEvents = Collections.synchronizedList(new ArrayList<>());
- final CountDownLatch finishLatch = new CountDownLatch(1);
- final CountDownLatch extProcCompletedLatch = new CountDownLatch(1);
- final ExecutorService sidecarResponseExecutor = Executors.newSingleThreadExecutor();
- final Metadata.Key immediateKey =
- Metadata.Key.of("x-immediate-header", Metadata.ASCII_STRING_MARSHALLER);
- final AtomicReference appTrailers = new AtomicReference<>();
+ 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()
+ .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
+ // External Processor Server
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) {
- if (request.hasRequestHeaders()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestHeaders(HeadersResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- } else if (request.hasResponseHeaders()) {
- try {
- Thread.sleep(500);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setImmediateResponse(ImmediateResponse.newBuilder()
- .setGrpcStatus(
- io.envoyproxy.envoy.service.ext_proc.v3.GrpcStatus.newBuilder()
- .setStatus(Status.UNAUTHENTICATED.getCode().value())
- .build())
- .setDetails("Immediate Auth Failure")
- .setHeaders(
- io.envoyproxy.envoy.service.ext_proc.v3.HeaderMutation.newBuilder()
- .addSetHeaders(
- io.envoyproxy.envoy.config.core.v3.HeaderValueOption
- .newBuilder()
- .setHeader(
- io.envoyproxy.envoy.config.core.v3.HeaderValue
- .newBuilder()
- .setKey("x-immediate-header")
- .setValue("true")
- .build())
- .build())
- .build())
- .build())
- .build());
- responseObserver.onCompleted();
- }
- }
- });
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestDrain(true)
+ .build());
+ }
}
@Override
@@ -8986,292 +9142,221 @@ public void onError(Throwable t) {
@Override
public void onCompleted() {
- extProcCompletedLatch.countDown();
}
};
}
};
-
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
- .addService(extProcImpl).directExecutor().build().start());
+ grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ InProcessChannelBuilder.forName(extProcServerName).directExecutor().build());
});
- ExternalProcessorFilter filter = new ExternalProcessorFilter(FAKE_CONTEXT, channelManager);
- ExternalProcessor proto = createBaseProto(extProcServerName)
- .setProcessingMode(ProcessingMode.newBuilder()
- .setRequestBodyMode(ProcessingMode.BodySendMode.NONE)
- .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
- .build())
- .build();
- ConfigOrError configOrError =
- provider.parseFilterConfig(Any.pack(proto), filterContext);
- ExternalProcessorFilterConfig filterConfig = configOrError.config;
-
- ClientInterceptor interceptor = filter.buildClientInterceptor(filterConfig, null, scheduler);
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
- MutableHandlerRegistry dataPlaneRegistry = new MutableHandlerRegistry();
- dataPlaneRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, (call, headers) -> {
- call.sendHeaders(new Metadata());
- call.request(1);
- return new ServerCall.Listener() {
- @Override
- public void onMessage(String message) {
- call.sendMessage("server-response");
- call.close(Status.OK, new Metadata());
- }
- };
- })
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
+ (request, responseObserver) -> {
+ responseObserver.onNext("Hello " + request);
+ responseObserver.onCompleted();
+ }))
.build());
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueDataPlaneServerName)
- .fallbackHandlerRegistry(dataPlaneRegistry)
- .executor(Executors.newSingleThreadExecutor())
- .build().start());
-
- ManagedChannel channel =
- grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueDataPlaneServerName).directExecutor().build());
- Channel interceptedChannel = io.grpc.ClientInterceptors.interceptForward(
- channel,
- Arrays.asList(new XdsNameResolver.RawMessageClientInterceptor(), interceptor));
-
- ClientCall call =
- interceptedChannel.newCall(
- METHOD_SAY_HELLO,
- DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()));
- call.start(new ClientCall.Listener() {
- @Override
- public void onHeaders(Metadata headers) {
- appEvents.add("HEADERS");
- }
+ 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());
- @Override
- public void onMessage(String message) {
- appEvents.add("MESSAGE");
- }
+ CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
- @Override
- public void onClose(Status status, Metadata trailers) {
- appEvents.add("CLOSE:" + status.getCode());
- appTrailers.set(trailers);
- finishLatch.countDown();
- }
- }, new Metadata());
+ // Wait for drain to be processed
+ long startTime = System.currentTimeMillis();
+ while (proxyCall.isReady() && System.currentTimeMillis() - startTime < 5000) {
+ }
+ assertThat(proxyCall.isReady()).isFalse();
- call.request(1);
- call.sendMessage("request-body");
- call.halfClose();
+ // App requests more messages
+ proxyCall.request(3);
- assertThat(finishLatch.await(5, TimeUnit.SECONDS)).isTrue();
- assertThat(appEvents).containsExactly("HEADERS", "MESSAGE", "CLOSE:UNAUTHENTICATED");
- assertThat(appTrailers.get().get(immediateKey)).isEqualTo("true");
- assertThat(extProcCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ // Verify requests are buffered and not sent to data plane
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
+ // proxyCall.isReady() should remain false during drain
+ assertThat(proxyCall.isReady()).isFalse();
- sidecarResponseExecutor.shutdown();
+ proxyCall.cancel("Cleanup", null);
channelManager.close();
}
-
@Test
- @SuppressWarnings("FutureReturnValueIgnored")
- public void
- givenStreamingCall_whenImmediateResponseReceivedDuringRequestStreaming_thenTerminatesCleanly()
+ @SuppressWarnings("unchecked")
+ public void givenBufferedRequests_whenExtProcStreamBecomesReady_thenDataPlaneDrained()
throws Exception {
- final String uniqueExtProcServerName = InProcessServerBuilder.generateName();
- final String uniqueDataPlaneServerName = InProcessServerBuilder.generateName();
- final List appEvents = Collections.synchronizedList(new ArrayList<>());
- final CountDownLatch finishLatch = new CountDownLatch(1);
- final CountDownLatch extProcCompletedLatch = new CountDownLatch(1);
- final ExecutorService sidecarResponseExecutor = Executors.newSingleThreadExecutor();
- final Metadata.Key immediateKey =
- Metadata.Key.of("x-immediate-header", Metadata.ASCII_STRING_MARSHALLER);
- final AtomicReference appTrailers = new AtomicReference<>();
- final AtomicInteger extProcRequestCount = new AtomicInteger(0);
+ 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())
+ .setObservabilityMode(true)
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ assertThat(configOrError.errorDetail).isNull();
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
- ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
- new ExternalProcessorGrpc.ExternalProcessorImplBase() {
+ // External Processor Server
+ 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 StreamObserver process(
- final StreamObserver responseObserver) {
- ((ServerCallStreamObserver) responseObserver).request(100);
- return new StreamObserver() {
- @Override
- public void onNext(ProcessingRequest request) {
- sidecarResponseExecutor.submit(() -> {
- synchronized (responseObserver) {
- if (request.hasRequestHeaders()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestHeaders(HeadersResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- } else if (request.hasRequestBody()) {
- int count = extProcRequestCount.incrementAndGet();
- if (count == 1) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestBody(BodyResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- } else if (count == 2) {
- try {
- Thread.sleep(500);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setImmediateResponse(ImmediateResponse.newBuilder()
- .setGrpcStatus(
- io.envoyproxy.envoy.service.ext_proc.v3.GrpcStatus.newBuilder()
- .setStatus(Status.UNAUTHENTICATED.getCode().value())
- .build())
- .setDetails("Immediate Auth Failure")
- .setHeaders(
- io.envoyproxy.envoy.service.ext_proc.v3.HeaderMutation
- .newBuilder()
- .addSetHeaders(
- io.envoyproxy.envoy.config.core.v3.HeaderValueOption
- .newBuilder()
- .setHeader(
- io.envoyproxy.envoy.config.core.v3.HeaderValue
- .newBuilder()
- .setKey("x-immediate-header")
- .setValue("true")
- .build())
- .build())
- .build())
- .build())
- .build());
- responseObserver.onCompleted();
- }
- }
- }
- });
- }
+ public void onNext(ProcessingRequest request) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder().build())
+ .build());
+ }
+ }
- @Override
- public void onError(Throwable t) {}
+ @Override
+ public void onError(Throwable t) {
+ }
- @Override
- public void onCompleted() {
- extProcCompletedLatch.countDown();
- }
- };
+ @Override
+ public void onCompleted() {
}
};
+ }
+ };
+ grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
+ .addService(extProcImpl)
+ .directExecutor()
+ .build().start());
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
- .addService(extProcImpl).directExecutor().build().start());
-
+ final AtomicBoolean sidecarReady = new AtomicBoolean(true);
+ final AtomicReference> sidecarListenerRef =
+ new AtomicReference<>();
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
- });
-
- ExternalProcessorFilter filter = new ExternalProcessorFilter(FAKE_CONTEXT, channelManager);
- ExternalProcessor proto = createBaseProto(uniqueExtProcServerName)
- .setProcessingMode(ProcessingMode.newBuilder()
- .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
- .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
- .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
- .build())
- .build();
- ConfigOrError configOrError =
- provider.parseFilterConfig(Any.pack(proto), filterContext);
- ExternalProcessorFilterConfig filterConfig = configOrError.config;
+ InProcessChannelBuilder.forName(extProcServerName)
+ .directExecutor()
+ .intercept(new ClientInterceptor() {
+ @Override
+ public ClientCall interceptCall(
+ MethodDescriptor method, CallOptions callOptions, Channel next) {
+ return new io.grpc.ForwardingClientCall.SimpleForwardingClientCall<
+ ReqT, RespT>(next.newCall(method, callOptions)) {
+ @Override
+ public void start(Listener responseListener, Metadata headers) {
+ sidecarListenerRef.set((Listener) responseListener);
+ super.start(responseListener, headers);
+ }
- ClientInterceptor interceptor = filter.buildClientInterceptor(filterConfig, null, scheduler);
+ @Override
+ public boolean isReady() {
+ return sidecarReady.get();
+ }
+ };
+ }
+ })
+ .build());
+ });
- MutableHandlerRegistry dataPlaneRegistry = new MutableHandlerRegistry();
- dataPlaneRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_BIDI_STREAMING, (call, headers) -> {
- call.sendHeaders(new Metadata());
- call.request(100);
- return new ServerCall.Listener() {
- @Override
- public void onMessage(String message) {
- call.sendMessage("server-response-" + message);
- }
+ ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
+ filterConfig, channelManager, scheduler, FAKE_CONTEXT);
- @Override
- public void onHalfClose() {
- call.close(Status.OK, new Metadata());
- }
- };
- })
+ dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
+ (request, responseObserver) -> {
+ responseObserver.onNext("Hello " + request);
+ responseObserver.onCompleted();
+ }))
.build());
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueDataPlaneServerName)
- .fallbackHandlerRegistry(dataPlaneRegistry)
- .executor(Executors.newSingleThreadExecutor())
- .build().start());
-
- ManagedChannel channel =
- grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueDataPlaneServerName).directExecutor().build());
- Channel interceptedChannel = io.grpc.ClientInterceptors.interceptForward(
- channel,
- Arrays.asList(new XdsNameResolver.RawMessageClientInterceptor(), interceptor));
-
- ClientCall call =
- interceptedChannel.newCall(
- METHOD_BIDI_STREAMING,
- DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()));
-
- call.start(new ClientCall.Listener() {
- @Override
- public void onHeaders(Metadata headers) {
- appEvents.add("HEADERS");
- }
+ 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());
- @Override
- public void onMessage(String message) {
- appEvents.add("MESSAGE:" + message);
- }
+ CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
+ ClientCall proxyCall =
+ interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
+ proxyCall.start(new ClientCall.Listener() {}, new Metadata());
- @Override
- public void onClose(Status status, Metadata trailers) {
- appEvents.add("CLOSE:" + status.getCode());
- appTrailers.set(trailers);
- finishLatch.countDown();
- }
- }, new Metadata());
+ // Wait for sidecar call to start
+ long startTime = System.currentTimeMillis();
+ while (sidecarListenerRef.get() == null && System.currentTimeMillis() - startTime < 5000) {
+ }
+ assertThat(sidecarListenerRef.get()).isNotNull();
- call.request(100);
-
- // 1. Send Message 1 (should succeed and be allowed)
- call.sendMessage("msg1");
-
- // 2. Send Message 2 (should trigger the delay and then ImmediateResponse on ext_proc)
- call.sendMessage("msg2");
+ // Sidecar is busy initially
+ sidecarReady.set(false);
- // 3. Concurrent write of Message 3 (while ext_proc is sleeping)
- try {
- call.sendMessage("msg3");
- } catch (IllegalStateException e) {
- appEvents.add("WRITE_FAILED");
- }
+ // Request from application
+ proxyCall.request(10);
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(0);
- call.halfClose();
+ // Sidecar becomes ready
+ sidecarReady.set(true);
+ sidecarListenerRef.get().onReady();
- assertThat(finishLatch.await(5, TimeUnit.SECONDS)).isTrue();
- assertThat(appEvents).contains("CLOSE:UNAUTHENTICATED");
- assertThat(appEvents).doesNotContain("WRITE_FAILED");
- assertThat(appTrailers.get().get(immediateKey)).isEqualTo("true");
- assertThat(extProcCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ // Verify buffered request drained
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(10);
+ assertThat(proxyCall.isReady()).isTrue();
- sidecarResponseExecutor.shutdown();
+ proxyCall.cancel("Cleanup", null);
channelManager.close();
}
@Test
@SuppressWarnings("unchecked")
- public void givenFailureModeAllowFalse_whenExtProcStreamFails_thenDataPlaneCallCancelled()
+ public void givenExtProcStreamCompleted_whenAppRequestsMessages_thenRequestsForwarded()
throws Exception {
ExternalProcessor proto = ExternalProcessor.newBuilder()
.setGrpcService(GrpcService.newBuilder()
@@ -9283,14 +9368,13 @@ public void givenFailureModeAllowFalse_whenExtProcStreamFails_thenDataPlaneCallC
.build())
.build())
.build())
- .setFailureModeAllow(false) // Fail Closed
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
assertThat(configOrError.errorDetail).isNull();
ExternalProcessorFilterConfig filterConfig = configOrError.config;
- // External Processor Server triggers error
+ // External Processor Server
ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl;
extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
@@ -9302,11 +9386,8 @@ public StreamObserver process(
@Override
public void onNext(ProcessingRequest request) {
if (request.hasRequestHeaders()) {
- // Fail the stream immediately on headers
- responseObserver.onError(
- Status.INTERNAL
- .withDescription("Simulated sidecar failure")
- .asRuntimeException());
+ // Immediately complete the stream from server side
+ responseObserver.onCompleted();
}
}
@@ -9333,72 +9414,125 @@ public void onCompleted() {
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());
+
+ final AtomicInteger dataPlaneRequestCount = new AtomicInteger(0);
ManagedChannel dataPlaneChannel = grpcCleanup.register(
- InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+ 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());
- final AtomicReference closedStatus = new AtomicReference<>();
- final CountDownLatch closedLatch = new CountDownLatch(1);
- ClientCall.Listener appListener = new ClientCall.Listener() {
- @Override
- public void onClose(Status status, Metadata trailers) {
- closedStatus.set(status);
- closedLatch.countDown();
- }
- };
-
+ final CountDownLatch readyLatch = new CountDownLatch(1);
CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
- proxyCall.start(appListener, new Metadata());
+ proxyCall.start(new ClientCall.Listener() {
+ @Override
+ public void onReady() {
+ readyLatch.countDown();
+ }
+ }, new Metadata());
- // Verify application receives INTERNAL due to sidecar failure
- assertThat(closedLatch.await(5, TimeUnit.SECONDS)).isTrue();
- assertThat(closedStatus.get().getCode()).isEqualTo(Status.Code.INTERNAL);
- assertThat(closedStatus.get().getDescription()).contains("External processor stream failed");
-
- proxyCall.cancel("Cleanup", null);
- channelManager.close();
+ // Wait for sidecar stream completion
+ assertThat(readyLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(proxyCall.isReady()).isTrue();
+
+ proxyCall.request(7);
+
+ // Verify request forwarded immediately
+ assertThat(dataPlaneRequestCount.get()).isEqualTo(7);
+ // proxyCall.isReady() should remain true as sidecar is gone
+ assertThat(proxyCall.isReady()).isTrue();
+
+ proxyCall.cancel("Cleanup", null);
+ channelManager.close();
}
+ // --- Category 17: Error Handling & Security ---
+
@Test
- @SuppressWarnings("unchecked")
- public void givenFailureModeAllowTrue_whenExtProcStreamFails_thenCallFailsOpen()
+ @SuppressWarnings("FutureReturnValueIgnored")
+ public void givenPendingData_whenImmediateResponseReceived_thenDeliversDataBeforeStatus()
throws Exception {
- 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())
- .setFailureModeAllow(true) // Fail Open
- .build();
- ConfigOrError configOrError =
- provider.parseFilterConfig(Any.pack(proto), filterContext);
- assertThat(configOrError.errorDetail).isNull();
- ExternalProcessorFilterConfig filterConfig = configOrError.config;
+ final String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ final String uniqueDataPlaneServerName = InProcessServerBuilder.generateName();
+ final List appEvents = Collections.synchronizedList(new ArrayList<>());
+ final CountDownLatch finishLatch = new CountDownLatch(1);
+ final CountDownLatch extProcCompletedLatch = new CountDownLatch(1);
+ final ExecutorService sidecarResponseExecutor = Executors.newSingleThreadExecutor();
+ final Metadata.Key immediateKey =
+ Metadata.Key.of("x-immediate-header", Metadata.ASCII_STRING_MARSHALLER);
+ final AtomicReference appTrailers = new AtomicReference<>();
- // External Processor Server
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) {
- if (request.hasRequestHeaders()) {
- new Thread(() -> {
- synchronized (responseObserver) {
- responseObserver.onError(Status.INTERNAL.asRuntimeException());
+ sidecarResponseExecutor.submit(() -> {
+ synchronized (responseObserver) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder().build())
+ .build())
+ .build());
+ } else if (request.hasResponseHeaders()) {
+ try {
+ Thread.sleep(500);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setImmediateResponse(ImmediateResponse.newBuilder()
+ .setGrpcStatus(
+ io.envoyproxy.envoy.service.ext_proc.v3.GrpcStatus.newBuilder()
+ .setStatus(Status.UNAUTHENTICATED.getCode().value())
+ .build())
+ .setDetails("Immediate Auth Failure")
+ .setHeaders(
+ io.envoyproxy.envoy.service.ext_proc.v3.HeaderMutation.newBuilder()
+ .addSetHeaders(
+ io.envoyproxy.envoy.config.core.v3.HeaderValueOption
+ .newBuilder()
+ .setHeader(
+ io.envoyproxy.envoy.config.core.v3.HeaderValue
+ .newBuilder()
+ .setKey("x-immediate-header")
+ .setValue("true")
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ responseObserver.onCompleted();
}
- }).start();
- }
+ }
+ });
}
@Override
@@ -9407,247 +9541,311 @@ public void onError(Throwable t) {
@Override
public void onCompleted() {
+ extProcCompletedLatch.countDown();
}
};
}
};
- grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
- .addService(extProcImpl)
- .directExecutor()
- .build().start());
+
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl).directExecutor().build().start());
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(extProcServerName).directExecutor().build());
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
});
- ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
- filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+ ExternalProcessorFilter filter = new ExternalProcessorFilter(FAKE_CONTEXT, channelManager);
+ ExternalProcessor proto = createBaseProto(extProcServerName)
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.NONE)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
- final CountDownLatch dataPlaneLatch = new CountDownLatch(1);
- final CountDownLatch headersReceivedLatch = new CountDownLatch(1);
- final CountDownLatch resumeAsyncThreadLatch = new CountDownLatch(1);
+ ClientInterceptor interceptor = filter.buildClientInterceptor(filterConfig, null, scheduler);
- ServerInterceptor dataPlaneInterceptor = new ServerInterceptor() {
- @Override
- public ServerCall.Listener interceptCall(
- ServerCall call, Metadata headers, ServerCallHandler next) {
- headersReceivedLatch.countDown();
- try {
- resumeAsyncThreadLatch.await();
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
- return next.startCall(call, headers);
- }
- };
+ MutableHandlerRegistry dataPlaneRegistry = new MutableHandlerRegistry();
+ dataPlaneRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_SAY_HELLO, (call, headers) -> {
+ call.sendHeaders(new Metadata());
+ call.request(1);
+ return new ServerCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ call.sendMessage("server-response");
+ call.close(Status.OK, new Metadata());
+ }
+ };
+ })
+ .build());
- dataPlaneServiceRegistry.addService(ServerInterceptors.intercept(
- ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello " + request);
- responseObserver.onCompleted();
- dataPlaneLatch.countDown();
- }))
- .build(),
- dataPlaneInterceptor));
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueDataPlaneServerName)
+ .fallbackHandlerRegistry(dataPlaneRegistry)
+ .executor(Executors.newSingleThreadExecutor())
+ .build().start());
- ManagedChannel dataPlaneChannel = grpcCleanup.register(
- InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+ ManagedChannel channel =
+ grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueDataPlaneServerName).directExecutor().build());
+ Channel interceptedChannel = io.grpc.ClientInterceptors.interceptForward(
+ channel,
+ Arrays.asList(new XdsNameResolver.RawMessageClientInterceptor(), interceptor));
- final AtomicReference statusRef = new AtomicReference<>();
- final CountDownLatch closedLatch = new CountDownLatch(1);
- ClientCall.Listener appListener = new ClientCall.Listener() {
+ ClientCall call =
+ interceptedChannel.newCall(
+ METHOD_SAY_HELLO,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()));
+ call.start(new ClientCall.Listener() {
@Override
- public void onClose(Status status, Metadata trailers) {
- statusRef.set(status);
- closedLatch.countDown();
+ public void onHeaders(Metadata headers) {
+ appEvents.add("HEADERS");
}
- };
-
- CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
- ClientCall proxyCall =
- interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
- proxyCall.start(appListener, new Metadata());
- // Trigger unary call. request(1) starts it.
- proxyCall.request(1);
-
- // Wait for the async sidecar thread to enter activateCall() and block inside interceptCall
- assertThat(headersReceivedLatch.await(5, TimeUnit.SECONDS)).isTrue();
-
- // Now, while the async thread is blocked (and passThroughMode is still false),
- // send a message and half-close.
- proxyCall.sendMessage("test");
- proxyCall.halfClose();
+ @Override
+ public void onMessage(String message) {
+ appEvents.add("MESSAGE");
+ }
- // Unblock the async thread
- resumeAsyncThreadLatch.countDown();
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ appEvents.add("CLOSE:" + status.getCode());
+ appTrailers.set(trailers);
+ finishLatch.countDown();
+ }
+ }, new Metadata());
- // Verify data plane call reached (failed open)
- assertThat(dataPlaneLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ call.request(1);
+ call.sendMessage("request-body");
+ call.halfClose();
- // Verify client call completes successfully
- assertThat(closedLatch.await(5, TimeUnit.SECONDS)).isTrue();
- assertThat(statusRef.get().isOk()).isTrue();
+ assertThat(finishLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(appEvents).containsExactly("HEADERS", "MESSAGE", "CLOSE:UNAUTHENTICATED");
+ assertThat(appTrailers.get().get(immediateKey)).isEqualTo("true");
+ assertThat(extProcCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ sidecarResponseExecutor.shutdown();
channelManager.close();
}
+
@Test
- @SuppressWarnings("unchecked")
- public void givenObservabilityMode_whenDataPlaneClosed_thenSidecarCloseIsDeferred()
+ @SuppressWarnings("FutureReturnValueIgnored")
+ public void
+ givenStreamingCall_whenImmediateResponseReceivedDuringRequestStreaming_thenTerminatesCleanly()
throws Exception {
- 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())
- .setObservabilityMode(true)
- .setDeferredCloseTimeout(
- com.google.protobuf.Duration.newBuilder().setSeconds(10).build())
- .build();
- ConfigOrError configOrError =
- provider.parseFilterConfig(Any.pack(proto), filterContext);
- assertThat(configOrError.errorDetail).isNull();
- ExternalProcessorFilterConfig filterConfig = configOrError.config;
+ final String uniqueExtProcServerName = InProcessServerBuilder.generateName();
+ final String uniqueDataPlaneServerName = InProcessServerBuilder.generateName();
+ final List appEvents = Collections.synchronizedList(new ArrayList<>());
+ final CountDownLatch finishLatch = new CountDownLatch(1);
+ final CountDownLatch extProcCompletedLatch = new CountDownLatch(1);
+ final ExecutorService sidecarResponseExecutor = Executors.newSingleThreadExecutor();
+ final Metadata.Key immediateKey =
+ Metadata.Key.of("x-immediate-header", Metadata.ASCII_STRING_MARSHALLER);
+ final AtomicReference appTrailers = new AtomicReference<>();
+ final AtomicInteger extProcRequestCount = new AtomicInteger(0);
- // External Processor Server
- final CountDownLatch sidecarCompletedLatch = 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() {
+ ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
+ new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
- public void onNext(ProcessingRequest request) {
- }
+ public StreamObserver process(
+ final StreamObserver responseObserver) {
+ ((ServerCallStreamObserver) responseObserver).request(100);
+ return new StreamObserver() {
+ @Override
+ public void onNext(ProcessingRequest request) {
+ sidecarResponseExecutor.submit(() -> {
+ synchronized (responseObserver) {
+ if (request.hasRequestHeaders()) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestHeaders(HeadersResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder().build())
+ .build())
+ .build());
+ } else if (request.hasRequestBody()) {
+ int count = extProcRequestCount.incrementAndGet();
+ if (count == 1) {
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setRequestBody(BodyResponse.newBuilder()
+ .setResponse(CommonResponse.newBuilder().build())
+ .build())
+ .build());
+ } else if (count == 2) {
+ try {
+ Thread.sleep(500);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ responseObserver.onNext(ProcessingResponse.newBuilder()
+ .setImmediateResponse(ImmediateResponse.newBuilder()
+ .setGrpcStatus(
+ io.envoyproxy.envoy.service.ext_proc.v3.GrpcStatus.newBuilder()
+ .setStatus(Status.UNAUTHENTICATED.getCode().value())
+ .build())
+ .setDetails("Immediate Auth Failure")
+ .setHeaders(
+ io.envoyproxy.envoy.service.ext_proc.v3.HeaderMutation
+ .newBuilder()
+ .addSetHeaders(
+ io.envoyproxy.envoy.config.core.v3.HeaderValueOption
+ .newBuilder()
+ .setHeader(
+ io.envoyproxy.envoy.config.core.v3.HeaderValue
+ .newBuilder()
+ .setKey("x-immediate-header")
+ .setValue("true")
+ .build())
+ .build())
+ .build())
+ .build())
+ .build());
+ responseObserver.onCompleted();
+ }
+ }
+ }
+ });
+ }
- @Override
- public void onError(Throwable t) {
- }
+ @Override
+ public void onError(Throwable t) {}
- @Override
- public void onCompleted() {
- sidecarCompletedLatch.countDown();
+ @Override
+ public void onCompleted() {
+ extProcCompletedLatch.countDown();
+ }
+ };
}
};
- }
- };
- final io.grpc.Server extProcServer =
- grpcCleanup.register(
- InProcessServerBuilder.forName(extProcServerName)
- .addService(extProcImpl)
- .executor(fakeClock.getScheduledExecutorService())
- .build()
- .start());
+
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ .addService(extProcImpl).directExecutor().build().start());
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(extProcServerName)
- .executor(fakeClock.getScheduledExecutorService())
- .build());
+ InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
});
- ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
- filterConfig, channelManager, scheduler, FAKE_CONTEXT);
+ ExternalProcessorFilter filter = new ExternalProcessorFilter(FAKE_CONTEXT, channelManager);
+ ExternalProcessor proto = createBaseProto(uniqueExtProcServerName)
+ .setProcessingMode(ProcessingMode.newBuilder()
+ .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
+ .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SKIP)
+ .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SKIP)
+ .build())
+ .build();
+ ConfigOrError configOrError =
+ provider.parseFilterConfig(Any.pack(proto), filterContext);
+ ExternalProcessorFilterConfig filterConfig = configOrError.config;
- ManagedChannel dataPlaneChannel = grpcCleanup.register(
- InProcessChannelBuilder.forName(dataPlaneServerName)
- .executor(fakeClock.getScheduledExecutorService())
- .build());
+ ClientInterceptor interceptor = filter.buildClientInterceptor(filterConfig, null, scheduler);
- try {
- final CountDownLatch appCloseLatch = new CountDownLatch(1);
- ClientCall.Listener appListener = new ClientCall.Listener() {
- @Override public void onClose(Status status, Metadata trailers) {
- appCloseLatch.countDown();
- }
- };
-
- CallOptions callOptions =
- DEFAULT_CALL_OPTIONS.withExecutor(fakeClock.getScheduledExecutorService());
- ClientCall proxyCall =
- interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
- proxyCall.start(appListener, new Metadata());
+ MutableHandlerRegistry dataPlaneRegistry = new MutableHandlerRegistry();
+ dataPlaneRegistry.addService(ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_BIDI_STREAMING, (call, headers) -> {
+ call.sendHeaders(new Metadata());
+ call.request(100);
+ return new ServerCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ call.sendMessage("server-response-" + message);
+ }
- // Data plane closes immediately
- proxyCall.halfClose();
- dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("test");
- responseObserver.onCompleted();
- }))
- .build());
- proxyCall.request(1);
+ @Override
+ public void onHalfClose() {
+ call.close(Status.OK, new Metadata());
+ }
+ };
+ })
+ .build());
- // Wait for app onClose
- for (int i = 0; i < 1000 && appCloseLatch.getCount() > 0; i++) {
- fakeClock.forwardTime(1, TimeUnit.SECONDS);
- }
- assertThat(appCloseLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ grpcCleanup.register(InProcessServerBuilder.forName(uniqueDataPlaneServerName)
+ .fallbackHandlerRegistry(dataPlaneRegistry)
+ .executor(Executors.newSingleThreadExecutor())
+ .build().start());
- // At this point, app received onClose, but sidecar should NOT be completed yet
- assertThat(sidecarCompletedLatch.getCount()).isEqualTo(1);
+ ManagedChannel channel =
+ grpcCleanup.register(
+ InProcessChannelBuilder.forName(uniqueDataPlaneServerName).directExecutor().build());
+ Channel interceptedChannel = io.grpc.ClientInterceptors.interceptForward(
+ channel,
+ Arrays.asList(new XdsNameResolver.RawMessageClientInterceptor(), interceptor));
- // Fast forward time to trigger deferred close
- fakeClock.forwardTime(10, TimeUnit.SECONDS);
-
- for (int i = 0; i < 100 && sidecarCompletedLatch.getCount() > 0; i++) {
- fakeClock.forwardTime(1, TimeUnit.SECONDS);
+ ClientCall call =
+ interceptedChannel.newCall(
+ METHOD_BIDI_STREAMING,
+ DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()));
+
+ call.start(new ClientCall.Listener() {
+ @Override
+ public void onHeaders(Metadata headers) {
+ appEvents.add("HEADERS");
}
- assertThat(sidecarCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue();
-
- proxyCall.cancel("Cleanup", null);
- } finally {
- dataPlaneChannel.shutdownNow();
- extProcServer.shutdownNow();
- for (int i = 0;
- i < 100 && (!dataPlaneChannel.isTerminated() || !extProcServer.isTerminated());
- i++) {
- fakeClock.forwardTime(1, TimeUnit.SECONDS);
+
+ @Override
+ public void onMessage(String message) {
+ appEvents.add("MESSAGE:" + message);
}
- channelManager.close();
+
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ appEvents.add("CLOSE:" + status.getCode());
+ appTrailers.set(trailers);
+ finishLatch.countDown();
+ }
+ }, new Metadata());
+
+ call.request(100);
+
+ // 1. Send Message 1 (should succeed and be allowed)
+ call.sendMessage("msg1");
+
+ // 2. Send Message 2 (should trigger the delay and then ImmediateResponse on ext_proc)
+ call.sendMessage("msg2");
+
+ // 3. Concurrent write of Message 3 (while ext_proc is sleeping)
+ try {
+ call.sendMessage("msg3");
+ } catch (IllegalStateException e) {
+ appEvents.add("WRITE_FAILED");
}
+
+ call.halfClose();
+
+ assertThat(finishLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(appEvents).contains("CLOSE:UNAUTHENTICATED");
+ assertThat(appEvents).doesNotContain("WRITE_FAILED");
+ assertThat(appTrailers.get().get(immediateKey)).isEqualTo("true");
+ assertThat(extProcCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue();
+
+ sidecarResponseExecutor.shutdown();
+ channelManager.close();
}
@Test
@SuppressWarnings("unchecked")
- public void givenUnsupportedCompressionInResponse_whenReceived_thenStreamErrored()
+ public void givenFailureModeAllowFalse_whenExtProcStreamFails_thenDataPlaneCallCancelled()
throws Exception {
- String uniqueExtProcServerName =
- "extProc-compression-" + InProcessServerBuilder.generateName();
- String uniqueDataPlaneServerName =
- "dataPlane-compression-" + InProcessServerBuilder.generateName();
ExternalProcessor proto = ExternalProcessor.newBuilder()
.setGrpcService(GrpcService.newBuilder()
.setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
- .setTargetUri("in-process:///" + uniqueExtProcServerName)
+ .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()
- .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC).build())
+ .setFailureModeAllow(false) // Fail Closed
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
assertThat(configOrError.errorDetail).isNull();
ExternalProcessorFilterConfig filterConfig = configOrError.config;
- // External Processor Server
+ // External Processor Server triggers error
ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl;
extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() {
@Override
@@ -9659,28 +9857,11 @@ public StreamObserver process(
@Override
public void onNext(ProcessingRequest request) {
if (request.hasRequestHeaders()) {
- synchronized (responseObserver) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestHeaders(HeadersResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- }
- } else if (request.hasRequestBody()) {
- // Simulate sidecar sending compressed body mutation (unsupported)
- synchronized (responseObserver) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestBody(BodyResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder()
- .setBodyMutation(BodyMutation.newBuilder()
- .setStreamedResponse(StreamedBodyResponse.newBuilder()
- .setGrpcMessageCompressed(true)
- .build())
- .build())
- .build())
- .build())
- .build());
- }
+ // Fail the stream immediately on headers
+ responseObserver.onError(
+ Status.INTERNAL
+ .withDescription("Simulated sidecar failure")
+ .asRuntimeException());
}
}
@@ -9690,57 +9871,25 @@ public void onError(Throwable t) {
@Override
public void onCompleted() {
- new Thread(() -> {
- synchronized (responseObserver) {
- responseObserver.onCompleted();
- }
- }).start();
}
};
}
};
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
.addService(extProcImpl)
- .executor(fakeClock.getScheduledExecutorService())
+ .directExecutor()
.build().start());
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueExtProcServerName)
- .executor(fakeClock.getScheduledExecutorService())
- .build());
+ InProcessChannelBuilder.forName(extProcServerName).directExecutor().build());
});
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
- final CountDownLatch dataPlaneLatch = new CountDownLatch(1);
- MutableHandlerRegistry uniqueRegistry = new MutableHandlerRegistry();
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueDataPlaneServerName)
- .fallbackHandlerRegistry(uniqueRegistry)
- .directExecutor()
- .build().start());
- uniqueRegistry.addService(ServerInterceptors.intercept(
- ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello " + request);
- responseObserver.onCompleted();
- dataPlaneLatch.countDown();
- }))
- .build(),
- new ServerInterceptor() {
- @Override
- public ServerCall.Listener interceptCall(
- ServerCall call, Metadata headers, ServerCallHandler next) {
- return next.startCall(call, headers);
- }
- }));
-
ManagedChannel dataPlaneChannel = grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueDataPlaneServerName)
- .executor(fakeClock.getScheduledExecutorService())
- .build());
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
final AtomicReference closedStatus = new AtomicReference<>();
final CountDownLatch closedLatch = new CountDownLatch(1);
@@ -9752,27 +9901,12 @@ public void onClose(Status status, Metadata trailers) {
}
};
- CallOptions callOptions =
- DEFAULT_CALL_OPTIONS.withExecutor(fakeClock.getScheduledExecutorService());
+ CallOptions callOptions = DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor());
ClientCall proxyCall =
interceptCall(interceptor, METHOD_SAY_HELLO, callOptions, dataPlaneChannel);
proxyCall.start(appListener, new Metadata());
- // 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
- proxyCall.request(1);
- proxyCall.sendMessage("test");
- proxyCall.halfClose();
-
- // Verify application receives INTERNAL with correct description
- for (int i = 0; i < 10000 && closedLatch.getCount() > 0; i++) {
- fakeClock.forwardTime(1, TimeUnit.MILLISECONDS);
- }
+ // Verify application receives INTERNAL due to sidecar failure
assertThat(closedLatch.await(5, TimeUnit.SECONDS)).isTrue();
assertThat(closedStatus.get().getCode()).isEqualTo(Status.Code.INTERNAL);
assertThat(closedStatus.get().getDescription()).contains("External processor stream failed");
@@ -9783,27 +9917,19 @@ public void onClose(Status status, Metadata trailers) {
@Test
@SuppressWarnings("unchecked")
- public void givenUnsupportedCompressionInResponseBody_whenReceived_thenStreamErrored()
+ public void givenFailureModeAllowTrue_whenExtProcStreamFails_thenCallFailsOpen()
throws Exception {
- String uniqueExtProcServerName =
- "extProc-resp-compression-" + InProcessServerBuilder.generateName();
- String uniqueDataPlaneServerName =
- "dataPlane-resp-compression-" + InProcessServerBuilder.generateName();
ExternalProcessor proto = ExternalProcessor.newBuilder()
.setGrpcService(GrpcService.newBuilder()
.setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
- .setTargetUri("in-process:///" + uniqueExtProcServerName)
+ .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()
- .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
- .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
- .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
- .build())
+ .setFailureModeAllow(true) // Fail Open
.build();
ConfigOrError configOrError =
provider.parseFilterConfig(Any.pack(proto), filterContext);
@@ -9822,36 +9948,11 @@ public StreamObserver process(
@Override
public void onNext(ProcessingRequest request) {
if (request.hasRequestHeaders()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestHeaders(HeadersResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- } else if (request.hasRequestBody()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setRequestBody(BodyResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- } else if (request.hasResponseHeaders()) {
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setResponseHeaders(HeadersResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder().build())
- .build())
- .build());
- } else if (request.hasResponseBody()) {
- // Simulate sidecar sending compressed body mutation (unsupported) for response body
- responseObserver.onNext(ProcessingResponse.newBuilder()
- .setResponseBody(BodyResponse.newBuilder()
- .setResponse(CommonResponse.newBuilder()
- .setBodyMutation(BodyMutation.newBuilder()
- .setStreamedResponse(StreamedBodyResponse.newBuilder()
- .setGrpcMessageCompressed(true)
- .build())
- .build())
- .build())
- .build())
- .build());
+ new Thread(() -> {
+ synchronized (responseObserver) {
+ responseObserver.onError(Status.INTERNAL.asRuntimeException());
+ }
+ }).start();
}
}
@@ -9861,46 +9962,61 @@ public void onError(Throwable t) {
@Override
public void onCompleted() {
- responseObserver.onCompleted();
}
};
}
};
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
+ grpcCleanup.register(InProcessServerBuilder.forName(extProcServerName)
.addService(extProcImpl)
.directExecutor()
.build().start());
CachedChannelManager channelManager = new CachedChannelManager(config -> {
return grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
+ InProcessChannelBuilder.forName(extProcServerName).directExecutor().build());
});
ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
filterConfig, channelManager, scheduler, FAKE_CONTEXT);
- MutableHandlerRegistry uniqueRegistry = new MutableHandlerRegistry();
- grpcCleanup.register(InProcessServerBuilder.forName(uniqueDataPlaneServerName)
- .fallbackHandlerRegistry(uniqueRegistry)
- .directExecutor()
- .build().start());
- uniqueRegistry.addService(ServerServiceDefinition.builder("test.TestService")
- .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
- (request, responseObserver) -> {
- responseObserver.onNext("Hello");
- responseObserver.onCompleted();
- }))
- .build());
+ final CountDownLatch dataPlaneLatch = new CountDownLatch(1);
+ final CountDownLatch headersReceivedLatch = new CountDownLatch(1);
+ final CountDownLatch resumeAsyncThreadLatch = new CountDownLatch(1);
- ManagedChannel dataPlaneChannel =
- grpcCleanup.register(
- InProcessChannelBuilder.forName(uniqueDataPlaneServerName).directExecutor().build());
+ ServerInterceptor dataPlaneInterceptor = new ServerInterceptor() {
+ @Override
+ public ServerCall.Listener interceptCall(
+ ServerCall call, Metadata headers, ServerCallHandler next) {
+ headersReceivedLatch.countDown();
+ try {
+ resumeAsyncThreadLatch.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ return next.startCall(call, headers);
+ }
+ };
- final AtomicReference closedStatus = new AtomicReference<>();
+ dataPlaneServiceRegistry.addService(ServerInterceptors.intercept(
+ ServerServiceDefinition.builder("test.TestService")
+ .addMethod(METHOD_SAY_HELLO, ServerCalls.asyncUnaryCall(
+ (request, responseObserver) -> {
+ responseObserver.onNext("Hello " + request);
+ responseObserver.onCompleted();
+ dataPlaneLatch.countDown();
+ }))
+ .build(),
+ dataPlaneInterceptor));
+
+ ManagedChannel dataPlaneChannel = grpcCleanup.register(
+ InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
+
+ final AtomicReference statusRef = new AtomicReference<>();
final CountDownLatch closedLatch = new CountDownLatch(1);
ClientCall.Listener appListener = new ClientCall.Listener