From 45b5af1d349fcf07bd50b2743690a425b15cfbc4 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 3 Aug 2026 16:35:25 +0000 Subject: [PATCH 01/17] Implement custom events framework in gRPC-Java server This adds triggerEvent/onEvent APIs to ServerCall and ServerCall.Listener, routing them through ServerStream transport to ensure thread-safety (especially for SerializeReentrantCallsDirectExecutor). TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- api/src/main/java/io/grpc/ServerCall.java | 25 +++++++++++++++ .../grpc/internal/AbstractServerStream.java | 12 +++++++ .../java/io/grpc/internal/ServerCallImpl.java | 13 ++++++++ .../java/io/grpc/internal/ServerImpl.java | 31 +++++++++++++++++++ .../java/io/grpc/internal/ServerStream.java | 6 ++++ .../grpc/internal/ServerStreamListener.java | 5 +++ .../internal/AbstractServerStreamTest.java | 3 ++ .../io/grpc/internal/ServerCallImplTest.java | 26 ++++++++++++++++ .../internal/ServerStreamListenerBase.java | 10 ++++++ .../io/grpc/inprocess/InProcessTransport.java | 19 ++++++++++++ 10 files changed, 150 insertions(+) diff --git a/api/src/main/java/io/grpc/ServerCall.java b/api/src/main/java/io/grpc/ServerCall.java index 3db8ac30e83..b92e04f02f8 100644 --- a/api/src/main/java/io/grpc/ServerCall.java +++ b/api/src/main/java/io/grpc/ServerCall.java @@ -100,6 +100,19 @@ public void onComplete() {} * another {@code onReady()} callback. */ public void onReady() {} + + /** + * A custom event has been triggered by the call. + * + *

This callback is guaranteed to run on the call's executor, serialized with other + * callbacks (like {@link #onMessage}, {@link #onHalfClose}). This means the implementation + * does not need internal synchronization to access call-specific state. + * + * @param event the triggered event. + */ + public void onEvent(Object event) { + // Default no-op + } } /** @@ -262,6 +275,18 @@ public String getAuthority() { return null; } + /** + * Triggers a custom event to be processed by the listener. + * The event will be delivered to {@link Listener#onEvent(Object)} on the call's executor. + * + *

This method is thread-safe and can be called from any thread. + * + * @param event the event to trigger. + */ + public void triggerEvent(Object event) { + // Default no-op + } + /** * The {@link MethodDescriptor} for the call. */ diff --git a/core/src/main/java/io/grpc/internal/AbstractServerStream.java b/core/src/main/java/io/grpc/internal/AbstractServerStream.java index c468cba978a..bc04ad6d4d9 100644 --- a/core/src/main/java/io/grpc/internal/AbstractServerStream.java +++ b/core/src/main/java/io/grpc/internal/AbstractServerStream.java @@ -173,6 +173,18 @@ public final void setListener(ServerStreamListener serverStreamListener) { transportState().setListener(serverStreamListener); } + @Override + public final void triggerEvent(final Object event) { + transportState().runOnTransportThread(new Runnable() { + @Override + public void run() { + if (transportState().listener != null) { + transportState().listener.triggerEvent(event); + } + } + }); + } + @Override public StatsTraceContext statsTraceContext() { return statsTraceCtx; diff --git a/core/src/main/java/io/grpc/internal/ServerCallImpl.java b/core/src/main/java/io/grpc/internal/ServerCallImpl.java index e224384ce8f..6e894371be1 100644 --- a/core/src/main/java/io/grpc/internal/ServerCallImpl.java +++ b/core/src/main/java/io/grpc/internal/ServerCallImpl.java @@ -254,6 +254,11 @@ public MethodDescriptor getMethodDescriptor() { return method; } + @Override + public void triggerEvent(Object event) { + stream.triggerEvent(event); + } + @Override public SecurityLevel getSecurityLevel() { final Attributes attributes = getAttributes(); @@ -395,5 +400,13 @@ public void onReady() { listener.onReady(); } } + + @Override + public void triggerEvent(Object event) { + if (call.cancelled) { + return; + } + listener.onEvent(event); + } } } diff --git a/core/src/main/java/io/grpc/internal/ServerImpl.java b/core/src/main/java/io/grpc/internal/ServerImpl.java index d9f64c2d473..767d85f443b 100644 --- a/core/src/main/java/io/grpc/internal/ServerImpl.java +++ b/core/src/main/java/io/grpc/internal/ServerImpl.java @@ -781,6 +781,9 @@ public void closed(Status status) {} @Override public void onReady() {} + + @Override + public void triggerEvent(Object event) {} } /** @@ -960,6 +963,34 @@ public void runInContext() { callExecutor.execute(new OnReady()); } } + + @Override + public void triggerEvent(final Object event) { + try (TaskCloseable ignore = PerfMark.traceTask("ServerStreamListener.triggerEvent")) { + PerfMark.attachTag(tag); + final Link link = PerfMark.linkOut(); + + final class TriggerEvent extends ContextRunnable { + TriggerEvent() { + super(context); + } + + @Override + public void runInContext() { + try (TaskCloseable ignore = PerfMark.traceTask("ServerCallListener(app).onEvent")) { + PerfMark.attachTag(tag); + PerfMark.linkIn(link); + getListener().triggerEvent(event); + } catch (Throwable t) { + internalClose(t); + throw t; + } + } + } + + callExecutor.execute(new TriggerEvent()); + } + } } @VisibleForTesting diff --git a/core/src/main/java/io/grpc/internal/ServerStream.java b/core/src/main/java/io/grpc/internal/ServerStream.java index aa5ba10329c..4c88e94e4d7 100644 --- a/core/src/main/java/io/grpc/internal/ServerStream.java +++ b/core/src/main/java/io/grpc/internal/ServerStream.java @@ -87,6 +87,12 @@ public interface ServerStream extends Stream { */ void setListener(ServerStreamListener serverStreamListener); + /** + * Triggers a custom event. Implementations must ensure this is propagated to the + * listener on the transport thread. + */ + void triggerEvent(Object event); + /** * The context for recording stats and traces for this stream. */ diff --git a/core/src/main/java/io/grpc/internal/ServerStreamListener.java b/core/src/main/java/io/grpc/internal/ServerStreamListener.java index e55217ab422..74de0f2079e 100644 --- a/core/src/main/java/io/grpc/internal/ServerStreamListener.java +++ b/core/src/main/java/io/grpc/internal/ServerStreamListener.java @@ -42,4 +42,9 @@ public interface ServerStreamListener extends StreamListener { * @param status details about the remote closure */ void closed(Status status); + + /** + * Propagates a custom event to the listener. Must be called on the transport thread. + */ + void triggerEvent(Object event); } diff --git a/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java index 137ba19bfea..93030d6936f 100644 --- a/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java @@ -391,6 +391,9 @@ public void halfClosed() {} @Override public void closed(Status status) {} + + @Override + public void triggerEvent(Object event) {} } private static class AbstractServerStreamBase extends AbstractServerStream { diff --git a/core/src/test/java/io/grpc/internal/ServerCallImplTest.java b/core/src/test/java/io/grpc/internal/ServerCallImplTest.java index 7394c83eab2..4a2de9f3936 100644 --- a/core/src/test/java/io/grpc/internal/ServerCallImplTest.java +++ b/core/src/test/java/io/grpc/internal/ServerCallImplTest.java @@ -493,6 +493,32 @@ public void streamListener_unexpectedRuntimeException() { assertThat(e).hasMessageThat().isEqualTo("unexpected exception"); } + @Test + public void triggerEvent_propagatesToStream() { + Object event = new Object(); + call.triggerEvent(event); + verify(stream).triggerEvent(event); + } + + @Test + public void streamListener_triggerEvent() { + ServerStreamListenerImpl streamListener = + new ServerCallImpl.ServerStreamListenerImpl<>(call, callListener, context); + Object event = new Object(); + streamListener.triggerEvent(event); + verify(callListener).onEvent(event); + } + + @Test + public void streamListener_triggerEvent_cancelled() { + ServerStreamListenerImpl streamListener = + new ServerCallImpl.ServerStreamListenerImpl<>(call, callListener, context); + Object event = new Object(); + streamListener.closed(Status.CANCELLED); + streamListener.triggerEvent(event); + verify(callListener, never()).onEvent(event); + } + private static class LongMarshaller implements Marshaller { @Override public InputStream stream(Long value) { diff --git a/core/src/testFixtures/java/io/grpc/internal/ServerStreamListenerBase.java b/core/src/testFixtures/java/io/grpc/internal/ServerStreamListenerBase.java index aaa70600542..e4ac01912e4 100644 --- a/core/src/testFixtures/java/io/grpc/internal/ServerStreamListenerBase.java +++ b/core/src/testFixtures/java/io/grpc/internal/ServerStreamListenerBase.java @@ -89,6 +89,8 @@ public void halfClosed() { halfClosedLatch.countDown(); } + public final BlockingQueue eventQueue = new LinkedBlockingQueue<>(); + @Override public void closed(Status status) { if (this.status.isDone()) { @@ -96,4 +98,12 @@ public void closed(Status status) { } this.status.set(status); } + + @Override + public void triggerEvent(Object event) { + if (this.status.isDone()) { + fail("triggerEvent invoked after closed"); + } + eventQueue.add(event); + } } diff --git a/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java b/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java index a92f10fd5c5..dc3a970d156 100644 --- a/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java +++ b/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java @@ -429,6 +429,11 @@ public void setListener(ServerStreamListener serverStreamListener) { clientStream.setListener(serverStreamListener); } + @Override + public void triggerEvent(Object event) { + clientStream.triggerServerEvent(event); + } + @Override public void request(int numMessages) { boolean onReady = clientStream.serverRequested(numMessages); @@ -732,6 +737,20 @@ private synchronized void setListener(ServerStreamListener listener) { this.serverStreamListener = listener; } + void triggerServerEvent(final Object event) { + synchronized (this) { + if (!closed && serverStreamListener != null) { + syncContext.executeLater(new Runnable() { + @Override + public void run() { + serverStreamListener.triggerEvent(event); + } + }); + } + } + syncContext.drain(); + } + @Override public void request(int numMessages) { boolean onReady = serverStream.clientRequested(numMessages); From e37245caad81c2a703cb7652c8db21365bc2966b Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 07:34:06 +0000 Subject: [PATCH 02/17] Fix forwarding listeners to propagate custom events --- api/src/main/java/io/grpc/Contexts.java | 10 ++++++++++ .../main/java/io/grpc/PartialForwardingServerCall.java | 5 +++++ .../io/grpc/PartialForwardingServerCallListener.java | 5 +++++ 3 files changed, 20 insertions(+) diff --git a/api/src/main/java/io/grpc/Contexts.java b/api/src/main/java/io/grpc/Contexts.java index c62ffc80a38..9c3697dd6d4 100644 --- a/api/src/main/java/io/grpc/Contexts.java +++ b/api/src/main/java/io/grpc/Contexts.java @@ -118,6 +118,16 @@ public void onReady() { context.detach(previous); } } + + @Override + public void onEvent(Object event) { + Context previous = context.attach(); + try { + super.onEvent(event); + } finally { + context.detach(previous); + } + } } /** diff --git a/api/src/main/java/io/grpc/PartialForwardingServerCall.java b/api/src/main/java/io/grpc/PartialForwardingServerCall.java index a313407b23e..8c8f53cf93c 100644 --- a/api/src/main/java/io/grpc/PartialForwardingServerCall.java +++ b/api/src/main/java/io/grpc/PartialForwardingServerCall.java @@ -87,6 +87,11 @@ public SecurityLevel getSecurityLevel() { return delegate().getSecurityLevel(); } + @Override + public void triggerEvent(Object event) { + delegate().triggerEvent(event); + } + @Override public String toString() { return MoreObjects.toStringHelper(this).add("delegate", delegate()).toString(); diff --git a/api/src/main/java/io/grpc/PartialForwardingServerCallListener.java b/api/src/main/java/io/grpc/PartialForwardingServerCallListener.java index ca2fd0058c9..23e93bb065e 100644 --- a/api/src/main/java/io/grpc/PartialForwardingServerCallListener.java +++ b/api/src/main/java/io/grpc/PartialForwardingServerCallListener.java @@ -50,6 +50,11 @@ public void onReady() { delegate().onReady(); } + @Override + public void onEvent(Object event) { + delegate().onEvent(event); + } + @Override public String toString() { return MoreObjects.toStringHelper(this).add("delegate", delegate()).toString(); From 2636ebb0ade2c02ca936107940f364df8275e77a Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 07:23:42 +0000 Subject: [PATCH 03/17] Address review from server interceptor executor design comments. --- api/src/main/java/io/grpc/ServerCall.java | 3 ++- .../java/io/grpc/internal/AbstractServerStream.java | 11 ++++++++--- 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/api/src/main/java/io/grpc/ServerCall.java b/api/src/main/java/io/grpc/ServerCall.java index b92e04f02f8..04c335fc380 100644 --- a/api/src/main/java/io/grpc/ServerCall.java +++ b/api/src/main/java/io/grpc/ServerCall.java @@ -279,7 +279,8 @@ public String getAuthority() { * Triggers a custom event to be processed by the listener. * The event will be delivered to {@link Listener#onEvent(Object)} on the call's executor. * - *

This method is thread-safe and can be called from any thread. + *

This method is thread-safe and can be called from any thread. No events will be delivered + * after the RPC is cancelled or completed. * * @param event the event to trigger. */ diff --git a/core/src/main/java/io/grpc/internal/AbstractServerStream.java b/core/src/main/java/io/grpc/internal/AbstractServerStream.java index bc04ad6d4d9..67dfdc93d42 100644 --- a/core/src/main/java/io/grpc/internal/AbstractServerStream.java +++ b/core/src/main/java/io/grpc/internal/AbstractServerStream.java @@ -178,9 +178,7 @@ public final void triggerEvent(final Object event) { transportState().runOnTransportThread(new Runnable() { @Override public void run() { - if (transportState().listener != null) { - transportState().listener.triggerEvent(event); - } + transportState().triggerEvent(event); } }); } @@ -271,6 +269,13 @@ public void deframerClosed(boolean hasPartialMessage) { + public final void triggerEvent(Object event) { + if (listenerClosed) { + return; + } + listener().triggerEvent(event); + } + @Override protected ServerStreamListener listener() { return listener; From 5caf23719810dd93eca41403add43f6e52880602 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 08:43:33 +0000 Subject: [PATCH 04/17] Add @ExperimentalApi annotation to triggerEvent and onEvent (issue #12979) --- api/src/main/java/io/grpc/ServerCall.java | 2 ++ .../src/main/java/io/grpc/inprocess/InProcessTransport.java | 2 +- 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/api/src/main/java/io/grpc/ServerCall.java b/api/src/main/java/io/grpc/ServerCall.java index 04c335fc380..2e6ee07a23f 100644 --- a/api/src/main/java/io/grpc/ServerCall.java +++ b/api/src/main/java/io/grpc/ServerCall.java @@ -110,6 +110,7 @@ public void onReady() {} * * @param event the triggered event. */ + @ExperimentalApi("https://github.com/grpc/grpc-java/issues/12979") public void onEvent(Object event) { // Default no-op } @@ -284,6 +285,7 @@ public String getAuthority() { * * @param event the event to trigger. */ + @ExperimentalApi("https://github.com/grpc/grpc-java/issues/12979") public void triggerEvent(Object event) { // Default no-op } diff --git a/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java b/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java index dc3a970d156..57820b396ad 100644 --- a/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java +++ b/inprocess/src/main/java/io/grpc/inprocess/InProcessTransport.java @@ -739,7 +739,7 @@ private synchronized void setListener(ServerStreamListener listener) { void triggerServerEvent(final Object event) { synchronized (this) { - if (!closed && serverStreamListener != null) { + if (!closed) { syncContext.executeLater(new Runnable() { @Override public void run() { From 2158b976295270e8bd7a65cca8403f4a5ec6f412 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 10:13:49 +0000 Subject: [PATCH 05/17] Add unit test for server stream custom events in AbstractTransportTest --- .../grpc/internal/AbstractTransportTest.java | 29 +++++++++++++++++++ 1 file changed, 29 insertions(+) diff --git a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java index 5d07de32df9..5127c7c2b0f 100644 --- a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java +++ b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java @@ -2088,6 +2088,35 @@ public void clientChecksInboundMetadataSize_trailer() throws Exception { assertNull(metadata.get(tellTaleKey)); } + @Test + public void serverStream_triggerEvent() throws Exception { + server.start(serverListener); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener + = serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), callOptions, noopTracers); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + + StreamCreation serverStreamCreation + = serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + ServerStream serverStream = serverStreamCreation.stream; + ServerStreamListenerBase serverStreamListener = serverStreamCreation.listener; + + Object event = new Object(); + serverStream.triggerEvent(event); + + Object receivedEvent = serverStreamListener.eventQueue.poll(TIMEOUT_MS, TimeUnit.MILLISECONDS); + assertEquals(event, receivedEvent); + + // Cleanup + clientStream.cancel(Status.CANCELLED); + } + /** * Helper that simply does an RPC. It can be used similar to a sleep for negative testing: to give * time for actions _not_ to happen. Since it is based on doing an actual RPC with actual From 605cc03e0cdacdfb68c12b85f2af50202b89f7b5 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 12:03:48 +0000 Subject: [PATCH 06/17] Fix compilation errors in binder, netty, and okhttp (issue #12979) --- .../io/grpc/binder/internal/MultiMessageServerStream.java | 5 +++++ .../io/grpc/binder/internal/SingleMessageServerStream.java | 5 +++++ .../test/java/io/grpc/netty/NettyClientTransportTest.java | 4 ++++ .../test/java/io/grpc/okhttp/OkHttpServerTransportTest.java | 4 ++++ 4 files changed, 18 insertions(+) diff --git a/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java b/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java index f54769caefa..40cebf40e21 100644 --- a/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java +++ b/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java @@ -175,6 +175,11 @@ public void setDecompressor(Decompressor decompressor) { // Ignore. } + @Override + public void triggerEvent(Object event) { + // Ignore. + } + @Override public void optimizeForDirectExecutor() { // Ignore. diff --git a/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java b/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java index 383bd7a2593..d9b177132d3 100644 --- a/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java +++ b/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java @@ -167,6 +167,11 @@ public void setDecompressor(Decompressor decompressor) { // Ignore. } + @Override + public void triggerEvent(Object event) { + // Ignore. + } + @Override public void optimizeForDirectExecutor() { // Ignore. diff --git a/netty/src/test/java/io/grpc/netty/NettyClientTransportTest.java b/netty/src/test/java/io/grpc/netty/NettyClientTransportTest.java index ef8d2e5efda..b22be2460f1 100644 --- a/netty/src/test/java/io/grpc/netty/NettyClientTransportTest.java +++ b/netty/src/test/java/io/grpc/netty/NettyClientTransportTest.java @@ -1343,6 +1343,10 @@ public void halfClosed() { @Override public void closed(Status status) { } + + @Override + public void triggerEvent(Object event) { + } } private final class EchoServerListener implements ServerListener { diff --git a/okhttp/src/test/java/io/grpc/okhttp/OkHttpServerTransportTest.java b/okhttp/src/test/java/io/grpc/okhttp/OkHttpServerTransportTest.java index 00db6e1d339..d3e50cf3821 100644 --- a/okhttp/src/test/java/io/grpc/okhttp/OkHttpServerTransportTest.java +++ b/okhttp/src/test/java/io/grpc/okhttp/OkHttpServerTransportTest.java @@ -1519,6 +1519,10 @@ public void closed(Status status) { public void onReady() { } + @Override + public void triggerEvent(Object event) { + } + static String getContent(InputStream message) throws IOException { try { return new String(ByteStreams.toByteArray(message), UTF_8); From dd549f391c47254d343067d7d590450e5baeba77 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 12:52:46 +0000 Subject: [PATCH 07/17] Implement `triggerEvent` for Binder transport. --- .../src/main/java/io/grpc/binder/internal/Inbound.java | 10 ++++++++++ .../grpc/binder/internal/MultiMessageServerStream.java | 2 +- .../binder/internal/SingleMessageServerStream.java | 2 +- 3 files changed, 12 insertions(+), 2 deletions(-) diff --git a/binder/src/main/java/io/grpc/binder/internal/Inbound.java b/binder/src/main/java/io/grpc/binder/internal/Inbound.java index 83fc8273d6f..5671d808114 100644 --- a/binder/src/main/java/io/grpc/binder/internal/Inbound.java +++ b/binder/src/main/java/io/grpc/binder/internal/Inbound.java @@ -668,6 +668,16 @@ protected void deliverCloseAbnormal(Status status) { listener.closed(status); } + void triggerEvent(Object event) { + ServerStreamListener localListener; + synchronized (this) { + localListener = listener; + } + if (localListener != null) { + localListener.triggerEvent(event); + } + } + @GuardedBy("this") void onCloseSent(Status status) { if (!isClosed()) { diff --git a/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java b/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java index 40cebf40e21..7a57138ce22 100644 --- a/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java +++ b/binder/src/main/java/io/grpc/binder/internal/MultiMessageServerStream.java @@ -177,7 +177,7 @@ public void setDecompressor(Decompressor decompressor) { @Override public void triggerEvent(Object event) { - // Ignore. + inbound.triggerEvent(event); } @Override diff --git a/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java b/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java index d9b177132d3..5f1dd511f73 100644 --- a/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java +++ b/binder/src/main/java/io/grpc/binder/internal/SingleMessageServerStream.java @@ -169,7 +169,7 @@ public void setDecompressor(Decompressor decompressor) { @Override public void triggerEvent(Object event) { - // Ignore. + inbound.triggerEvent(event); } @Override From d5f0182a640f5e148178f23e4b647abf8ec7a262 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 10 Aug 2026 13:35:11 +0000 Subject: [PATCH 08/17] Add unit test coverage for ServerCall triggerEvent and custom events framework. - Added unit tests in AbstractServerStreamTest for triggerEvent propagation and close behavior. - Updated ContextsTest to cover onEvent propagation in ContextualizedServerCallListener. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- api/src/test/java/io/grpc/ContextsTest.java | 17 ++++++++++++- .../internal/AbstractServerStreamTest.java | 25 +++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/api/src/test/java/io/grpc/ContextsTest.java b/api/src/test/java/io/grpc/ContextsTest.java index ec9dc3929a2..974b1aff0c2 100644 --- a/api/src/test/java/io/grpc/ContextsTest.java +++ b/api/src/test/java/io/grpc/ContextsTest.java @@ -82,6 +82,11 @@ public void interceptCall_basic() { assertSame(uniqueContext, Context.current()); methodCalls.add(5); } + + @Override public void onEvent(Object event) { + assertSame(uniqueContext, Context.current()); + methodCalls.add(6); + } }; ServerCall.Listener wrapped = interceptCall(uniqueContext, call, headers, new ServerCallHandler() { @@ -101,7 +106,8 @@ public ServerCall.Listener startCall( wrapped.onCancel(); wrapped.onComplete(); wrapped.onReady(); - assertEquals(Arrays.asList(1, 2, 3, 4, 5), methodCalls); + wrapped.onEvent(new Object()); + assertEquals(Arrays.asList(1, 2, 3, 4, 5, 6), methodCalls); assertSame(origContext, Context.current()); } @@ -145,6 +151,10 @@ public void interceptCall_restoresIfListenerThrows() { @Override public void onReady() { throw new RuntimeException(); } + + @Override public void onEvent(Object event) { + throw new RuntimeException(); + } }; ServerCall.Listener wrapped = interceptCall(uniqueContext, call, headers, new ServerCallHandler() { @@ -180,6 +190,11 @@ public ServerCall.Listener startCall( fail("Exception expected"); } catch (RuntimeException expected) { } + try { + wrapped.onEvent(new Object()); + fail("Exception expected"); + } catch (RuntimeException expected) { + } assertSame(origContext, Context.current()); } diff --git a/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java b/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java index 93030d6936f..5defd17fdd0 100644 --- a/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java +++ b/core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java @@ -361,6 +361,31 @@ public void close_sendTrailersClearsReservedFields() { assertEquals("bad", metadataCaptor.getValue().get(InternalStatus.MESSAGE_KEY)); } + @Test + public void triggerEvent_propagatesToListener() { + ServerStreamListener listener = mock(ServerStreamListener.class); + stream.transportState().setListener(listener); + + Object event = new Object(); + stream.triggerEvent(event); + + verify(listener).triggerEvent(event); + } + + @Test + public void triggerEvent_ignoredAfterClose() { + ServerStreamListener listener = mock(ServerStreamListener.class); + stream.transportState().setListener(listener); + + stream.close(Status.OK, new Metadata()); + stream.transportState().complete(); + + Object event = new Object(); + stream.triggerEvent(event); + + verify(listener, never()).triggerEvent(any()); + } + @Test public void changeOnReadyThreshold() { stream.setListener(new ServerStreamListenerBase()); From 958fddc9f74c00ac048a85a965653121b2b3b31b Mon Sep 17 00:00:00 2001 From: Kannan J Date: Tue, 11 Aug 2026 06:00:08 +0000 Subject: [PATCH 09/17] Fix race condition in AsyncSecurityPoliciesTest. Synchronized with the executor before asserting cancellation of the delegate future to ensure that transformAsync has finished processing the delegate future and propagated the cancellation. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- .../src/test/java/io/grpc/binder/AsyncSecurityPoliciesTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/binder/src/test/java/io/grpc/binder/AsyncSecurityPoliciesTest.java b/binder/src/test/java/io/grpc/binder/AsyncSecurityPoliciesTest.java index b0d84f1be74..e556954fd40 100644 --- a/binder/src/test/java/io/grpc/binder/AsyncSecurityPoliciesTest.java +++ b/binder/src/test/java/io/grpc/binder/AsyncSecurityPoliciesTest.java @@ -289,6 +289,7 @@ public ListenableFuture checkAuthorizationAsync(int uid) { ListenableFuture authFuture = asyncPolicy.checkAuthorizationAsync(SOME_UID); assertThat(awaitResult(settableUid)).isEqualTo(SOME_UID); authFuture.cancel(false); + executor.submit(() -> {}).get(10, TimeUnit.SECONDS); assertThat(delegateAuthFuture.isCancelled()).isTrue(); } From feeab1ee1742ae572da986d79c703857d52e0143 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Thu, 13 Aug 2026 08:46:09 +0000 Subject: [PATCH 10/17] 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 --- .../java/io/grpc/binder/internal/Inbound.java | 3 ++ .../java/io/grpc/internal/ServerImplTest.java | 46 +++++++++++++++++++ .../grpc/internal/AbstractTransportTest.java | 30 ++++++++++++ 3 files changed, 79 insertions(+) diff --git a/binder/src/main/java/io/grpc/binder/internal/Inbound.java b/binder/src/main/java/io/grpc/binder/internal/Inbound.java index 5671d808114..83decf4a89a 100644 --- a/binder/src/main/java/io/grpc/binder/internal/Inbound.java +++ b/binder/src/main/java/io/grpc/binder/internal/Inbound.java @@ -671,6 +671,9 @@ protected void deliverCloseAbnormal(Status status) { void triggerEvent(Object event) { ServerStreamListener localListener; synchronized (this) { + if (isClosed()) { + return; + } localListener = listener; } if (localListener != null) { diff --git a/core/src/test/java/io/grpc/internal/ServerImplTest.java b/core/src/test/java/io/grpc/internal/ServerImplTest.java index 91969dd6910..da2e9646042 100644 --- a/core/src/test/java/io/grpc/internal/ServerImplTest.java +++ b/core/src/test/java/io/grpc/internal/ServerImplTest.java @@ -1689,6 +1689,52 @@ public void onReady_runtimeExceptionCancelsCall() { } } + @Test + public void triggerEvent_delegatesToListener() { + JumpToApplicationThreadServerStreamListener listener + = new JumpToApplicationThreadServerStreamListener( + executor.getScheduledExecutorService(), + executor.getScheduledExecutorService(), + stream, + Context.ROOT.withCancellation(), + PerfMark.createTag()); + ServerStreamListener mockListener = mock(ServerStreamListener.class); + listener.setListener(mockListener); + + Object event = new Object(); + listener.triggerEvent(event); + + verify(mockListener, never()).triggerEvent(any()); + + executor.runDueTasks(); + verify(mockListener).triggerEvent(event); + } + + @Test + public void triggerEvent_errorCancelsCall() { + JumpToApplicationThreadServerStreamListener listener + = new JumpToApplicationThreadServerStreamListener( + executor.getScheduledExecutorService(), + executor.getScheduledExecutorService(), + stream, + Context.ROOT.withCancellation(), + PerfMark.createTag()); + ServerStreamListener mockListener = mock(ServerStreamListener.class); + listener.setListener(mockListener); + + TestError expectedT = new TestError(); + doThrow(expectedT).when(mockListener).triggerEvent(any()); + + listener.triggerEvent(new Object()); + try { + executor.runDueTasks(); + fail("Expected exception"); + } catch (TestError t) { + assertSame(expectedT, t); + ensureServerStateNotLeaked(); + } + } + @Test public void binaryLogInstalled() throws Exception { final SettableFuture intercepted = SettableFuture.create(); diff --git a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java index 5127c7c2b0f..88b680b996f 100644 --- a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java +++ b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java @@ -2117,6 +2117,36 @@ public void serverStream_triggerEvent() throws Exception { clientStream.cancel(Status.CANCELLED); } + @Test + public void serverStream_triggerEvent_afterClose() throws Exception { + server.start(serverListener); + client = newClientTransport(server); + startTransport(client, mockClientTransportListener); + MockServerTransportListener serverTransportListener + = serverListener.takeListenerOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + serverTransport = serverTransportListener.transport; + + ClientStream clientStream = client.newStream( + methodDescriptor, new Metadata(), callOptions, noopTracers); + ClientStreamListenerBase clientStreamListener = new ClientStreamListenerBase(); + clientStream.start(clientStreamListener); + + StreamCreation serverStreamCreation + = serverTransportListener.takeStreamOrFail(TIMEOUT_MS, TimeUnit.MILLISECONDS); + ServerStream serverStream = serverStreamCreation.stream; + ServerStreamListenerBase serverStreamListener = serverStreamCreation.listener; + + // Close the stream from client side + clientStream.cancel(Status.CANCELLED); + + Object event = new Object(); + serverStream.triggerEvent(event); + + // Verify listener did NOT receive the event + Object receivedEvent = serverStreamListener.eventQueue.poll(100, TimeUnit.MILLISECONDS); + assertNull(receivedEvent); + } + /** * Helper that simply does an RPC. It can be used similar to a sleep for negative testing: to give * time for actions _not_ to happen. Since it is based on doing an actual RPC with actual From b9e1e2be0805cba2cd24c9a3bee0853d312e2e98 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Thu, 13 Aug 2026 09:15:02 +0000 Subject: [PATCH 11/17] Fix race condition in serverStream_triggerEvent_afterClose test. Wait for the server stream to be fully closed (via awaitClose) before calling triggerEvent, to ensure the transport has processed the cancellation and marked the listener as closed. This fixes flakiness in slower transports like Jetty. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- .../java/io/grpc/internal/AbstractTransportTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java index 88b680b996f..3898eb8be30 100644 --- a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java +++ b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java @@ -2139,6 +2139,8 @@ public void serverStream_triggerEvent_afterClose() throws Exception { // Close the stream from client side clientStream.cancel(Status.CANCELLED); + serverStreamListener.awaitClose(TIMEOUT_MS, TimeUnit.MILLISECONDS); + Object event = new Object(); serverStream.triggerEvent(event); From e8324225e963fb4635ce6fb7c636f51959331f42 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 31 Aug 2026 10:23:21 +0000 Subject: [PATCH 12/17] Make Binder transport triggerEvent thread-safe and guarantee ordering. Updated ServerInbound.triggerEvent to invoke the listener's triggerEvent callback inside the synchronized(this) block. This ensures that the check for isClosed() and the invocation of the listener are atomic relative to stream closure (which also runs under the same lock). This prevents a race where triggerEvent could be called on the listener after the stream has been closed, which would result in out-of-order events delivered to the application. This is consistent with how other listener callbacks (like closed and halfClosed) are delivered in Inbound.java. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- binder/src/main/java/io/grpc/binder/internal/Inbound.java | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/binder/src/main/java/io/grpc/binder/internal/Inbound.java b/binder/src/main/java/io/grpc/binder/internal/Inbound.java index 83decf4a89a..a96f94751d8 100644 --- a/binder/src/main/java/io/grpc/binder/internal/Inbound.java +++ b/binder/src/main/java/io/grpc/binder/internal/Inbound.java @@ -669,15 +669,13 @@ protected void deliverCloseAbnormal(Status status) { } void triggerEvent(Object event) { - ServerStreamListener localListener; synchronized (this) { if (isClosed()) { return; } - localListener = listener; - } - if (localListener != null) { - localListener.triggerEvent(event); + if (listener != null) { + listener.triggerEvent(event); + } } } From 48aee63fa85fa4cc63930a7c399d7f77d5acf214 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Wed, 2 Sep 2026 06:47:32 +0000 Subject: [PATCH 13/17] Propagate custom events in PendingAuthListener, TransmitStatusRuntimeExceptionInterceptor, and OpenTelemetryTracingModule. - binder: Implement onEvent in PendingAuthListener to buffer and replay custom events to the delegate once auth completes, preventing events from being dropped. - util: Handle onEvent in TransmitStatusRuntimeExceptionInterceptor listener wrapper to catch StatusRuntimeException and close the call. Serialize triggerEvent on SerializingServerCall's executor. - opentelemetry: Implement onEvent in ContextServerCallListener to attach OpenTelemetry trace context and scope during delegate invocation. - Add unit tests for all updated implementations. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- .../binder/internal/PendingAuthListener.java | 6 ++++ .../internal/PendingAuthListenerTest.java | 4 +++ .../OpenTelemetryTracingModule.java | 7 ++++ .../OpenTelemetryTracingModuleTest.java | 10 ++++++ ...smitStatusRuntimeExceptionInterceptor.java | 19 ++++++++++ .../grpc/util/UtilServerInterceptorsTest.java | 36 ++++++++++++++++++- 6 files changed, 81 insertions(+), 1 deletion(-) diff --git a/binder/src/main/java/io/grpc/binder/internal/PendingAuthListener.java b/binder/src/main/java/io/grpc/binder/internal/PendingAuthListener.java index ad993b8c93b..3cedce837f6 100644 --- a/binder/src/main/java/io/grpc/binder/internal/PendingAuthListener.java +++ b/binder/src/main/java/io/grpc/binder/internal/PendingAuthListener.java @@ -88,6 +88,12 @@ public void onReady() { maybeRunPendingSteps(); } + @Override + public void onEvent(Object event) { + pendingSteps.offer(delegate -> delegate.onEvent(event)); + maybeRunPendingSteps(); + } + /** * Similar to Java8's {@link java.util.function.Consumer}, but redeclared in order to support * Android SDK 21. diff --git a/binder/src/test/java/io/grpc/binder/internal/PendingAuthListenerTest.java b/binder/src/test/java/io/grpc/binder/internal/PendingAuthListenerTest.java index 9cdf123033b..4c17dae302d 100644 --- a/binder/src/test/java/io/grpc/binder/internal/PendingAuthListenerTest.java +++ b/binder/src/test/java/io/grpc/binder/internal/PendingAuthListenerTest.java @@ -45,6 +45,7 @@ public void setUp() { public void onCallbacks_noOpBeforeStartCall() { listener.onReady(); listener.onMessage("foo"); + listener.onEvent("bar"); listener.onHalfClose(); listener.onComplete(); @@ -54,16 +55,19 @@ public void onCallbacks_noOpBeforeStartCall() { @Test public void onCallbacks_runsPendingCallbacksAfterStartCall() { String message = "foo"; + String event = "bar"; // Act 1 listener.onReady(); listener.onMessage(message); + listener.onEvent(event); listener.startCall(call, headers, next); // Assert 1 InOrder order = Mockito.inOrder(delegate); order.verify(delegate).onReady(); order.verify(delegate).onMessage(message); + order.verify(delegate).onEvent(event); // Act 2 listener.onHalfClose(); diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java index 32aab870f0f..088d6dc9845 100644 --- a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java +++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java @@ -452,6 +452,13 @@ public void onReady() { delegate().onReady(); } } + + @Override + public void onEvent(Object event) { + try (Scope scope = context.makeCurrent()) { + delegate().onEvent(event); + } + } } @VisibleForTesting diff --git a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java index 0b5bff1d036..ee7e86e05cc 100644 --- a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java +++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java @@ -948,6 +948,11 @@ public void onCancel() { public void onComplete() { callbackSpan.set(Span.fromContext(Context.current())); } + + @Override + public void onEvent(Object event) { + callbackSpan.set(Span.fromContext(Context.current())); + } }; ServerInterceptor interceptor = tracingModule.getServerSpanPropagationInterceptor(); @SuppressWarnings("unchecked") @@ -967,6 +972,8 @@ public void onComplete() { assertEquals(callbackSpan.get(), Span.getInvalid()); listener.onComplete(); assertEquals(callbackSpan.get(), Span.getInvalid()); + listener.onEvent(new Object()); + assertEquals(callbackSpan.get(), Span.getInvalid()); Span parentSpan = tracerRule.spanBuilder("parent-span").startSpan(); io.grpc.Context context = io.grpc.Context.current().withValue( @@ -990,6 +997,9 @@ public void onComplete() { listener.onComplete(); assertEquals(callbackSpan.get().getSpanContext().getTraceId(), parentSpan.getSpanContext().getTraceId()); + listener.onEvent(new Object()); + assertEquals(callbackSpan.get().getSpanContext().getTraceId(), + parentSpan.getSpanContext().getTraceId()); } finally { context.detach(previous); } diff --git a/util/src/main/java/io/grpc/util/TransmitStatusRuntimeExceptionInterceptor.java b/util/src/main/java/io/grpc/util/TransmitStatusRuntimeExceptionInterceptor.java index b477ae1fdfb..e4b364bd532 100644 --- a/util/src/main/java/io/grpc/util/TransmitStatusRuntimeExceptionInterceptor.java +++ b/util/src/main/java/io/grpc/util/TransmitStatusRuntimeExceptionInterceptor.java @@ -104,6 +104,15 @@ public void onReady() { } } + @Override + public void onEvent(Object event) { + try { + super.onEvent(event); + } catch (StatusRuntimeException e) { + closeWithException(e); + } + } + private void closeWithException(StatusRuntimeException t) { Metadata metadata = t.getTrailers(); if (metadata == null) { @@ -276,5 +285,15 @@ public void run() { throw new RuntimeException(ERROR_MSG, e); } } + + @Override + public void triggerEvent(final Object event) { + serializingExecutor.execute(new Runnable() { + @Override + public void run() { + SerializingServerCall.super.triggerEvent(event); + } + }); + } } } diff --git a/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java b/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java index a4691d8bdec..4cc229c1851 100644 --- a/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java +++ b/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java @@ -104,6 +104,11 @@ public void onComplete() { public void onReady() { throw exception; } + + @Override + public void onEvent(Object event) { + throw exception; + } }; ServerServiceDefinition intercepted = ServerInterceptors.intercept( @@ -116,7 +121,36 @@ public void onReady() { getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onComplete(); getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onHalfClose(); getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onReady(); - assertEquals(5, call.numCloses); + getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers) + .onEvent(new Object()); + assertEquals(6, call.numCloses); + } + + @Test + public void statusRuntimeExceptionTransmitter_serializingServerCall_triggerEvent() { + final java.util.concurrent.atomic.AtomicReference eventRef = + new java.util.concurrent.atomic.AtomicReference<>(); + FakeServerCall call = new FakeServerCall(Status.OK, new Metadata()) { + @Override + public void triggerEvent(Object event) { + eventRef.set(event); + } + }; + final java.util.concurrent.atomic.AtomicReference> interceptedCall = + new java.util.concurrent.atomic.AtomicReference<>(); + listener = new VoidCallListener() { + @Override + public void onCall(ServerCall call, Metadata headers) { + interceptedCall.set(call); + } + }; + ServerServiceDefinition intercepted = ServerInterceptors.intercept( + serviceDefinition, + Arrays.asList(TransmitStatusRuntimeExceptionInterceptor.instance())); + getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers); + Object testEvent = new Object(); + interceptedCall.get().triggerEvent(testEvent); + assertEquals(testEvent, eventRef.get()); } @Test From 1b8a2308a63e0efa27f50c9be89f5d37ccc0c9ae Mon Sep 17 00:00:00 2001 From: Kannan J Date: Wed, 2 Sep 2026 08:01:44 +0000 Subject: [PATCH 14/17] core: add PerfMark tracing and tags to ServerCallImpl and ServerStreamListenerImpl triggerEvent Wrap ServerCallImpl.triggerEvent and ServerStreamListenerImpl.triggerEvent in PerfMark.traceTask with PerfMark.attachTag, aligning them with sendMessage, sendHeaders, close, request, and listener callbacks. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- .../main/java/io/grpc/internal/ServerCallImpl.java | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/core/src/main/java/io/grpc/internal/ServerCallImpl.java b/core/src/main/java/io/grpc/internal/ServerCallImpl.java index 6e894371be1..1f321d0d8dc 100644 --- a/core/src/main/java/io/grpc/internal/ServerCallImpl.java +++ b/core/src/main/java/io/grpc/internal/ServerCallImpl.java @@ -256,7 +256,10 @@ public MethodDescriptor getMethodDescriptor() { @Override public void triggerEvent(Object event) { - stream.triggerEvent(event); + try (TaskCloseable ignore = PerfMark.traceTask("ServerCall.triggerEvent")) { + PerfMark.attachTag(tag); + stream.triggerEvent(event); + } } @Override @@ -403,10 +406,13 @@ public void onReady() { @Override public void triggerEvent(Object event) { - if (call.cancelled) { - return; + try (TaskCloseable ignore = PerfMark.traceTask("ServerStreamListener.triggerEvent")) { + PerfMark.attachTag(call.tag); + if (call.cancelled) { + return; + } + listener.onEvent(event); } - listener.onEvent(event); } } } From bdfdafe5e6ca3d48c7693344665badfdff463caf Mon Sep 17 00:00:00 2001 From: Kannan J Date: Thu, 3 Sep 2026 04:47:00 +0000 Subject: [PATCH 15/17] core: fast-path return in ServerCallImpl.triggerEvent when closeCalled is true Check closeCalled before dispatching triggerEvent to the transport stream, avoiding unnecessary task allocations and transport hops if the call has already been closed. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- core/src/main/java/io/grpc/internal/ServerCallImpl.java | 3 +++ .../test/java/io/grpc/internal/ServerCallImplTest.java | 8 ++++++++ 2 files changed, 11 insertions(+) diff --git a/core/src/main/java/io/grpc/internal/ServerCallImpl.java b/core/src/main/java/io/grpc/internal/ServerCallImpl.java index 1f321d0d8dc..46b60a82b8c 100644 --- a/core/src/main/java/io/grpc/internal/ServerCallImpl.java +++ b/core/src/main/java/io/grpc/internal/ServerCallImpl.java @@ -256,6 +256,9 @@ public MethodDescriptor getMethodDescriptor() { @Override public void triggerEvent(Object event) { + if (closeCalled) { + return; + } try (TaskCloseable ignore = PerfMark.traceTask("ServerCall.triggerEvent")) { PerfMark.attachTag(tag); stream.triggerEvent(event); diff --git a/core/src/test/java/io/grpc/internal/ServerCallImplTest.java b/core/src/test/java/io/grpc/internal/ServerCallImplTest.java index 4a2de9f3936..d28c353145c 100644 --- a/core/src/test/java/io/grpc/internal/ServerCallImplTest.java +++ b/core/src/test/java/io/grpc/internal/ServerCallImplTest.java @@ -500,6 +500,14 @@ public void triggerEvent_propagatesToStream() { verify(stream).triggerEvent(event); } + @Test + public void triggerEvent_afterClose_noop() { + call.close(Status.OK, new Metadata()); + Object event = new Object(); + call.triggerEvent(event); + verify(stream, never()).triggerEvent(event); + } + @Test public void streamListener_triggerEvent() { ServerStreamListenerImpl streamListener = From b63b4a1f7d0fe0419bb0720b4762b2a32b1e95d1 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Thu, 3 Sep 2026 05:38:40 +0000 Subject: [PATCH 16/17] util: add unit tests for TransmitStatusRuntimeExceptionInterceptor custom event changes - Test onEvent throwing StatusRuntimeException closes the call with status and trailers. - Test onEvent throwing StatusRuntimeException on an already closed call does not trigger duplicate close. - Test SerializingServerCall executes triggerEvent sequentially in FIFO order on serializingExecutor. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- .../grpc/util/UtilServerInterceptorsTest.java | 107 ++++++++++++++++++ 1 file changed, 107 insertions(+) diff --git a/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java b/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java index 4cc229c1851..0054aa46031 100644 --- a/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java +++ b/util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java @@ -31,7 +31,10 @@ import io.grpc.Status; import io.grpc.StatusRuntimeException; import io.grpc.testing.TestMethodDescriptors; +import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; +import java.util.List; import org.junit.Test; import org.junit.runner.RunWith; import org.junit.runners.JUnit4; @@ -183,6 +186,7 @@ public void onHalfClose() { getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers); callDoubleSreListener.onMessage(null); // the only close with our exception callDoubleSreListener.onHalfClose(); // should not trigger a close + callDoubleSreListener.onEvent(new Object()); // should not trigger a close // this listener closes the call when it is initialized with startCall listener = new VoidCallListener() { @@ -195,13 +199,116 @@ public void onCall(ServerCall call, Metadata headers) { public void onHalfClose() { throw exception; } + + @Override + public void onEvent(Object event) { + throw exception; + } }; ServerCall.Listener callClosedListener = getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers); // call is already closed, does not match exception callClosedListener.onHalfClose(); // should not trigger a close + callClosedListener.onEvent(new Object()); // should not trigger a close + assertEquals(1, call.numCloses); + } + + @Test + public void statusRuntimeExceptionTransmitter_onEvent_transmitsStatusAndTrailers() { + final Status expectedStatus = Status.RESOURCE_EXHAUSTED.withDescription("rate limited"); + final Metadata expectedMetadata = new Metadata(); + Metadata.Key key = + Metadata.Key.of("custom-trailer", Metadata.ASCII_STRING_MARSHALLER); + expectedMetadata.put(key, "val"); + + final java.util.concurrent.atomic.AtomicReference closedStatus = + new java.util.concurrent.atomic.AtomicReference<>(); + final java.util.concurrent.atomic.AtomicReference closedTrailers = + new java.util.concurrent.atomic.AtomicReference<>(); + + FakeServerCall call = + new FakeServerCall(expectedStatus, expectedMetadata) { + @Override + public void close(Status status, Metadata trailers) { + closedStatus.set(status); + closedTrailers.set(trailers); + super.close(status, trailers); + } + }; + + final StatusRuntimeException exception = + new StatusRuntimeException(expectedStatus, expectedMetadata); + + listener = new VoidCallListener() { + @Override + public void onEvent(Object event) { + throw exception; + } + }; + + ServerServiceDefinition intercepted = ServerInterceptors.intercept( + serviceDefinition, + Arrays.asList(TransmitStatusRuntimeExceptionInterceptor.instance())); + + // When onEvent throws StatusRuntimeException, it should close the call with status and trailers + getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onEvent("event"); + assertEquals(1, call.numCloses); + assertEquals(expectedStatus, closedStatus.get()); + assertEquals("val", closedTrailers.get().get(key)); + } + + @Test + public void statusRuntimeExceptionTransmitter_serializingServerCall_serializesTriggerEvent() { + final List executionOrder = Collections.synchronizedList(new ArrayList()); + FakeServerCall call = new FakeServerCall(Status.OK, new Metadata()) { + @Override + public void sendHeaders(Metadata headers) { + executionOrder.add("sendHeaders"); + } + + @Override + public void triggerEvent(Object event) { + executionOrder.add("triggerEvent:" + event); + } + + @Override + public void sendMessage(Void message) { + executionOrder.add("sendMessage"); + } + + @Override + public void close(Status status, Metadata trailers) { + executionOrder.add("close"); + } + }; + + final java.util.concurrent.atomic.AtomicReference> interceptedCall = + new java.util.concurrent.atomic.AtomicReference<>(); + listener = new VoidCallListener() { + @Override + public void onCall(ServerCall call, Metadata headers) { + interceptedCall.set(call); + } + }; + + ServerServiceDefinition intercepted = ServerInterceptors.intercept( + serviceDefinition, + Arrays.asList(TransmitStatusRuntimeExceptionInterceptor.instance())); + getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers); + + ServerCall sc = interceptedCall.get(); + sc.sendHeaders(new Metadata()); + sc.triggerEvent("event1"); + sc.sendMessage(null); + sc.triggerEvent("event2"); + sc.close(Status.OK, new Metadata()); + + assertEquals( + Arrays.asList( + "sendHeaders", "triggerEvent:event1", "sendMessage", "triggerEvent:event2", "close"), + executionOrder); } private static class FakeServerCall extends NoopServerCall { From 91ec839b470a283427fd0f836e898e50904ef9c0 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Thu, 3 Sep 2026 07:14:19 +0000 Subject: [PATCH 17/17] testing: Fix race condition in MockServerTransportListener and improve stream cleanup in AbstractTransportTest In MockServerTransportListener.streamCreated(), stream.setListener(listener) was called after streams.add(StreamCreation(...)). This created a race condition where a test thread calling takeStreamOrFail() could dequeue the stream and call serverStream.triggerEvent() before stream.setListener() was called by the transport/container thread. When this occurred (e.g. in TomcatTransportTest on multi-core runners), ServletServerStream invoked transportState.triggerEvent() on the test thread, saw a null listener, threw a NullPointerException (swallowed by SerializingExecutor), and never enqueued the event into the listener queue, leading to a timeout and assertion failure: expected:<...Object@...> but was: Setting stream.setListener(listener) before enqueuing to streams guarantees that any thread consuming the StreamCreation will always observe a fully initialized listener. Additionally, in AbstractTransportTest.serverStream_triggerEvent(), replace clientStream.cancel(Status.CANCELLED) with serverStream.close(Status.OK, ...) for clean stream closure instead of leaving an uncoordinated client RST_STREAM in flight during container tearDown. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- .../java/io/grpc/internal/AbstractTransportTest.java | 2 +- .../java/io/grpc/internal/MockServerTransportListener.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java index 3898eb8be30..ec4e71e07f7 100644 --- a/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java +++ b/core/src/testFixtures/java/io/grpc/internal/AbstractTransportTest.java @@ -2114,7 +2114,7 @@ public void serverStream_triggerEvent() throws Exception { assertEquals(event, receivedEvent); // Cleanup - clientStream.cancel(Status.CANCELLED); + serverStream.close(Status.OK, new Metadata()); } @Test diff --git a/core/src/testFixtures/java/io/grpc/internal/MockServerTransportListener.java b/core/src/testFixtures/java/io/grpc/internal/MockServerTransportListener.java index e6c4e2f578e..be9436dd8d0 100644 --- a/core/src/testFixtures/java/io/grpc/internal/MockServerTransportListener.java +++ b/core/src/testFixtures/java/io/grpc/internal/MockServerTransportListener.java @@ -45,8 +45,8 @@ public MockServerTransportListener(ServerTransport transport) { @Override public void streamCreated(ServerStream stream, String method, Metadata headers) { ServerStreamListenerBase listener = new ServerStreamListenerBase(); - streams.add(new StreamCreation(stream, method, headers, listener)); stream.setListener(listener); + streams.add(new StreamCreation(stream, method, headers, listener)); } @Override