@@ -12014,364 +12014,7 @@ public void onMessage(String message) {
1201412014 channelManager.close();
1201512015 }
1201612016
12017- @Test
12018- @SuppressWarnings("unchecked")
12019- public void testSidestreamToDownstreamFlowControl_Violations() throws Exception {
12020- ExternalProcessor proto = ExternalProcessor.newBuilder()
12021- .setGrpcService(GrpcService.newBuilder()
12022- .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
12023- .setTargetUri("in-process:///" + extProcServerName)
12024- .addChannelCredentialsPlugin(Any.newBuilder()
12025- .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
12026- .build())
12027- .build())
12028- .build())
12029- .setProcessingMode(ProcessingMode.newBuilder()
12030- .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
12031- .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
12032- .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
12033- .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
12034- .build())
12035- .build();
12036- ConfigOrError<ExternalProcessorFilterConfig> configOrError =
12037- provider.parseFilterConfig(Any.pack(proto), filterContext);
12038- assertThat(configOrError.errorDetail).isNull();
12039- ExternalProcessorFilterConfig filterConfig = configOrError.config;
12040-
12041- final String mutatedMessageTooLarge = new String(new char[70000]).replace('\0', 'd');
12042- final CountDownLatch callClosedLatch = new CountDownLatch(1);
12043- final AtomicReference<Status> capturedStatus = new AtomicReference<>();
12044-
12045- ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
12046- new ExternalProcessorGrpc.ExternalProcessorImplBase() {
12047- @Override
12048- public StreamObserver<ProcessingRequest> process(
12049- final StreamObserver<ProcessingResponse> responseObserver) {
12050- ((ServerCallStreamObserver<ProcessingResponse>) responseObserver).request(100);
12051- return new StreamObserver<ProcessingRequest>() {
12052- @Override
12053- public void onNext(ProcessingRequest request) {
12054- if (request.hasRequestHeaders()) {
12055- responseObserver.onNext(ProcessingResponse.newBuilder()
12056- .setRequestHeaders(HeadersResponse.newBuilder().build())
12057- .build());
12058- } else if (request.hasRequestBody()) {
12059- responseObserver.onNext(ProcessingResponse.newBuilder()
12060- .setRequestBody(BodyResponse.newBuilder()
12061- .setResponse(CommonResponse.newBuilder()
12062- .setBodyMutation(BodyMutation.newBuilder()
12063- .setStreamedResponse(StreamedBodyResponse.newBuilder()
12064- .setBody(request.getRequestBody().getBody())
12065- .build())
12066- .build())
12067- .build())
12068- .build())
12069- .build());
12070- } else if (request.hasResponseHeaders()) {
12071- responseObserver.onNext(ProcessingResponse.newBuilder()
12072- .setResponseHeaders(HeadersResponse.newBuilder().build())
12073- .build());
12074- } else if (request.hasResponseBody()) {
12075- // Send Chunk 1 (70k) -> Consumed by client app's pending request count of 1.
12076- // Replenishes window.
12077- responseObserver.onNext(ProcessingResponse.newBuilder()
12078- .setResponseBody(BodyResponse.newBuilder()
12079- .setResponse(CommonResponse.newBuilder()
12080- .setBodyMutation(BodyMutation.newBuilder()
12081- .setStreamedResponse(StreamedBodyResponse.newBuilder()
12082- .setBody(ByteString.copyFromUtf8(mutatedMessageTooLarge))
12083- .build())
12084- .build())
12085- .build())
12086- .build())
12087- .build());
12088-
12089- // Send Chunk 2 (70k) -> Excess chunk. Stored in queue. Window stays negative.
12090- responseObserver.onNext(ProcessingResponse.newBuilder()
12091- .setResponseBody(BodyResponse.newBuilder()
12092- .setResponse(CommonResponse.newBuilder()
12093- .setBodyMutation(BodyMutation.newBuilder()
12094- .setStreamedResponse(StreamedBodyResponse.newBuilder()
12095- .setBody(ByteString.copyFromUtf8(mutatedMessageTooLarge))
12096- .build())
12097- .build())
12098- .build())
12099- .build())
12100- .build());
12101-
12102- // Send Chunk 3 (1 byte) -> Arrives when window is negative, triggering violation.
12103- responseObserver.onNext(ProcessingResponse.newBuilder()
12104- .setResponseBody(BodyResponse.newBuilder()
12105- .setResponse(CommonResponse.newBuilder()
12106- .setBodyMutation(BodyMutation.newBuilder()
12107- .setStreamedResponse(StreamedBodyResponse.newBuilder()
12108- .setBody(ByteString.copyFromUtf8("a"))
12109- .build())
12110- .build())
12111- .build())
12112- .build())
12113- .build());
12114- }
12115- }
12116-
12117- @Override
12118- public void onError(Throwable t) {}
12119-
12120- @Override
12121- public void onCompleted() {
12122- responseObserver.onCompleted();
12123- }
12124- };
12125- }
12126- };
12127-
12128- String uniqueExtProcServerName = InProcessServerBuilder.generateName();
12129- grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
12130- .addService(extProcImpl)
12131- .directExecutor()
12132- .build().start());
12133-
12134- CachedChannelManager channelManager = new CachedChannelManager(config -> {
12135- return grpcCleanup.register(
12136- InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
12137- });
12138-
12139- ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
12140- filterConfig, channelManager, scheduler, FAKE_CONTEXT);
12141-
12142- final AtomicReference<StreamObserver<String>> dataPlaneResponseObserverRef =
12143- new AtomicReference<>();
12144- dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
12145- .addMethod(METHOD_BIDI_STREAMING, ServerCalls.asyncBidiStreamingCall(
12146- new ServerCalls.BidiStreamingMethod<String, String>() {
12147- @Override
12148- public StreamObserver<String> invoke(StreamObserver<String> responseObserver) {
12149- dataPlaneResponseObserverRef.set(responseObserver);
12150- return new StreamObserver<String>() {
12151- @Override
12152- public void onNext(String value) {}
12153-
12154- @Override
12155- public void onError(Throwable t) {}
12156-
12157- @Override
12158- public void onCompleted() {}
12159- };
12160- }
12161- }))
12162- .build());
12163-
12164- ManagedChannel dataPlaneChannel = grpcCleanup.register(
12165- InProcessChannelBuilder.forName(dataPlaneServerName).directExecutor().build());
12166-
12167- ClientCall<String, String> proxyCall =
12168- interceptCall(interceptor, METHOD_BIDI_STREAMING,
12169- DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()),
12170- dataPlaneChannel);
12171-
12172- proxyCall.start(new ClientCall.Listener<String>() {
12173- @Override
12174- public void onClose(Status status, Metadata trailers) {
12175- capturedStatus.set(status);
12176- callClosedLatch.countDown();
12177- }
12178- }, new Metadata());
12179- proxyCall.request(1);
12180-
12181- proxyCall.sendMessage("Client Msg");
12182-
12183- // Send a response from upstream to trigger headers and then the body response
12184- StreamObserver<String> upstreamResponseObserver = dataPlaneResponseObserverRef.get();
12185- upstreamResponseObserver.onNext("Response Msg");
12186-
12187- assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue();
12188-
12189- // The call should fail immediately with INTERNAL error code due to flow control violation
12190- assertThat(capturedStatus.get().getCode()).isEqualTo(Status.Code.INTERNAL);
12191- assertThat(capturedStatus.get().getDescription())
12192-
12193- .isEqualTo("External processor stream failed");
12194- assertThat(capturedStatus.get().getCause()).isInstanceOf(io.grpc.StatusRuntimeException.class);
12195- assertThat(capturedStatus.get().getCause().getMessage())
12196- .contains("Flow control violation: received server body from ext_proc "
12197- + "when window is closed");
12198-
12199- channelManager.close();
12200- }
12201-
12202- @Test
12203- @SuppressWarnings("unchecked")
12204- public void testSidestreamToUpstreamFlowControl_Violations() throws Exception {
12205- ExternalProcessor proto = ExternalProcessor.newBuilder()
12206- .setGrpcService(GrpcService.newBuilder()
12207- .setGoogleGrpc(GrpcService.GoogleGrpc.newBuilder()
12208- .setTargetUri("in-process:///" + extProcServerName)
12209- .addChannelCredentialsPlugin(Any.newBuilder()
12210- .setTypeUrl(INSECURE_CREDENTIALS_TYPE_URL)
12211- .build())
12212- .build())
12213- .build())
12214- .setProcessingMode(ProcessingMode.newBuilder()
12215- .setRequestBodyMode(ProcessingMode.BodySendMode.GRPC)
12216- .setResponseBodyMode(ProcessingMode.BodySendMode.GRPC)
12217- .setResponseHeaderMode(ProcessingMode.HeaderSendMode.SEND)
12218- .setResponseTrailerMode(ProcessingMode.HeaderSendMode.SEND)
12219- .build())
12220- .build();
12221- ConfigOrError<ExternalProcessorFilterConfig> configOrError =
12222- provider.parseFilterConfig(Any.pack(proto), filterContext);
12223- assertThat(configOrError.errorDetail).isNull();
12224- ExternalProcessorFilterConfig filterConfig = configOrError.config;
12225-
12226- final String mutatedMessageTooLarge = new String(new char[70000]).replace('\0', 'c');
12227- final CountDownLatch callClosedLatch = new CountDownLatch(1);
12228- final AtomicReference<Status> capturedStatus = new AtomicReference<>();
12229-
12230- ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl =
12231- new ExternalProcessorGrpc.ExternalProcessorImplBase() {
12232- @Override
12233- public StreamObserver<ProcessingRequest> process(
12234- final StreamObserver<ProcessingResponse> responseObserver) {
12235- ((ServerCallStreamObserver<ProcessingResponse>) responseObserver).request(100);
12236- return new StreamObserver<ProcessingRequest>() {
12237- @Override
12238- public void onNext(ProcessingRequest request) {
12239- if (request.hasRequestHeaders()) {
12240- responseObserver.onNext(ProcessingResponse.newBuilder()
12241- .setRequestHeaders(HeadersResponse.newBuilder().build())
12242- .build());
12243- } else if (request.hasRequestBody()) {
12244- ByteString original = request.getRequestBody().getBody();
12245- if (original.toStringUtf8().equals("Message 1")) {
12246- responseObserver.onNext(ProcessingResponse.newBuilder()
12247- .setRequestBody(BodyResponse.newBuilder()
12248- .setResponse(CommonResponse.newBuilder()
12249- .setBodyMutation(BodyMutation.newBuilder()
12250- .setStreamedResponse(StreamedBodyResponse.newBuilder()
12251- .setBody(ByteString.copyFromUtf8(mutatedMessageTooLarge))
12252- .build())
12253- .build())
12254- .build())
12255- .build())
12256- .build());
12257- } else {
12258- responseObserver.onNext(ProcessingResponse.newBuilder()
12259- .setRequestBody(BodyResponse.newBuilder()
12260- .setResponse(CommonResponse.newBuilder()
12261- .setBodyMutation(BodyMutation.newBuilder()
12262- .setStreamedResponse(StreamedBodyResponse.newBuilder()
12263- .setBody(ByteString.copyFromUtf8("a"))
12264- .build())
12265- .build())
12266- .build())
12267- .build())
12268- .build());
12269- }
12270- }
12271- }
12272-
12273- @Override
12274- public void onError(Throwable t) {}
1227512017
12276- @Override
12277- public void onCompleted() {
12278- responseObserver.onCompleted();
12279- }
12280- };
12281- }
12282- };
12283-
12284- String uniqueExtProcServerName = InProcessServerBuilder.generateName();
12285- grpcCleanup.register(InProcessServerBuilder.forName(uniqueExtProcServerName)
12286- .addService(extProcImpl)
12287- .directExecutor()
12288- .build().start());
12289-
12290- CachedChannelManager channelManager = new CachedChannelManager(config -> {
12291- return grpcCleanup.register(
12292- InProcessChannelBuilder.forName(uniqueExtProcServerName).directExecutor().build());
12293- });
12294-
12295- ExternalProcessorClientInterceptor interceptor = new ExternalProcessorClientInterceptor(
12296- filterConfig, channelManager, scheduler, FAKE_CONTEXT);
12297-
12298- dataPlaneServiceRegistry.addService(ServerServiceDefinition.builder("test.TestService")
12299- .addMethod(METHOD_CLIENT_STREAMING, ServerCalls.asyncClientStreamingCall(
12300- new ServerCalls.ClientStreamingMethod<String, String>() {
12301- @Override
12302- public StreamObserver<String> invoke(StreamObserver<String> responseObserver) {
12303- return new StreamObserver<String>() {
12304- @Override
12305- public void onNext(String value) {}
12306-
12307- @Override
12308- public void onError(Throwable t) {}
12309-
12310- @Override
12311- public void onCompleted() {
12312- responseObserver.onNext("Response");
12313- responseObserver.onCompleted();
12314- }
12315- };
12316- }
12317- }))
12318- .build());
12319-
12320- // Build the data plane channel with a custom ClientInterceptor.
12321- // This interceptor overrides isReady() to always return false.
12322- // This simulates a blocked backend server
12323- // (transport flow control buffer is full) and blocks the filter's upstream
12324- // window replenishment.
12325- ManagedChannel dataPlaneChannel = grpcCleanup.register(
12326- InProcessChannelBuilder.forName(dataPlaneServerName)
12327- .intercept(new ClientInterceptor() {
12328- @Override
12329- public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
12330- MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
12331- return new io.grpc.ForwardingClientCall.SimpleForwardingClientCall<ReqT, RespT>(
12332- next.newCall(method, callOptions)) {
12333- @Override
12334- public boolean isReady() {
12335- return false;
12336- }
12337- };
12338- }
12339- })
12340- .directExecutor()
12341- .build());
12342-
12343- ClientCall<String, String> proxyCall =
12344- interceptCall(interceptor, METHOD_CLIENT_STREAMING,
12345- DEFAULT_CALL_OPTIONS.withExecutor(MoreExecutors.directExecutor()), dataPlaneChannel);
12346-
12347- proxyCall.start(new ClientCall.Listener<String>() {
12348- @Override
12349- public void onClose(Status status, Metadata trailers) {
12350- capturedStatus.set(status);
12351- callClosedLatch.countDown();
12352- }
12353- }, new Metadata());
12354-
12355- // Send first message. Window becomes negative.
12356- proxyCall.sendMessage("Message 1");
12357-
12358- // Send second message. Window is still negative, triggering flow control violation.
12359- proxyCall.sendMessage("Message 2");
12360-
12361- assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue();
12362-
12363- // The call should fail immediately with INTERNAL error code due to flow control violation
12364- assertThat(capturedStatus.get().getCode()).isEqualTo(Status.Code.INTERNAL);
12365- assertThat(capturedStatus.get().getDescription())
12366-
12367- .isEqualTo("External processor stream failed");
12368- assertThat(capturedStatus.get().getCause()).isInstanceOf(io.grpc.StatusRuntimeException.class);
12369- assertThat(capturedStatus.get().getCause().getMessage())
12370- .contains("Flow control violation: received client body from ext_proc "
12371- + "when window is closed");
12372-
12373- channelManager.close();
12374- }
1237512018
1237612019 @Test
1237712020 @SuppressWarnings("unchecked")
0 commit comments