@@ -563,7 +563,11 @@ private void runTest(TestCases testCase) throws Exception {
563563 tester .testOrcaOob ();
564564 break ;
565565 }
566-
566+
567+ case MAX_CONCURRENT_STREAMS_CONNECTION_SCALING : {
568+ tester .testMcs ();
569+ break ;
570+ }
567571 default :
568572 throw new IllegalArgumentException ("Unknown test case: " + testCase );
569573 }
@@ -596,6 +600,7 @@ private ClientInterceptor maybeCreateAdditionalMetadataInterceptor(
596600 }
597601
598602 private class Tester extends AbstractInteropTest {
603+
599604 @ Override
600605 protected ManagedChannelBuilder <?> createChannelBuilder () {
601606 boolean useGeneric = false ;
@@ -979,31 +984,17 @@ public void testOrcaOob() throws Exception {
979984 .build ();
980985
981986 final int retryLimit = 5 ;
982- BlockingQueue <Object > queue = new LinkedBlockingQueue <>();
983987 final Object lastItem = new Object ();
988+ StreamingOutputCallResponseObserver streamingOutputCallResponseObserver =
989+ new StreamingOutputCallResponseObserver (lastItem );
984990 StreamObserver <StreamingOutputCallRequest > streamObserver =
985- asyncStub .fullDuplexCall (new StreamObserver <StreamingOutputCallResponse >() {
986-
987- @ Override
988- public void onNext (StreamingOutputCallResponse value ) {
989- queue .add (value );
990- }
991-
992- @ Override
993- public void onError (Throwable t ) {
994- queue .add (t );
995- }
996-
997- @ Override
998- public void onCompleted () {
999- queue .add (lastItem );
1000- }
1001- });
991+ asyncStub .fullDuplexCall (streamingOutputCallResponseObserver );
1002992
1003993 streamObserver .onNext (StreamingOutputCallRequest .newBuilder ()
1004994 .setOrcaOobReport (answer )
1005995 .addResponseParameters (ResponseParameters .newBuilder ().setSize (1 ).build ()).build ());
1006- assertThat (queue .take ()).isInstanceOf (StreamingOutputCallResponse .class );
996+ assertThat (streamingOutputCallResponseObserver .take ())
997+ .isInstanceOf (StreamingOutputCallResponse .class );
1007998 int i = 0 ;
1008999 for (; i < retryLimit ; i ++) {
10091000 Thread .sleep (1000 );
@@ -1016,7 +1007,8 @@ public void onCompleted() {
10161007 streamObserver .onNext (StreamingOutputCallRequest .newBuilder ()
10171008 .setOrcaOobReport (answer2 )
10181009 .addResponseParameters (ResponseParameters .newBuilder ().setSize (1 ).build ()).build ());
1019- assertThat (queue .take ()).isInstanceOf (StreamingOutputCallResponse .class );
1010+ assertThat (streamingOutputCallResponseObserver .take ())
1011+ .isInstanceOf (StreamingOutputCallResponse .class );
10201012
10211013 for (i = 0 ; i < retryLimit ; i ++) {
10221014 Thread .sleep (1000 );
@@ -1027,7 +1019,7 @@ public void onCompleted() {
10271019 }
10281020 assertThat (i ).isLessThan (retryLimit );
10291021 streamObserver .onCompleted ();
1030- assertThat (queue .take ()).isSameInstanceAs (lastItem );
1022+ assertThat (streamingOutputCallResponseObserver .take ()).isSameInstanceAs (lastItem );
10311023 }
10321024
10331025 @ Override
@@ -1054,6 +1046,85 @@ protected ServerBuilder<?> getHandshakerServerBuilder() {
10541046 protected int operationTimeoutMillis () {
10551047 return 15000 ;
10561048 }
1049+
1050+ class StreamingOutputCallResponseObserver implements
1051+ StreamObserver <StreamingOutputCallResponse > {
1052+ private final Object lastItem ;
1053+ private final BlockingQueue <Object > queue = new LinkedBlockingQueue <>();
1054+
1055+ public StreamingOutputCallResponseObserver (Object lastItem ) {
1056+ this .lastItem = lastItem ;
1057+ }
1058+
1059+ @ Override
1060+ public void onNext (StreamingOutputCallResponse value ) {
1061+ queue .add (value );
1062+ }
1063+
1064+ @ Override
1065+ public void onError (Throwable t ) {
1066+ queue .add (t );
1067+ }
1068+
1069+ @ Override
1070+ public void onCompleted () {
1071+ queue .add (lastItem );
1072+ }
1073+
1074+ Object take () throws InterruptedException {
1075+ return queue .take ();
1076+ }
1077+ }
1078+
1079+ public void testMcs () throws Exception {
1080+ final Object lastItem = new Object ();
1081+ StreamingOutputCallResponseObserver responseObserver1 =
1082+ new StreamingOutputCallResponseObserver (lastItem );
1083+ StreamObserver <StreamingOutputCallRequest > streamObserver1 =
1084+ asyncStub .fullDuplexCall (responseObserver1 );
1085+ StreamingOutputCallRequest request = StreamingOutputCallRequest .newBuilder ()
1086+ .addResponseParameters (ResponseParameters .newBuilder ()
1087+ .setFillPeerSocketAddress (
1088+ Messages .BoolValue .newBuilder ().setValue (true ).build ())
1089+ .build ())
1090+ .build ();
1091+ streamObserver1 .onNext (request );
1092+ Object responseObj = responseObserver1 .take ();
1093+ StreamingOutputCallResponse callResponse = (StreamingOutputCallResponse ) responseObj ;
1094+ String clientSocketAddressInCall1 = callResponse .getPeerSocketAddress ();
1095+ assertThat (clientSocketAddressInCall1 ).isNotEmpty ();
1096+
1097+ StreamingOutputCallResponseObserver responseObserver2 =
1098+ new StreamingOutputCallResponseObserver (lastItem );
1099+ StreamObserver <StreamingOutputCallRequest > streamObserver2 =
1100+ asyncStub .fullDuplexCall (responseObserver2 );
1101+ streamObserver2 .onNext (request );
1102+ callResponse = (StreamingOutputCallResponse ) responseObserver2 .take ();
1103+ String clientSocketAddressInCall2 = callResponse .getPeerSocketAddress ();
1104+
1105+ assertThat (clientSocketAddressInCall1 ).isEqualTo (clientSocketAddressInCall2 );
1106+
1107+ // The first connection is at max rpc call count of 2, so the 3rd rpc will cause a new
1108+ // connection to be created in the same subchannel and not get queued.
1109+ StreamingOutputCallResponseObserver responseObserver3 =
1110+ new StreamingOutputCallResponseObserver (lastItem );
1111+ StreamObserver <StreamingOutputCallRequest > streamObserver3 =
1112+ asyncStub .fullDuplexCall (responseObserver3 );
1113+ streamObserver3 .onNext (request );
1114+ callResponse = (StreamingOutputCallResponse ) responseObserver3 .take ();
1115+ String clientSocketAddressInCall3 = callResponse .getPeerSocketAddress ();
1116+
1117+ // This assertion is currently failing because connection scaling when MCS limit has been
1118+ // reached is not yet implemented in gRPC Java.
1119+ assertThat (clientSocketAddressInCall3 ).isNotEqualTo (clientSocketAddressInCall1 );
1120+
1121+ streamObserver1 .onCompleted ();
1122+ streamObserver2 .onCompleted ();
1123+ streamObserver3 .onCompleted ();
1124+ assertThat (responseObserver1 .take ()).isSameInstanceAs (lastItem );
1125+ assertThat (responseObserver2 .take ()).isSameInstanceAs (lastItem );
1126+ assertThat (responseObserver3 .take ()).isSameInstanceAs (lastItem );
1127+ }
10571128 }
10581129
10591130 private static String validTestCasesHelpText () {
0 commit comments