From 44075af5008af5729b02423fbada330356562214 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Tue, 21 Jul 2026 09:02:40 +0000 Subject: [PATCH] xds: Fix concurrent StreamObserver usage in ExternalProcessorClientInterceptorTest Synchronize calls to the gRPC StreamObserver (responseObserver) inside mock ExternalProcessor implementations. The mock servers in several tests spawn background threads to handle requests, but they invoke responseObserver.onNext(), onCompleted(), and onError() concurrently without synchronization. Since standard StreamObservers are not thread-safe, this leads to ThreadSanitizer-reported data races in gRPC internals (such as ServerCallImpl). This change wraps all asynchronous/background invocations of responseObserver methods inside synchronized blocks on the observer instance to ensure serialized execution. --- ...xternalProcessorClientInterceptorTest.java | 439 +++++++++++------- 1 file changed, 259 insertions(+), 180 deletions(-) diff --git a/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java b/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java index 9b07cae3477..90c82fb26bf 100644 --- a/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java +++ b/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java @@ -1915,24 +1915,26 @@ public void onNext(ProcessingRequest request) { } catch (InterruptedException e) { Thread.currentThread().interrupt(); } - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder() - .setResponse(CommonResponse.newBuilder() - .setHeaderMutation(HeaderMutation.newBuilder() - .addSetHeaders( - io.envoyproxy.envoy.config.core.v3.HeaderValueOption - .newBuilder() - .setHeader( - io.envoyproxy.envoy.config.core.v3.HeaderValue - .newBuilder() - .setKey("x-mutated") - .setValue("true") - .build()) - .build()) - .build()) - .build()) - .build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder() + .setResponse(CommonResponse.newBuilder() + .setHeaderMutation(HeaderMutation.newBuilder() + .addSetHeaders( + io.envoyproxy.envoy.config.core.v3.HeaderValueOption + .newBuilder() + .setHeader( + io.envoyproxy.envoy.config.core.v3.HeaderValue + .newBuilder() + .setKey("x-mutated") + .setValue("true") + .build()) + .build()) + .build()) + .build()) + .build()) + .build()); + } } }).start(); } @@ -1943,7 +1945,11 @@ public void onError(Throwable t) { @Override public void onCompleted() { - new Thread(() -> responseObserver.onCompleted()).start(); + new Thread(() -> { + synchronized (responseObserver) { + responseObserver.onCompleted(); + } + }).start(); } }; } @@ -2058,10 +2064,12 @@ public void onNext(ProcessingRequest request) { capturedRequest.set(request); } new Thread(() -> { - if (request.hasRequestHeaders()) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .build()); + synchronized (responseObserver) { + if (request.hasRequestHeaders()) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .build()); + } } extProcLatch.countDown(); }).start(); @@ -2073,7 +2081,9 @@ public void onError(Throwable t) { @Override public void onCompleted() { - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } }; } @@ -2399,38 +2409,40 @@ public StreamObserver process( @Override public void onNext(ProcessingRequest request) { new Thread(() -> { - if (request.hasRequestHeaders()) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .build()); - } else if (request.hasRequestBody()) { - if (capturedRequest.get() == null - && !request.getRequestBody().getBody().isEmpty()) { - capturedRequest.set(request); - bodySentLatch.countDown(); - } - BodyResponse.Builder bodyResponse = BodyResponse.newBuilder(); - if (request.getRequestBody().getBody().isEmpty() - && request.getRequestBody().getEndOfStreamWithoutMessage()) { - bodyResponse.setResponse(CommonResponse.newBuilder() - .setBodyMutation(BodyMutation.newBuilder() - .setStreamedResponse(StreamedBodyResponse.newBuilder() - .setEndOfStream(true) - .build()) - .build()) + synchronized (responseObserver) { + if (request.hasRequestHeaders()) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) .build()); - } else { - bodyResponse.setResponse(CommonResponse.newBuilder() - .setBodyMutation(BodyMutation.newBuilder() - .setStreamedResponse(StreamedBodyResponse.newBuilder() - .setEndOfStream(request.getRequestBody().getEndOfStream()) - .build()) - .build()) + } else if (request.hasRequestBody()) { + if (capturedRequest.get() == null + && !request.getRequestBody().getBody().isEmpty()) { + capturedRequest.set(request); + bodySentLatch.countDown(); + } + BodyResponse.Builder bodyResponse = BodyResponse.newBuilder(); + if (request.getRequestBody().getBody().isEmpty() + && request.getRequestBody().getEndOfStreamWithoutMessage()) { + bodyResponse.setResponse(CommonResponse.newBuilder() + .setBodyMutation(BodyMutation.newBuilder() + .setStreamedResponse(StreamedBodyResponse.newBuilder() + .setEndOfStream(true) + .build()) + .build()) + .build()); + } else { + bodyResponse.setResponse(CommonResponse.newBuilder() + .setBodyMutation(BodyMutation.newBuilder() + .setStreamedResponse(StreamedBodyResponse.newBuilder() + .setEndOfStream(request.getRequestBody().getEndOfStream()) + .build()) + .build()) + .build()); + } + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestBody(bodyResponse.build()) .build()); } - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestBody(bodyResponse.build()) - .build()); } }).start(); } @@ -2441,7 +2453,11 @@ public void onError(Throwable t) { @Override public void onCompleted() { - new Thread(() -> responseObserver.onCompleted()).start(); + new Thread(() -> { + synchronized (responseObserver) { + responseObserver.onCompleted(); + } + }).start(); } }; } @@ -6242,14 +6258,18 @@ public StreamObserver process( public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { new Thread(() -> { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestDrain(true) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestDrain(true) + .build()); + } sidecarOnNextLatch.countDown(); try { if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) { sidecarOnCompletedLatch.countDown(); - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -6506,12 +6526,16 @@ public StreamObserver process( public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { new Thread(() -> { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestDrain(true) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestDrain(true) + .build()); + } try { if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -6664,13 +6688,17 @@ public StreamObserver process( public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { new Thread(() -> { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .setRequestDrain(true) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .setRequestDrain(true) + .build()); + } try { if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -6832,27 +6860,33 @@ public StreamObserver process( @Override public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .build()); + } } else if (request.hasRequestBody()) { extProcReceivedBodyLatch.countDown(); new Thread(() -> { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestBody(BodyResponse.newBuilder() - .setResponse(CommonResponse.newBuilder() - .setBodyMutation(BodyMutation.newBuilder() - .setStreamedResponse(StreamedBodyResponse.newBuilder() - .setBody(ByteString.copyFromUtf8("Mutated Message 1")) - .build()) - .build()) - .build()) - .build()) - .setRequestDrain(true) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestBody(BodyResponse.newBuilder() + .setResponse(CommonResponse.newBuilder() + .setBodyMutation(BodyMutation.newBuilder() + .setStreamedResponse(StreamedBodyResponse.newBuilder() + .setBody(ByteString.copyFromUtf8("Mutated Message 1")) + .build()) + .build()) + .build()) + .build()) + .setRequestDrain(true) + .build()); + } try { if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -7023,13 +7057,17 @@ public StreamObserver process( public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { new Thread(() -> { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .setRequestDrain(true) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .setRequestDrain(true) + .build()); + } try { if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -7204,14 +7242,18 @@ public StreamObserver process( @Override public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .build()); + } reqHeadersLatch.countDown(); } else if (request.hasResponseHeaders()) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setResponseHeaders(HeadersResponse.newBuilder().build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setResponseHeaders(HeadersResponse.newBuilder().build()) + .build()); + } respHeadersLatch.countDown(); } else if (request.hasResponseBody()) { String msgStr = request.getResponseBody().getBody().toStringUtf8(); @@ -7220,18 +7262,21 @@ public void onNext(ProcessingRequest request) { try { // Wait until M2 is received by sidecar so both M1 and M2 are in flight if (m2ReceivedLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setResponseBody(BodyResponse.newBuilder() - .setResponse(CommonResponse.newBuilder() - .setBodyMutation(BodyMutation.newBuilder() - .setStreamedResponse(StreamedBodyResponse.newBuilder() - .setBody(ByteString.copyFromUtf8("Mutated Message 1")) - .build()) - .build()) - .build()) - .build()) - .setRequestDrain(true) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setResponseBody(BodyResponse.newBuilder() + .setResponse(CommonResponse.newBuilder() + .setBodyMutation(BodyMutation.newBuilder() + .setStreamedResponse(StreamedBodyResponse.newBuilder() + .setBody(ByteString.copyFromUtf8( + "Mutated Message 1")) + .build()) + .build()) + .build()) + .build()) + .setRequestDrain(true) + .build()); + } respBody1Latch.countDown(); } } catch (InterruptedException e) { @@ -7244,21 +7289,26 @@ public void onNext(ProcessingRequest request) { try { // Wait until M3 is sent by upstream concurrently during drain if (m3SentLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setResponseBody(BodyResponse.newBuilder() - .setResponse(CommonResponse.newBuilder() - .setBodyMutation(BodyMutation.newBuilder() - .setStreamedResponse(StreamedBodyResponse.newBuilder() - .setBody(ByteString.copyFromUtf8("Mutated Message 2")) - .build()) - .build()) - .build()) - .build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setResponseBody(BodyResponse.newBuilder() + .setResponse(CommonResponse.newBuilder() + .setBodyMutation(BodyMutation.newBuilder() + .setStreamedResponse(StreamedBodyResponse.newBuilder() + .setBody(ByteString.copyFromUtf8( + "Mutated Message 2")) + .build()) + .build()) + .build()) + .build()) + .build()); + } respBody2Latch.countDown(); } if (sidecarFinishLatch.await(5, TimeUnit.SECONDS)) { - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onCompleted(); + } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -8431,7 +8481,9 @@ public StreamObserver process( public void onNext(ProcessingRequest request) { if (request.hasRequestHeaders()) { new Thread(() -> { - responseObserver.onError(Status.INTERNAL.asRuntimeException()); + synchronized (responseObserver) { + responseObserver.onError(Status.INTERNAL.asRuntimeException()); + } }).start(); } } @@ -8694,24 +8746,28 @@ 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()); + 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) - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestBody(BodyResponse.newBuilder() - .setResponse(CommonResponse.newBuilder() - .setBodyMutation(BodyMutation.newBuilder() - .setStreamedResponse(StreamedBodyResponse.newBuilder() - .setGrpcMessageCompressed(true) - .build()) - .build()) - .build()) - .build()) - .build()); + 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()); + } } } @@ -8721,7 +8777,11 @@ public void onError(Throwable t) { @Override public void onCompleted() { - new Thread(() -> responseObserver.onCompleted()).start(); + new Thread(() -> { + synchronized (responseObserver) { + responseObserver.onCompleted(); + } + }).start(); } }; } @@ -9596,45 +9656,52 @@ 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()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder() + .setResponse(CommonResponse.newBuilder().build()) + .build()) + .build()); + } } else if (request.hasResponseHeaders()) { - responseObserver.onNext(ProcessingResponse.newBuilder() - .setResponseHeaders(HeadersResponse.newBuilder() - .setResponse(CommonResponse.newBuilder().build()) - .build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setResponseHeaders(HeadersResponse.newBuilder() + .setResponse(CommonResponse.newBuilder().build()) + .build()) + .build()); + } } else if (request.hasResponseTrailers()) { new Thread(() -> { - responseObserver.onNext( - ProcessingResponse.newBuilder() - .setImmediateResponse( - ImmediateResponse.newBuilder() - .setGrpcStatus( - io.envoyproxy.envoy.service.ext_proc.v3.GrpcStatus.newBuilder() - .setStatus(Status.DATA_LOSS.getCode().value()) - .build()) - .setDetails("Sidecar detected data loss") - .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-sidecar-extra") - .setValue("true") - .build()) - .build()) - .build()) - .build()) - .build()); - responseObserver.onCompleted(); + synchronized (responseObserver) { + responseObserver.onNext( + ProcessingResponse.newBuilder() + .setImmediateResponse( + ImmediateResponse.newBuilder() + .setGrpcStatus( + io.envoyproxy.envoy.service.ext_proc.v3.GrpcStatus + .newBuilder() + .setStatus(Status.DATA_LOSS.getCode().value()) + .build()) + .setDetails("Sidecar detected data loss") + .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-sidecar-extra") + .setValue("true") + .build()) + .build()) + .build()) + .build()) + .build()); + responseObserver.onCompleted(); + } }).start(); } } @@ -9860,9 +9927,11 @@ public void onNext(ProcessingRequest request) { new Thread(() -> { if (request.hasRequestHeaders()) { sidecarActionLatch.countDown(); - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .build()); + } } }).start(); } @@ -9873,7 +9942,11 @@ public void onError(Throwable t) { @Override public void onCompleted() { - new Thread(() -> responseObserver.onCompleted()).start(); + new Thread(() -> { + synchronized (responseObserver) { + responseObserver.onCompleted(); + } + }).start(); } }; } @@ -10002,9 +10075,11 @@ public void onNext(ProcessingRequest request) { new Thread(() -> { if (request.hasRequestHeaders()) { sidecarActionLatch.countDown(); - responseObserver.onNext(ProcessingResponse.newBuilder() - .setRequestHeaders(HeadersResponse.newBuilder().build()) - .build()); + synchronized (responseObserver) { + responseObserver.onNext(ProcessingResponse.newBuilder() + .setRequestHeaders(HeadersResponse.newBuilder().build()) + .build()); + } } }).start(); } @@ -10015,7 +10090,11 @@ public void onError(Throwable t) { @Override public void onCompleted() { - new Thread(() -> responseObserver.onCompleted()).start(); + new Thread(() -> { + synchronized (responseObserver) { + responseObserver.onCompleted(); + } + }).start(); } }; }