@@ -768,6 +768,11 @@ public void request(int numMessages) {
768768 super .request (numMessages );
769769 return ;
770770 }
771+ if (!config .getObservabilityMode ()
772+ && currentProcessingMode .getResponseBodyMode () != ProcessingMode .BodySendMode .GRPC ) {
773+ super .request (numMessages );
774+ return ;
775+ }
771776 if (!isSidecarReady ()) {
772777 pendingRequests .addAndGet (numMessages );
773778 return ;
@@ -795,6 +800,10 @@ public void sendMessage(InputStream message) {
795800
796801 ExtProcStreamState state = extProcStreamState .get ();
797802 if (state .isDraining () || state .isCompleted ()) {
803+ if (currentProcessingMode .getRequestBodyMode () == ProcessingMode .BodySendMode .NONE ) {
804+ super .sendMessage (message );
805+ return ;
806+ }
798807 try {
799808 ByteString copiedBody = ByteString .readFrom (message );
800809 pendingDrainingMessages .add (new KnownLengthInputStream (copiedBody ));
@@ -855,6 +864,18 @@ public void halfClose() {
855864 }
856865
857866 if (extProcStreamState .get ().isDraining ()) {
867+ boolean canProceed = false ;
868+ synchronized (streamLock ) {
869+ if (currentProcessingMode .getRequestBodyMode () == ProcessingMode .BodySendMode .NONE
870+ || (!bodyMessageSentToExtProc .get () && pendingDrainingMessages .isEmpty ())) {
871+ canProceed = true ;
872+ }
873+ }
874+ if (canProceed ) {
875+ if (requestSideClosed .compareAndSet (false , true )) {
876+ proceedWithHalfClose ();
877+ }
878+ }
858879 return ;
859880 }
860881
@@ -1072,16 +1093,17 @@ public void onReady() {
10721093 public void onHeaders (Metadata headers ) {
10731094 dataPlaneClientCall .setServerHeadersStartNanos (System .nanoTime ());
10741095 responseHeadersSent .set (true );
1075- if (dataPlaneClientCall .getExtProcStreamState ().get ().isDraining ()) {
1076- this .savedHeaders = headers ;
1077- return ;
1078- }
10791096 boolean sendResponseHeaders =
10801097 dataPlaneClientCall .getCurrentProcessingMode ().getResponseHeaderMode ()
10811098 == ProcessingMode .HeaderSendMode .SEND
10821099 || dataPlaneClientCall .getCurrentProcessingMode ().getResponseHeaderMode ()
10831100 == ProcessingMode .HeaderSendMode .DEFAULT ;
10841101
1102+ if (dataPlaneClientCall .getExtProcStreamState ().get ().isDraining () && sendResponseHeaders ) {
1103+ this .savedHeaders = headers ;
1104+ return ;
1105+ }
1106+
10851107 if (dataPlaneClientCall .getPassThroughMode ().get ()
10861108 || dataPlaneClientCall .getExtProcStreamState ().get ().isCompleted ()
10871109 || !sendResponseHeaders ) {
@@ -1110,8 +1132,11 @@ public void onMessage(InputStream message) {
11101132 return ;
11111133 }
11121134
1113- if (savedHeaders != null
1114- || dataPlaneClientCall .getExtProcStreamState ().get ().isDraining ()) {
1135+ boolean checkDrain = dataPlaneClientCall .getExtProcStreamState ().get ().isDraining ()
1136+ && dataPlaneClientCall .getCurrentProcessingMode ().getResponseBodyMode ()
1137+ == ProcessingMode .BodySendMode .GRPC ;
1138+
1139+ if (savedHeaders != null || checkDrain ) {
11151140 try {
11161141 ByteString copiedBody = ByteString .readFrom (message );
11171142 savedMessages .add (new KnownLengthInputStream (copiedBody ));
@@ -1184,7 +1209,11 @@ public void onClose(Status status, Metadata trailers) {
11841209 return ;
11851210 }
11861211
1187- if (dataPlaneClientCall .getExtProcStreamState ().get ().isDraining ()) {
1212+ boolean sendResponseTrailers =
1213+ dataPlaneClientCall .getCurrentProcessingMode ().getResponseTrailerMode ()
1214+ == ProcessingMode .HeaderSendMode .SEND ;
1215+
1216+ if (dataPlaneClientCall .getExtProcStreamState ().get ().isDraining () && sendResponseTrailers ) {
11881217 return ;
11891218 }
11901219
@@ -1193,15 +1222,6 @@ public void onClose(Status status, Metadata trailers) {
11931222 }
11941223
11951224 triggerCloseHandshake ();
1196-
1197- if (dataPlaneClientCall .getConfig ().getObservabilityMode ()) {
1198- proceedWithClose ();
1199- @ SuppressWarnings ("unused" )
1200- ScheduledFuture <?> unused = dataPlaneClientCall .getScheduler ().schedule (
1201- dataPlaneClientCall ::closeExtProcStream ,
1202- dataPlaneClientCall .getConfig ().getDeferredCloseTimeoutNanos (),
1203- TimeUnit .NANOSECONDS );
1204- }
12051225 }
12061226
12071227 void onReadyNotify () {
@@ -1330,6 +1350,15 @@ private void triggerCloseHandshake() {
13301350 dataPlaneClientCall .closeExtProcStream ();
13311351 }
13321352 }
1353+
1354+ if (dataPlaneClientCall .getConfig ().getObservabilityMode ()) {
1355+ proceedWithClose ();
1356+ @ SuppressWarnings ("unused" )
1357+ ScheduledFuture <?> unused = dataPlaneClientCall .getScheduler ().schedule (
1358+ dataPlaneClientCall ::closeExtProcStream ,
1359+ dataPlaneClientCall .getConfig ().getDeferredCloseTimeoutNanos (),
1360+ TimeUnit .NANOSECONDS );
1361+ }
13331362 }
13341363
13351364 private void sendResponseBodyToExtProc (
0 commit comments