@@ -96,43 +96,12 @@ public void primeChannel(ManagedChannel managedChannel) {
9696 }
9797
9898 private void primeChannelUnsafe (ManagedChannel managedChannel ) throws IOException {
99- sendPrimeRequests (managedChannel );
99+ sendPrimeRequestsBlocking (managedChannel );
100100 }
101101
102- private void sendPrimeRequests (ManagedChannel managedChannel ) {
102+ private void sendPrimeRequestsBlocking (ManagedChannel managedChannel ) {
103103 try {
104- ClientCall <PingAndWarmRequest , PingAndWarmResponse > clientCall =
105- managedChannel .newCall (
106- BigtableGrpc .getPingAndWarmMethod (),
107- CallOptions .DEFAULT
108- .withCallCredentials (callCredentials )
109- .withDeadline (Deadline .after (1 , TimeUnit .MINUTES )));
110-
111- SettableApiFuture <PingAndWarmResponse > future = SettableApiFuture .create ();
112- clientCall .start (
113- new ClientCall .Listener <PingAndWarmResponse >() {
114- PingAndWarmResponse response ;
115-
116- @ Override
117- public void onMessage (PingAndWarmResponse message ) {
118- response = message ;
119- }
120-
121- @ Override
122- public void onClose (Status status , Metadata trailers ) {
123- if (status .isOk ()) {
124- future .set (response );
125- } else {
126- future .setException (status .asException ());
127- }
128- }
129- },
130- createMetadata (headers , request ));
131- clientCall .sendMessage (request );
132- clientCall .halfClose ();
133- clientCall .request (Integer .MAX_VALUE );
134-
135- future .get (1 , TimeUnit .MINUTES );
104+ sendPrimeRequestsAsync (managedChannel ).get (1 , TimeUnit .MINUTES );
136105 } catch (Throwable e ) {
137106 // TODO: Not sure if we should swallow the error here. We are pre-emptively swapping
138107 // channels if the new
@@ -141,6 +110,53 @@ public void onClose(Status status, Metadata trailers) {
141110 }
142111 }
143112
113+ public SettableApiFuture <PingAndWarmResponse > sendPrimeRequestsAsync (
114+ ManagedChannel managedChannel ) {
115+ ClientCall <PingAndWarmRequest , PingAndWarmResponse > clientCall =
116+ managedChannel .newCall (
117+ BigtableGrpc .getPingAndWarmMethod (),
118+ CallOptions .DEFAULT
119+ .withCallCredentials (callCredentials )
120+ .withDeadline (Deadline .after (1 , TimeUnit .MINUTES )));
121+
122+ SettableApiFuture <PingAndWarmResponse > future = SettableApiFuture .create ();
123+ clientCall .start (
124+ new ClientCall .Listener <PingAndWarmResponse >() {
125+ private PingAndWarmResponse response ;
126+
127+ @ Override
128+ public void onMessage (PingAndWarmResponse message ) {
129+ response = message ;
130+ }
131+
132+ @ Override
133+ public void onClose (Status status , Metadata trailers ) {
134+ if (status .isOk ()) {
135+ future .set (response );
136+ } else {
137+ // Propagate the gRPC error to the future.
138+ future .setException (status .asException (trailers ));
139+ }
140+ }
141+ },
142+ createMetadata (headers , request ));
143+
144+ try {
145+ // Send the request message.
146+ clientCall .sendMessage (request );
147+ // Signal that no more messages will be sent.
148+ clientCall .halfClose ();
149+ // Request the response from the server.
150+ clientCall .request (Integer .MAX_VALUE );
151+ } catch (Throwable t ) {
152+ // If sending fails, cancel the call and notify the future.
153+ clientCall .cancel ("Failed to send priming request" , t );
154+ future .setException (t );
155+ }
156+
157+ return future ;
158+ }
159+
144160 private static Metadata createMetadata (Map <String , String > headers , PingAndWarmRequest request ) {
145161 Metadata metadata = new Metadata ();
146162
0 commit comments