Repository navigation
core: Delayed deserialization for unary/server-streaming calls #13004
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
09d8c90
553b33a
50197f1
f237d2f
072e999
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -288,6 +288,7 @@ static final class ServerStreamListenerImpl<ReqT> implements ServerStreamListene | |
| private final ServerCallImpl<ReqT, ?> call; | ||
| private final ServerCall.Listener<ReqT> listener; | ||
| private final Context.CancellableContext context; | ||
| private InputStream delayedMessage; | ||
|
|
||
| public ServerStreamListenerImpl( | ||
| ServerCallImpl<ReqT, ?> call, ServerCall.Listener<ReqT> listener, | ||
|
|
@@ -330,13 +331,27 @@ private void messagesAvailableInternal(final MessageProducer producer) { | |
| InputStream message; | ||
| try { | ||
| while ((message = producer.next()) != null) { | ||
| try { | ||
| listener.onMessage(call.method.parseRequest(message)); | ||
| } catch (Throwable t) { | ||
| GrpcUtil.closeQuietly(message); | ||
| throw t; | ||
| // TODO: Consider forcing this check to be done in the transport (MessageDeframer) | ||
| // https://github.com/grpc/grpc-java/pull/13004/changes#r3939373996 | ||
| if (call.method.getType().clientSendsOneMessage()) { | ||
| if (delayedMessage != null) { | ||
| GrpcUtil.closeQuietly(message); | ||
| call.stream.cancel(Status.INTERNAL.withDescription("Too many requests")); | ||
| GrpcUtil.closeQuietly(delayedMessage); | ||
| delayedMessage = null; | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This puts the call back into a normal state, so if other events happen after this one (e.g., message, or half close), that could end up propagating to the application before the cancel is processed. I don't know the easiest way to handle that though; obviously we could set some more state/booleans. It is probably worth looking into the exception handling in the executor see what would happen if we throw here.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. If an exception is thrown from
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Now doesn't this code call the application's onCancel() twice? It calls it once here with the call to
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes. All we need to do instead of calling |
||
| call.cancelled = true; | ||
| return; | ||
| } | ||
| delayedMessage = message; | ||
| } else { | ||
| try { | ||
| listener.onMessage(call.method.parseRequest(message)); | ||
| } catch (Throwable t) { | ||
| GrpcUtil.closeQuietly(message); | ||
| throw t; | ||
| } | ||
| message.close(); | ||
| } | ||
| message.close(); | ||
| } | ||
| } catch (Throwable t) { | ||
| GrpcUtil.closeQuietly(producer); | ||
|
|
@@ -353,6 +368,19 @@ public void halfClosed() { | |
| return; | ||
| } | ||
|
|
||
| if (delayedMessage != null) { | ||
| InputStream message = delayedMessage; | ||
| delayedMessage = null; | ||
| try { | ||
| listener.onMessage(call.method.parseRequest(message)); | ||
| } catch (Throwable t) { | ||
| GrpcUtil.closeQuietly(message); | ||
| Throwables.throwIfUnchecked(t); | ||
| throw new RuntimeException(t); | ||
| } | ||
| GrpcUtil.closeQuietly(message); | ||
| } | ||
|
|
||
| listener.onHalfClose(); | ||
| } | ||
| } | ||
|
|
@@ -366,6 +394,10 @@ public void closed(Status status) { | |
| } | ||
|
|
||
| private void closedInternal(Status status) { | ||
| if (delayedMessage != null) { | ||
| GrpcUtil.closeQuietly(delayedMessage); | ||
| delayedMessage = null; | ||
| } | ||
| Throwable cancelCause = null; | ||
| try { | ||
| if (status.isOk()) { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
FYI, It'd be fine to not support mark (hard-code this method and mark()). Decompression streams wouldn't support it anyway.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Done.