3131import io .grpc .Status ;
3232import io .grpc .StatusRuntimeException ;
3333import io .grpc .testing .TestMethodDescriptors ;
34+ import java .util .ArrayList ;
3435import java .util .Arrays ;
36+ import java .util .Collections ;
37+ import java .util .List ;
3538import org .junit .Test ;
3639import org .junit .runner .RunWith ;
3740import org .junit .runners .JUnit4 ;
@@ -183,6 +186,7 @@ public void onHalfClose() {
183186 getSoleMethod (intercepted ).getServerCallHandler ().startCall (call , headers );
184187 callDoubleSreListener .onMessage (null ); // the only close with our exception
185188 callDoubleSreListener .onHalfClose (); // should not trigger a close
189+ callDoubleSreListener .onEvent (new Object ()); // should not trigger a close
186190
187191 // this listener closes the call when it is initialized with startCall
188192 listener = new VoidCallListener () {
@@ -195,13 +199,116 @@ public void onCall(ServerCall<Void, Void> call, Metadata headers) {
195199 public void onHalfClose () {
196200 throw exception ;
197201 }
202+
203+ @ Override
204+ public void onEvent (Object event ) {
205+ throw exception ;
206+ }
198207 };
199208
200209 ServerCall .Listener <Void > callClosedListener =
201210 getSoleMethod (intercepted ).getServerCallHandler ().startCall (call , headers );
202211 // call is already closed, does not match exception
203212 callClosedListener .onHalfClose (); // should not trigger a close
213+ callClosedListener .onEvent (new Object ()); // should not trigger a close
214+ assertEquals (1 , call .numCloses );
215+ }
216+
217+ @ Test
218+ public void statusRuntimeExceptionTransmitter_onEvent_transmitsStatusAndTrailers () {
219+ final Status expectedStatus = Status .RESOURCE_EXHAUSTED .withDescription ("rate limited" );
220+ final Metadata expectedMetadata = new Metadata ();
221+ Metadata .Key <String > key =
222+ Metadata .Key .of ("custom-trailer" , Metadata .ASCII_STRING_MARSHALLER );
223+ expectedMetadata .put (key , "val" );
224+
225+ final java .util .concurrent .atomic .AtomicReference <Status > closedStatus =
226+ new java .util .concurrent .atomic .AtomicReference <>();
227+ final java .util .concurrent .atomic .AtomicReference <Metadata > closedTrailers =
228+ new java .util .concurrent .atomic .AtomicReference <>();
229+
230+ FakeServerCall <Void , Void > call =
231+ new FakeServerCall <Void , Void >(expectedStatus , expectedMetadata ) {
232+ @ Override
233+ public void close (Status status , Metadata trailers ) {
234+ closedStatus .set (status );
235+ closedTrailers .set (trailers );
236+ super .close (status , trailers );
237+ }
238+ };
239+
240+ final StatusRuntimeException exception =
241+ new StatusRuntimeException (expectedStatus , expectedMetadata );
242+
243+ listener = new VoidCallListener () {
244+ @ Override
245+ public void onEvent (Object event ) {
246+ throw exception ;
247+ }
248+ };
249+
250+ ServerServiceDefinition intercepted = ServerInterceptors .intercept (
251+ serviceDefinition ,
252+ Arrays .asList (TransmitStatusRuntimeExceptionInterceptor .instance ()));
253+
254+ // When onEvent throws StatusRuntimeException, it should close the call with status and trailers
255+ getSoleMethod (intercepted ).getServerCallHandler ().startCall (call , headers ).onEvent ("event" );
256+
204257 assertEquals (1 , call .numCloses );
258+ assertEquals (expectedStatus , closedStatus .get ());
259+ assertEquals ("val" , closedTrailers .get ().get (key ));
260+ }
261+
262+ @ Test
263+ public void statusRuntimeExceptionTransmitter_serializingServerCall_serializesTriggerEvent () {
264+ final List <String > executionOrder = Collections .synchronizedList (new ArrayList <String >());
265+ FakeServerCall <Void , Void > call = new FakeServerCall <Void , Void >(Status .OK , new Metadata ()) {
266+ @ Override
267+ public void sendHeaders (Metadata headers ) {
268+ executionOrder .add ("sendHeaders" );
269+ }
270+
271+ @ Override
272+ public void triggerEvent (Object event ) {
273+ executionOrder .add ("triggerEvent:" + event );
274+ }
275+
276+ @ Override
277+ public void sendMessage (Void message ) {
278+ executionOrder .add ("sendMessage" );
279+ }
280+
281+ @ Override
282+ public void close (Status status , Metadata trailers ) {
283+ executionOrder .add ("close" );
284+ }
285+ };
286+
287+ final java .util .concurrent .atomic .AtomicReference <ServerCall <Void , Void >> interceptedCall =
288+ new java .util .concurrent .atomic .AtomicReference <>();
289+ listener = new VoidCallListener () {
290+ @ Override
291+ public void onCall (ServerCall <Void , Void > call , Metadata headers ) {
292+ interceptedCall .set (call );
293+ }
294+ };
295+
296+ ServerServiceDefinition intercepted = ServerInterceptors .intercept (
297+ serviceDefinition ,
298+ Arrays .asList (TransmitStatusRuntimeExceptionInterceptor .instance ()));
299+ getSoleMethod (intercepted ).getServerCallHandler ().startCall (call , headers );
300+
301+ ServerCall <Void , Void > sc = interceptedCall .get ();
302+ sc .sendHeaders (new Metadata ());
303+ sc .triggerEvent ("event1" );
304+ sc .sendMessage (null );
305+ sc .triggerEvent ("event2" );
306+ sc .close (Status .OK , new Metadata ());
307+
308+ assertEquals (
309+ Arrays .asList (
310+ "sendHeaders" , "triggerEvent:event1" , "sendMessage" , "triggerEvent:event2" , "close" ),
311+ executionOrder );
205312 }
206313
207314 private static class FakeServerCall <ReqT , RespT > extends NoopServerCall <ReqT , RespT > {
0 commit comments