Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions api/src/main/java/io/grpc/Contexts.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}

/**
Expand Down
5 changes: 5 additions & 0 deletions api/src/main/java/io/grpc/PartialForwardingServerCall.java
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
28 changes: 28 additions & 0 deletions api/src/main/java/io/grpc/ServerCall.java
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,20 @@ public void onComplete() {}
* <em>another</em> {@code onReady()} callback.
*/
public void onReady() {}

/**
* A custom event has been triggered by the call.
*
* <p>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.
*/
@ExperimentalApi("https://github.com/grpc/grpc-java/issues/12979")
public void onEvent(Object event) {
// Default no-op
}
}

/**
Expand Down Expand Up @@ -262,6 +276,20 @@ 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.
*
* <p>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.
*/
@ExperimentalApi("https://github.com/grpc/grpc-java/issues/12979")
public void triggerEvent(Object event) {
// Default no-op
}

/**
* The {@link MethodDescriptor} for the call.
*/
Expand Down
17 changes: 16 additions & 1 deletion api/src/test/java/io/grpc/ContextsTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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<Object> wrapped = interceptCall(uniqueContext, call, headers,
new ServerCallHandler<Object, Object>() {
Expand All @@ -101,7 +106,8 @@ public ServerCall.Listener<Object> 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());
}

Expand Down Expand Up @@ -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<Object> wrapped = interceptCall(uniqueContext, call, headers,
new ServerCallHandler<Object, Object>() {
Expand Down Expand Up @@ -180,6 +190,11 @@ public ServerCall.Listener<Object> startCall(
fail("Exception expected");
} catch (RuntimeException expected) {
}
try {
wrapped.onEvent(new Object());
fail("Exception expected");
} catch (RuntimeException expected) {
}
assertSame(origContext, Context.current());
}

Expand Down
10 changes: 10 additions & 0 deletions binder/src/main/java/io/grpc/binder/internal/Inbound.java
Original file line number Diff line number Diff line change
Expand Up @@ -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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,11 @@ public void setDecompressor(Decompressor decompressor) {
// Ignore.
}

@Override
public void triggerEvent(Object event) {
inbound.triggerEvent(event);
}

@Override
public void optimizeForDirectExecutor() {
// Ignore.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,11 @@ public void setDecompressor(Decompressor decompressor) {
// Ignore.
}

@Override
public void triggerEvent(Object event) {
inbound.triggerEvent(event);
}

@Override
public void optimizeForDirectExecutor() {
// Ignore.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -289,6 +289,7 @@ public ListenableFuture<Status> checkAuthorizationAsync(int uid) {
ListenableFuture<Status> authFuture = asyncPolicy.checkAuthorizationAsync(SOME_UID);
assertThat(awaitResult(settableUid)).isEqualTo(SOME_UID);
authFuture.cancel(false);
executor.submit(() -> {}).get(10, TimeUnit.SECONDS);

assertThat(delegateAuthFuture.isCancelled()).isTrue();
}
Expand Down
17 changes: 17 additions & 0 deletions core/src/main/java/io/grpc/internal/AbstractServerStream.java
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,16 @@ 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() {
transportState().triggerEvent(event);
}
});
}

@Override
public StatsTraceContext statsTraceContext() {
return statsTraceCtx;
Expand Down Expand Up @@ -259,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;
Expand Down
13 changes: 13 additions & 0 deletions core/src/main/java/io/grpc/internal/ServerCallImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,11 @@ public MethodDescriptor<ReqT, RespT> getMethodDescriptor() {
return method;
}

@Override
public void triggerEvent(Object event) {
stream.triggerEvent(event);
}

@Override
public SecurityLevel getSecurityLevel() {
final Attributes attributes = getAttributes();
Expand Down Expand Up @@ -395,5 +400,13 @@ public void onReady() {
listener.onReady();
}
}

@Override
public void triggerEvent(Object event) {
if (call.cancelled) {
return;
}
listener.onEvent(event);
}
}
}
31 changes: 31 additions & 0 deletions core/src/main/java/io/grpc/internal/ServerImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -781,6 +781,9 @@ public void closed(Status status) {}

@Override
public void onReady() {}

@Override
public void triggerEvent(Object event) {}
}

/**
Expand Down Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions core/src/main/java/io/grpc/internal/ServerStream.java
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
5 changes: 5 additions & 0 deletions core/src/main/java/io/grpc/internal/ServerStreamListener.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
28 changes: 28 additions & 0 deletions core/src/test/java/io/grpc/internal/AbstractServerStreamTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down Expand Up @@ -391,6 +416,9 @@ public void halfClosed() {}

@Override
public void closed(Status status) {}

@Override
public void triggerEvent(Object event) {}
}

private static class AbstractServerStreamBase extends AbstractServerStream {
Expand Down
26 changes: 26 additions & 0 deletions core/src/test/java/io/grpc/internal/ServerCallImplTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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<Long> 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<Long> 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<Long> {
@Override
public InputStream stream(Long value) {
Expand Down
Loading
Loading