Skip to content

Commit 48aee63

Browse files
committed
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
1 parent e832422 commit 48aee63

6 files changed

Lines changed: 81 additions & 1 deletion

File tree

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,12 @@ public void onReady() {
8888
maybeRunPendingSteps();
8989
}
9090

91+
@Override
92+
public void onEvent(Object event) {
93+
pendingSteps.offer(delegate -> delegate.onEvent(event));
94+
maybeRunPendingSteps();
95+
}
96+
9197
/**
9298
* Similar to Java8's {@link java.util.function.Consumer}, but redeclared in order to support
9399
* Android SDK 21.

binder/src/test/java/io/grpc/binder/internal/PendingAuthListenerTest.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ public void setUp() {
4545
public void onCallbacks_noOpBeforeStartCall() {
4646
listener.onReady();
4747
listener.onMessage("foo");
48+
listener.onEvent("bar");
4849
listener.onHalfClose();
4950
listener.onComplete();
5051

@@ -54,16 +55,19 @@ public void onCallbacks_noOpBeforeStartCall() {
5455
@Test
5556
public void onCallbacks_runsPendingCallbacksAfterStartCall() {
5657
String message = "foo";
58+
String event = "bar";
5759

5860
// Act 1
5961
listener.onReady();
6062
listener.onMessage(message);
63+
listener.onEvent(event);
6164
listener.startCall(call, headers, next);
6265

6366
// Assert 1
6467
InOrder order = Mockito.inOrder(delegate);
6568
order.verify(delegate).onReady();
6669
order.verify(delegate).onMessage(message);
70+
order.verify(delegate).onEvent(event);
6771

6872
// Act 2
6973
listener.onHalfClose();

opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -452,6 +452,13 @@ public void onReady() {
452452
delegate().onReady();
453453
}
454454
}
455+
456+
@Override
457+
public void onEvent(Object event) {
458+
try (Scope scope = context.makeCurrent()) {
459+
delegate().onEvent(event);
460+
}
461+
}
455462
}
456463

457464
@VisibleForTesting

opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -948,6 +948,11 @@ public void onCancel() {
948948
public void onComplete() {
949949
callbackSpan.set(Span.fromContext(Context.current()));
950950
}
951+
952+
@Override
953+
public void onEvent(Object event) {
954+
callbackSpan.set(Span.fromContext(Context.current()));
955+
}
951956
};
952957
ServerInterceptor interceptor = tracingModule.getServerSpanPropagationInterceptor();
953958
@SuppressWarnings("unchecked")
@@ -967,6 +972,8 @@ public void onComplete() {
967972
assertEquals(callbackSpan.get(), Span.getInvalid());
968973
listener.onComplete();
969974
assertEquals(callbackSpan.get(), Span.getInvalid());
975+
listener.onEvent(new Object());
976+
assertEquals(callbackSpan.get(), Span.getInvalid());
970977

971978
Span parentSpan = tracerRule.spanBuilder("parent-span").startSpan();
972979
io.grpc.Context context = io.grpc.Context.current().withValue(
@@ -990,6 +997,9 @@ public void onComplete() {
990997
listener.onComplete();
991998
assertEquals(callbackSpan.get().getSpanContext().getTraceId(),
992999
parentSpan.getSpanContext().getTraceId());
1000+
listener.onEvent(new Object());
1001+
assertEquals(callbackSpan.get().getSpanContext().getTraceId(),
1002+
parentSpan.getSpanContext().getTraceId());
9931003
} finally {
9941004
context.detach(previous);
9951005
}

util/src/main/java/io/grpc/util/TransmitStatusRuntimeExceptionInterceptor.java

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,15 @@ public void onReady() {
104104
}
105105
}
106106

107+
@Override
108+
public void onEvent(Object event) {
109+
try {
110+
super.onEvent(event);
111+
} catch (StatusRuntimeException e) {
112+
closeWithException(e);
113+
}
114+
}
115+
107116
private void closeWithException(StatusRuntimeException t) {
108117
Metadata metadata = t.getTrailers();
109118
if (metadata == null) {
@@ -276,5 +285,15 @@ public void run() {
276285
throw new RuntimeException(ERROR_MSG, e);
277286
}
278287
}
288+
289+
@Override
290+
public void triggerEvent(final Object event) {
291+
serializingExecutor.execute(new Runnable() {
292+
@Override
293+
public void run() {
294+
SerializingServerCall.super.triggerEvent(event);
295+
}
296+
});
297+
}
279298
}
280299
}

util/src/test/java/io/grpc/util/UtilServerInterceptorsTest.java

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,11 @@ public void onComplete() {
104104
public void onReady() {
105105
throw exception;
106106
}
107+
108+
@Override
109+
public void onEvent(Object event) {
110+
throw exception;
111+
}
107112
};
108113

109114
ServerServiceDefinition intercepted = ServerInterceptors.intercept(
@@ -116,7 +121,36 @@ public void onReady() {
116121
getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onComplete();
117122
getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onHalfClose();
118123
getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers).onReady();
119-
assertEquals(5, call.numCloses);
124+
getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers)
125+
.onEvent(new Object());
126+
assertEquals(6, call.numCloses);
127+
}
128+
129+
@Test
130+
public void statusRuntimeExceptionTransmitter_serializingServerCall_triggerEvent() {
131+
final java.util.concurrent.atomic.AtomicReference<Object> eventRef =
132+
new java.util.concurrent.atomic.AtomicReference<>();
133+
FakeServerCall<Void, Void> call = new FakeServerCall<Void, Void>(Status.OK, new Metadata()) {
134+
@Override
135+
public void triggerEvent(Object event) {
136+
eventRef.set(event);
137+
}
138+
};
139+
final java.util.concurrent.atomic.AtomicReference<ServerCall<Void, Void>> interceptedCall =
140+
new java.util.concurrent.atomic.AtomicReference<>();
141+
listener = new VoidCallListener() {
142+
@Override
143+
public void onCall(ServerCall<Void, Void> call, Metadata headers) {
144+
interceptedCall.set(call);
145+
}
146+
};
147+
ServerServiceDefinition intercepted = ServerInterceptors.intercept(
148+
serviceDefinition,
149+
Arrays.asList(TransmitStatusRuntimeExceptionInterceptor.instance()));
150+
getSoleMethod(intercepted).getServerCallHandler().startCall(call, headers);
151+
Object testEvent = new Object();
152+
interceptedCall.get().triggerEvent(testEvent);
153+
assertEquals(testEvent, eventRef.get());
120154
}
121155

122156
@Test

0 commit comments

Comments
 (0)