Skip to content

Commit feeab1e

Browse files
committed
Add unit tests for JTATSSL triggerEvent and cover closed stream event behavior.
- Added unit tests in ServerImplTest for JumpToApplicationThreadServerStreamListener.triggerEvent. - Added serverStream_triggerEvent_afterClose in AbstractTransportTest to verify events are ignored after stream closure. - Updated Inbound.ServerInbound to check isClosed() before triggering events. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be
1 parent 958fddc commit feeab1e

3 files changed

Lines changed: 79 additions & 0 deletions

File tree

binder/src/main/java/io/grpc/binder/internal/Inbound.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -671,6 +671,9 @@ protected void deliverCloseAbnormal(Status status) {
671671
void triggerEvent(Object event) {
672672
ServerStreamListener localListener;
673673
synchronized (this) {
674+
if (isClosed()) {
675+
return;
676+
}
674677
localListener = listener;
675678
}
676679
if (localListener != null) {

core/src/test/java/io/grpc/internal/ServerImplTest.java

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1689,6 +1689,52 @@ public void onReady_runtimeExceptionCancelsCall() {
16891689
}
16901690
}
16911691

1692+
@Test
1693+
public void triggerEvent_delegatesToListener() {
1694+
JumpToApplicationThreadServerStreamListener listener
1695+
= new JumpToApplicationThreadServerStreamListener(
1696+
executor.getScheduledExecutorService(),
1697+
executor.getScheduledExecutorService(),
1698+
stream,
1699+
Context.ROOT.withCancellation(),
1700+
PerfMark.createTag());
1701+
ServerStreamListener mockListener = mock(ServerStreamListener.class);
1702+
listener.setListener(mockListener);
1703+
1704+
Object event = new Object();
1705+
listener.triggerEvent(event);
1706+
1707+
verify(mockListener, never()).triggerEvent(any());
1708+
1709+
executor.runDueTasks();
1710+
verify(mockListener).triggerEvent(event);
1711+
}
1712+
1713+
@Test
1714+
public void triggerEvent_errorCancelsCall() {
1715+
JumpToApplicationThreadServerStreamListener listener
1716+
= new JumpToApplicationThreadServerStreamListener(
1717+
executor.getScheduledExecutorService(),
1718+
executor.getScheduledExecutorService(),
1719+
stream,
1720+
Context.ROOT.withCancellation(),
1721+
PerfMark.createTag());
1722+
ServerStreamListener mockListener = mock(ServerStreamListener.class);
1723+
listener.setListener(mockListener);
1724+
1725+
TestError expectedT = new TestError();
1726+
doThrow(expectedT).when(mockListener).triggerEvent(any());
1727+
1728+
listener.triggerEvent(new Object());
1729+
try {
1730+
executor.runDueTasks();
1731+
fail("Expected exception");
1732+
} catch (TestError t) {
1733+
assertSame(expectedT, t);
1734+
ensureServerStateNotLeaked();
1735+
}
1736+
}
1737+
16921738
@Test
16931739
public void binaryLogInstalled() throws Exception {
16941740
final SettableFuture<Boolean> intercepted = SettableFuture.create();

core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2117,6 +2117,36 @@ public void serverStream_triggerEvent() throws Exception {
21172117
clientStream.cancel(Status.CANCELLED);
21182118
}
21192119

2120+
@Test
2121+
public void serverStream_triggerEvent_afterClose() throws Exception {
2122+
server.start(serverListener);
2123+
client = newClientTransport(server);
2124+
startTransport(client, mockClientTransportListener);
2125+
MockServerTransportListener serverTransportListener
2126+
= serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS);
2127+
serverTransport = serverTransportListener.transport;
2128+
2129+
ClientStream clientStream = client.newStream(
2130+
methodDescriptor, new Metadata(), callOptions, noopTracers);
2131+
ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase();
2132+
clientStream.start(clientStreamListener);
2133+
2134+
StreamCreation serverStreamCreation
2135+
= serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS);
2136+
ServerStream serverStream = serverStreamCreation.stream;
2137+
ServerStreamListenerBase serverStreamListener = serverStreamCreation.listener;
2138+
2139+
// Close the stream from client side
2140+
clientStream.cancel(Status.CANCELLED);
2141+
2142+
Object event = new Object();
2143+
serverStream.triggerEvent(event);
2144+
2145+
// Verify listener did NOT receive the event
2146+
Object receivedEvent = serverStreamListener.eventQueue.poll(100, TimeUnit.MILLISECONDS);
2147+
assertNull(receivedEvent);
2148+
}
2149+
21202150
/**
21212151
* Helper that simply does an RPC. It can be used similar to a sleep for negative testing: to give
21222152
* time for actions _not_ to happen. Since it is based on doing an actual RPC with actual

0 commit comments

Comments
 (0)