@@ -197,15 +197,7 @@ public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
197197 MethodDescriptor <ReqT , RespT > method ,
198198 CallOptions callOptions ,
199199 Channel next ) {
200- java .util .concurrent .Executor callExecutor = callOptions .getExecutor ();
201- if (callExecutor == null ) {
202- callExecutor = new java .util .concurrent .Executor () {
203- @ Override
204- public void execute (Runnable command ) {
205- command .run ();
206- }
207- };
208- }
200+ Executor callExecutor = callOptions .getExecutor ();
209201 SerializingExecutor serializingExecutor = new SerializingExecutor (callExecutor );
210202
211203 ExternalProcessorGrpc .ExternalProcessorStub extProcStub = ExternalProcessorGrpc .newStub (
@@ -314,7 +306,7 @@ private static class DataPlaneClientCall
314306 private final Queue <ByteString > pendingRequestBodyMessages =
315307 new java .util .concurrent .ConcurrentLinkedQueue <>();
316308
317- // Buffered mutated response bodies from sidecar
309+ // Buffered mutated response bodies from ext_proc server
318310 private int downstreamRequestsPending = 0 ;
319311 private final Queue <ByteString > pendingMutatedResponseBodies =
320312 new java .util .concurrent .ConcurrentLinkedQueue <>();
@@ -488,10 +480,12 @@ public void onNext(ProcessingResponse response) {
488480 upstreamToSidestreamWindow += update .getWindowIncrementUpstreamToSidestream ();
489481 drainPendingRequestBodyMessages ();
490482 drainPendingRequests ();
483+ if (wrappedListener != null ) {
484+ wrappedListener .drainSavedMessages ();
485+ }
491486 }
492- if (wrappedListener != null ) {
493- wrappedListener .drainSavedMessages ();
494- }
487+ // If isReady() becomes true (depends on updated downstreamToSidestreamWindow),
488+ // notify the client application via onReadyNotify() (runs unlocked).
495489 if (!wasReady && isReady ()) {
496490 onReadyNotify ();
497491 }
@@ -811,7 +805,7 @@ void drainPendingRequests() {
811805 }
812806
813807 // Normal mode flow control: pull 1 message at a time
814- if (isSidecarReady () && upstreamToSidestreamWindow > 0 && pendingRequests .get () > 0 ) {
808+ if (isExtProcReady () && upstreamToSidestreamWindow > 0 && pendingRequests .get () > 0 ) {
815809 super .request (1 );
816810 pendingRequests .decrementAndGet ();
817811 }
@@ -863,7 +857,7 @@ private void onReadyNotify() {
863857 wrappedListener .onReadyNotify ();
864858 }
865859
866- boolean isSidecarReady () {
860+ boolean isExtProcReady () {
867861 ExtProcStreamState state = extProcStreamState .get ();
868862 if (state .isCompleted ()) {
869863 return true ;
@@ -889,11 +883,11 @@ public boolean isReady() {
889883 return false ;
890884 }
891885 synchronized (streamLock ) {
892- boolean sidecarReady = isSidecarReady ();
886+ boolean extProcReady = isExtProcReady ();
893887 if (config .getObservabilityMode ()) {
894- return super .isReady () && sidecarReady ;
888+ return super .isReady () && extProcReady ;
895889 }
896- return downstreamToSidestreamWindow > 0 && sidecarReady
890+ return downstreamToSidestreamWindow > 0 && extProcReady
897891 && pendingRequestBodyMessages .isEmpty ();
898892 }
899893 }
@@ -904,22 +898,34 @@ public void request(int numMessages) {
904898 super .request (numMessages );
905899 return ;
906900 }
907- if (!config .getObservabilityMode ()
908- && currentProcessingMode .getResponseBodyMode () != ProcessingMode .BodySendMode .GRPC ) {
909- super .request (numMessages );
910- return ;
911- }
912- if (config .getObservabilityMode ()
913- || currentProcessingMode .getResponseBodyMode () != ProcessingMode .BodySendMode .GRPC ) {
914- super .request (numMessages );
915- return ;
916- }
917901 synchronized (streamLock ) {
918- pendingRequests .addAndGet (numMessages );
919- downstreamRequestsPending += numMessages ;
920- if (isSidecarReady ()) {
921- drainPendingMutatedResponseBodies ();
922- drainPendingRequests ();
902+ boolean sendResponseBodiesToExtProc = config .getObservabilityMode ()
903+ || currentProcessingMode .getResponseBodyMode () == ProcessingMode .BodySendMode .GRPC ;
904+
905+ if (!sendResponseBodiesToExtProc ) {
906+ // We do not send response bodies to ext_proc server at all. Bypassed.
907+ super .request (numMessages );
908+ return ;
909+ }
910+
911+ // We send response bodies to ext_proc server (either in normal GRPC mode or observability mode).
912+ // Gated by ext_proc server readiness.
913+ boolean normalFlowControl = !config .getObservabilityMode (); // i.e. normal GRPC response body mode
914+
915+ if (normalFlowControl ) {
916+ pendingRequests .addAndGet (numMessages );
917+ downstreamRequestsPending += numMessages ;
918+ if (isExtProcReady ()) {
919+ drainPendingMutatedResponseBodies ();
920+ drainPendingRequests ();
921+ }
922+ } else {
923+ // Observability mode: gate on readiness but pull all at once
924+ if (isExtProcReady ()) {
925+ super .request (numMessages );
926+ } else {
927+ pendingRequests .addAndGet (numMessages );
928+ }
923929 }
924930 }
925931 }
@@ -1443,7 +1449,7 @@ public void onMessage(InputStream message) {
14431449
14441450 void drainSavedMessages () {
14451451 synchronized (dataPlaneClientCall .getStreamLock ()) {
1446- while (dataPlaneClientCall .isSidecarReady ()
1452+ while (dataPlaneClientCall .isExtProcReady ()
14471453 && dataPlaneClientCall .upstreamToSidestreamWindow > 0
14481454 && !savedMessages .isEmpty ()) {
14491455 InputStream msg = savedMessages .poll ();
0 commit comments