From 45b5af1d349fcf07bd50b2743690a425b15cfbc4 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Mon, 3 Aug 2026 16:35:25 +0000 Subject: [PATCH 01/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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/20] 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 From 47b21316761be99aa2393d0a8cf7300cf37906dc Mon Sep 17 00:00:00 2001 From: Kannan J Date: Tue, 8 Sep 2026 09:14:11 +0000 Subject: [PATCH 18/20] core: remove closeCalled check in ServerCallImpl.triggerEvent Consistent with ServerCallImpl.request(int), do not short-circuit on the non-volatile closeCalled boolean in triggerEvent. This ensures ServerCallImpl delegates triggerEvent to the underlying ServerStream, where stream lifecycle state and serialization are authoritatively managed in the transport layer. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- core/src/main/java/io/grpc/internal/ServerCallImpl.java | 3 --- core/src/test/java/io/grpc/internal/ServerCallImplTest.java | 5 +++-- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/core/src/main/java/io/grpc/internal/ServerCallImpl.java b/core/src/main/java/io/grpc/internal/ServerCallImpl.java index 46b60a82b8c..1f321d0d8dc 100644 --- a/core/src/main/java/io/grpc/internal/ServerCallImpl.java +++ b/core/src/main/java/io/grpc/internal/ServerCallImpl.java @@ -256,9 +256,6 @@ 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 d28c353145c..b29b6a1f6b7 100644 --- a/core/src/test/java/io/grpc/internal/ServerCallImplTest.java +++ b/core/src/test/java/io/grpc/internal/ServerCallImplTest.java @@ -501,13 +501,14 @@ public void triggerEvent_propagatesToStream() { } @Test - public void triggerEvent_afterClose_noop() { + public void triggerEvent_afterClose_propagatesToStream() { call.close(Status.OK, new Metadata()); Object event = new Object(); call.triggerEvent(event); - verify(stream, never()).triggerEvent(event); + verify(stream).triggerEvent(event); } + @Test public void streamListener_triggerEvent() { ServerStreamListenerImpl streamListener = From ab458cca71df42ee920d399e267e04f060ce6d0f Mon Sep 17 00:00:00 2001 From: Kannan J Date: Tue, 8 Sep 2026 12:10:51 +0000 Subject: [PATCH 19/20] api: document deadlock avoidance in ServerCall.triggerEvent and Listener.onEvent Add deadlock avoidance notes to the Javadoc of ServerCall.triggerEvent() and ServerCall.Listener.onEvent(). Transports such as Binder may hold transport locks (e.g. Inbound.this) while synchronously dispatching triggerEvent and onEvent. If application or interceptor code holds internal locks when calling ServerCall methods or acquires them inside onEvent(), lock order inversion deadlocks can occur. TAG=agy CONV=b39523be-9572-456e-a428-8b0e255067da --- api/src/main/java/io/grpc/ServerCall.java | 23 +++++++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/api/src/main/java/io/grpc/ServerCall.java b/api/src/main/java/io/grpc/ServerCall.java index 2e6ee07a23f..55579032edc 100644 --- a/api/src/main/java/io/grpc/ServerCall.java +++ b/api/src/main/java/io/grpc/ServerCall.java @@ -108,6 +108,14 @@ public void onReady() {} * callbacks (like {@link #onMessage}, {@link #onHalfClose}). This means the implementation * does not need internal synchronization to access call-specific state. * + *

Deadlock avoidance: In some transports (such as Binder transport) + * or when using a direct executor, this callback may be invoked while transport-level + * locks are held. Implementations should avoid acquiring locks that are held by callers of + * {@link ServerCall} methods (such as {@link ServerCall#triggerEvent}, + * {@link ServerCall#close}, or {@link ServerCall#request}), and should avoid calling + * {@link ServerCall} methods while holding application-level locks, as this can lead to + * deadlocks from lock-order inversion. + * * @param event the triggered event. */ @ExperimentalApi("https://github.com/grpc/grpc-java/issues/12979") @@ -280,8 +288,19 @@ 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. No events will be delivered - * after the RPC is cancelled or completed. + *

This method is safe to call from multiple threads without external synchronization. No + * events will be delivered after the RPC is cancelled or completed. + * + *

Deadlock avoidance: Callers should avoid holding application-level or + * interceptor locks when calling this method. Depending on the transport and executor + * configuration (such as {@code directExecutor()} or transports like Binder), + * {@code triggerEvent} may acquire transport-level locks and may dispatch + * {@link Listener#onEvent(Object)} synchronously on the calling thread. + * If the caller holds an application lock while calling {@code triggerEvent}, and + * {@code onEvent} or a concurrent transport operation (such as {@link #close} or + * {@link #request}) attempts to acquire that same lock, a deadlock can occur from + * lock-order inversion. Applications should mutate internal state under lock, release + * the lock, and only then invoke {@code triggerEvent}. * * @param event the event to trigger. */ From 573aff80228b8f85265cc1cc76b47dbe75ab45e5 Mon Sep 17 00:00:00 2001 From: Kannan J Date: Thu, 10 Sep 2026 06:11:25 +0000 Subject: [PATCH 20/20] xds: ensure stream completion before teardown in context propagation tests In ExternalProcessorClientInterceptorTest, tests in Category 27 (clientInterceptor_contextPropagated*) configure dedicated single-thread executors for thread-boundary context propagation testing and register the in-process channels and servers with GrpcCleanupRule. Previously, the tests exited their try-blocks after observing only intermediate latches (such as downstream start or ext-proc call start), then immediately triggered call cancellation and entered finally-block executor shutdown. When custom executors were shut down while stream cancellation or completion tasks were still in flight or being queued, in-process transport streams could not cleanly complete. Consequently, GrpcCleanupRule.after() timed out awaiting channel/server termination, failing with: java.lang.AssertionError: Resources could not be released in time Fix this by awaiting completion/closure latches on both the client call listener (onClose) and the mock external processor server (onError / onCompleted) before exiting the try block and shutting down the executors. TAG=agy CONV=e1bfa5a2-e855-4f79-abdd-ef2b264977be --- ...xternalProcessorClientInterceptorTest.java | 44 ++++++++++++++++--- 1 file changed, 38 insertions(+), 6 deletions(-) diff --git a/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java b/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java index cd6de138a48..992e85d5151 100644 --- a/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java +++ b/xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java @@ -14301,6 +14301,7 @@ public void onClose(Status status, Metadata trailers) { @Test public void clientInterceptor_contextPropagatedToStartCall() throws Exception { String uniqueExtProcServerName = InProcessServerBuilder.generateName(); + final CountDownLatch extProcCompletedLatch = new CountDownLatch(1); ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() { @Override @@ -14321,11 +14322,14 @@ public void onNext(ProcessingRequest request) { } @Override - public void onError(Throwable t) {} + public void onError(Throwable t) { + extProcCompletedLatch.countDown(); + } @Override public void onCompleted() { responseObserver.onCompleted(); + extProcCompletedLatch.countDown(); } }; } @@ -14390,6 +14394,13 @@ public void start(ClientCall.Listener responseListener, Metadata headers) final AtomicReference> proxyCallRef = new AtomicReference<>(); ExecutorService callExecutor = Executors.newSingleThreadExecutor(); + final CountDownLatch callClosedLatch = new CountDownLatch(1); + ClientCall.Listener callListener = new ClientCall.Listener() { + @Override + public void onClose(Status status, Metadata trailers) { + callClosedLatch.countDown(); + } + }; try { testContext.run(() -> { ClientCall proxyCall = interceptCall( @@ -14398,7 +14409,7 @@ public void start(ClientCall.Listener responseListener, Metadata headers) DEFAULT_CALL_OPTIONS.withExecutor(callExecutor), dataPlaneChannel); proxyCallRef.set(proxyCall); - proxyCall.start(new ClientCall.Listener() {}, new Metadata()); + proxyCall.start(callListener, new Metadata()); }); ClientCall proxyCall = proxyCallRef.get(); @@ -14409,6 +14420,8 @@ public void start(ClientCall.Listener responseListener, Metadata headers) assertThat(downstreamStartLatch.await(5, TimeUnit.SECONDS)).isTrue(); assertThat(contextValueAtDownstreamStart.get()).isEqualTo("test-value"); + assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(extProcCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue(); proxyCall.cancel("cleanup", null); } finally { @@ -14423,6 +14436,7 @@ public void start(ClientCall.Listener responseListener, Metadata headers) @Test public void clientInterceptor_contextPropagatedToListenerCallbacks() throws Exception { String uniqueExtProcServerName = InProcessServerBuilder.generateName(); + final CountDownLatch extProcCompletedLatch = new CountDownLatch(1); ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() { @Override @@ -14443,11 +14457,14 @@ public void onNext(ProcessingRequest request) { } @Override - public void onError(Throwable t) {} + public void onError(Throwable t) { + extProcCompletedLatch.countDown(); + } @Override public void onCompleted() { responseObserver.onCompleted(); + extProcCompletedLatch.countDown(); } }; } @@ -14542,6 +14559,7 @@ public void onReady() { proxyCall.halfClose(); assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(extProcCompletedLatch.await(5, TimeUnit.SECONDS)).isTrue(); assertThat(onHeadersContext.get()).isEqualTo("test-value"); assertThat(onMessageContext.get()).isEqualTo("test-value"); @@ -14561,6 +14579,7 @@ public void onReady() { @Test public void clientInterceptor_contextPropagatedToExtProcStub() throws Exception { String uniqueExtProcServerName = InProcessServerBuilder.generateName(); + final CountDownLatch extProcClosedLatch = new CountDownLatch(1); ExternalProcessorGrpc.ExternalProcessorImplBase extProcImpl = new ExternalProcessorGrpc.ExternalProcessorImplBase() { @Override @@ -14571,10 +14590,14 @@ public StreamObserver process( public void onNext(ProcessingRequest request) {} @Override - public void onError(Throwable t) {} + public void onError(Throwable t) { + extProcClosedLatch.countDown(); + } @Override - public void onCompleted() {} + public void onCompleted() { + extProcClosedLatch.countDown(); + } }; } }; @@ -14623,6 +14646,13 @@ public ClientCall interceptCall( final AtomicReference> proxyCallRef = new AtomicReference<>(); ExecutorService callExecutor = Executors.newSingleThreadExecutor(); + final CountDownLatch callClosedLatch = new CountDownLatch(1); + ClientCall.Listener callListener = new ClientCall.Listener() { + @Override + public void onClose(Status status, Metadata trailers) { + callClosedLatch.countDown(); + } + }; try { testContext.run(() -> { ClientCall proxyCall = interceptCall( @@ -14631,7 +14661,7 @@ public ClientCall interceptCall( DEFAULT_CALL_OPTIONS.withExecutor(callExecutor), dataPlaneChannel); proxyCallRef.set(proxyCall); - proxyCall.start(new ClientCall.Listener() {}, new Metadata()); + proxyCall.start(callListener, new Metadata()); }); ClientCall proxyCall = proxyCallRef.get(); @@ -14640,6 +14670,8 @@ public ClientCall interceptCall( assertThat(contextAtExtProcCall.get()).isEqualTo("test-value"); proxyCall.cancel("cleanup", null); + assertThat(callClosedLatch.await(5, TimeUnit.SECONDS)).isTrue(); + assertThat(extProcClosedLatch.await(5, TimeUnit.SECONDS)).isTrue(); } finally { channelManager.close(); shutdownAndAwaitTermination(extProcServerExecutor);