lryan | 56e307f | 2014-12-05 13:25:08 -0800 | [diff] [blame] | 1 | /* |
| 2 | * Copyright 2014, Google Inc. All rights reserved. |
| 3 | * |
| 4 | * Redistribution and use in source and binary forms, with or without |
| 5 | * modification, are permitted provided that the following conditions are |
| 6 | * met: |
| 7 | * |
| 8 | * * Redistributions of source code must retain the above copyright |
| 9 | * notice, this list of conditions and the following disclaimer. |
| 10 | * * Redistributions in binary form must reproduce the above |
| 11 | * copyright notice, this list of conditions and the following disclaimer |
| 12 | * in the documentation and/or other materials provided with the |
| 13 | * distribution. |
| 14 | * |
| 15 | * * Neither the name of Google Inc. nor the names of its |
| 16 | * contributors may be used to endorse or promote products derived from |
| 17 | * this software without specific prior written permission. |
| 18 | * |
| 19 | * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS |
| 20 | * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT |
| 21 | * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR |
| 22 | * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT |
| 23 | * OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, |
| 24 | * SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT |
| 25 | * LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, |
| 26 | * DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY |
| 27 | * THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT |
| 28 | * (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE |
| 29 | * OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. |
| 30 | */ |
| 31 | |
nathanmittler | 29cbef1 | 2014-10-27 11:33:19 -0700 | [diff] [blame] | 32 | package com.google.net.stubby.transport; |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 33 | |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 34 | import com.google.common.base.MoreObjects; |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 35 | import com.google.common.base.Preconditions; |
lryan | e4bd1c7 | 2014-09-08 14:03:35 -0700 | [diff] [blame] | 36 | import com.google.net.stubby.Metadata; |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 37 | import com.google.net.stubby.Status; |
| 38 | |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 39 | import java.io.InputStream; |
| 40 | import java.nio.ByteBuffer; |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 41 | import java.util.logging.Level; |
| 42 | import java.util.logging.Logger; |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 43 | |
lryan | 28497e3 | 2014-10-17 16:14:38 -0700 | [diff] [blame] | 44 | import javax.annotation.Nullable; |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 45 | |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 46 | /** |
| 47 | * The abstract base class for {@link ClientStream} implementations. |
| 48 | */ |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 49 | public abstract class AbstractClientStream<IdT> extends AbstractStream<IdT> |
| 50 | implements ClientStream { |
| 51 | |
| 52 | private static final Logger log = Logger.getLogger(AbstractClientStream.class.getName()); |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 53 | |
ejona | 6fc356b | 2014-09-22 12:49:20 -0700 | [diff] [blame] | 54 | private final ClientStreamListener listener; |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 55 | private boolean listenerClosed; |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 56 | |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 57 | // Stored status & trailers to report when deframer completes or |
| 58 | // transportReportStatus is directly called. |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 59 | private Status status; |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 60 | private Metadata.Trailers trailers; |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 61 | private Runnable closeListenerTask; |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 62 | |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 63 | |
nmittler | de3a131 | 2015-01-16 11:54:24 -0800 | [diff] [blame^] | 64 | protected AbstractClientStream(ClientStreamListener listener) { |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 65 | this.listener = Preconditions.checkNotNull(listener); |
| 66 | } |
| 67 | |
| 68 | @Override |
nmittler | de3a131 | 2015-01-16 11:54:24 -0800 | [diff] [blame^] | 69 | protected void receiveMessage(InputStream is, int length) { |
| 70 | if (!listenerClosed) { |
| 71 | listener.messageRead(is, length); |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 72 | } |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 73 | } |
| 74 | |
lryan | 28497e3 | 2014-10-17 16:14:38 -0700 | [diff] [blame] | 75 | @Override |
| 76 | public final void writeMessage(InputStream message, int length, @Nullable Runnable accepted) { |
| 77 | super.writeMessage(message, length, accepted); |
| 78 | } |
| 79 | |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 80 | /** |
| 81 | * The transport implementation has detected a protocol error on the stream. Transports are |
| 82 | * responsible for properly closing streams when protocol errors occur. |
| 83 | * |
| 84 | * @param errorStatus the error to report |
| 85 | */ |
| 86 | protected void inboundTransportError(Status errorStatus) { |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 87 | if (inboundPhase() == Phase.STATUS) { |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 88 | log.log(Level.INFO, "Received transport error on closed stream {0} {1}", |
| 89 | new Object[]{id(), errorStatus}); |
| 90 | return; |
| 91 | } |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 92 | // For transport errors we immediately report status to the application layer |
| 93 | // and do not wait for additional payloads. |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 94 | transportReportStatus(errorStatus, false, new Metadata.Trailers()); |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 95 | } |
| 96 | |
| 97 | /** |
| 98 | * Called by transport implementations when they receive headers. When receiving headers |
| 99 | * a transport may determine that there is an error in the protocol at this phase which is |
| 100 | * why this method takes an error {@link Status}. If a transport reports an |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 101 | * {@link com.google.net.stubby.Status.Code#INTERNAL} error |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 102 | * |
| 103 | * @param headers the parsed headers |
| 104 | */ |
| 105 | protected void inboundHeadersReceived(Metadata.Headers headers) { |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 106 | if (inboundPhase() == Phase.STATUS) { |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 107 | log.log(Level.INFO, "Received headers on closed stream {0} {1}", |
| 108 | new Object[]{id(), headers}); |
| 109 | } |
| 110 | inboundPhase(Phase.MESSAGE); |
nmittler | de3a131 | 2015-01-16 11:54:24 -0800 | [diff] [blame^] | 111 | listener.headersRead(headers); |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 112 | } |
| 113 | |
| 114 | /** |
| 115 | * Process the contents of a received data frame from the server. |
| 116 | */ |
| 117 | protected void inboundDataReceived(Buffer frame) { |
| 118 | Preconditions.checkNotNull(frame, "frame"); |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 119 | if (inboundPhase() == Phase.STATUS) { |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 120 | frame.close(); |
| 121 | return; |
| 122 | } |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 123 | if (inboundPhase() == Phase.HEADERS) { |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 124 | // Have not received headers yet so error |
| 125 | inboundTransportError(Status.INTERNAL.withDescription("headers not received before payload")); |
| 126 | frame.close(); |
| 127 | return; |
| 128 | } |
| 129 | inboundPhase(Phase.MESSAGE); |
nathanmittler | 5d953e8 | 2014-11-18 13:12:42 -0800 | [diff] [blame] | 130 | |
| 131 | deframe(frame, false); |
| 132 | } |
| 133 | |
| 134 | @Override |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 135 | protected void inboundDeliveryPaused() { |
| 136 | runCloseListenerTask(); |
| 137 | } |
| 138 | |
| 139 | @Override |
nathanmittler | 5d953e8 | 2014-11-18 13:12:42 -0800 | [diff] [blame] | 140 | protected final void deframeFailed(Throwable cause) { |
| 141 | log.log(Level.WARNING, "Exception processing message", cause); |
| 142 | cancel(); |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 143 | } |
| 144 | |
| 145 | /** |
| 146 | * Called by transport implementations when they receive trailers. |
| 147 | */ |
| 148 | protected void inboundTrailersReceived(Metadata.Trailers trailers, Status status) { |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 149 | Preconditions.checkNotNull(trailers, "trailers"); |
| 150 | if (inboundPhase() == Phase.STATUS) { |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 151 | log.log(Level.INFO, "Received trailers on closed stream {0}\n {1}\n {3}", |
| 152 | new Object[]{id(), status, trailers}); |
| 153 | } |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 154 | // Stash the status & trailers so they can be delivered by the deframer calls |
| 155 | // remoteEndClosed |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 156 | this.status = status; |
ejona | c0f4192 | 2014-12-23 12:16:17 -0800 | [diff] [blame] | 157 | this.trailers = trailers; |
nathanmittler | 5d953e8 | 2014-11-18 13:12:42 -0800 | [diff] [blame] | 158 | deframe(Buffers.empty(), true); |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 159 | } |
| 160 | |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 161 | @Override |
| 162 | protected void remoteEndClosed() { |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 163 | transportReportStatus(status, true, trailers); |
ejona | 9d50299 | 2014-09-22 12:23:19 -0700 | [diff] [blame] | 164 | } |
| 165 | |
| 166 | @Override |
| 167 | protected final void internalSendFrame(ByteBuffer frame, boolean endOfStream) { |
| 168 | sendFrame(frame, endOfStream); |
| 169 | } |
| 170 | |
| 171 | /** |
| 172 | * Sends an outbound frame to the remote end point. |
| 173 | * |
| 174 | * @param frame a buffer containing the chunk of data to be sent. |
| 175 | * @param endOfStream if {@code true} indicates that no more data will be sent on the stream by |
| 176 | * this endpoint. |
| 177 | */ |
| 178 | protected abstract void sendFrame(ByteBuffer frame, boolean endOfStream); |
| 179 | |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 180 | /** |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 181 | * Report stream closure with status to the application layer if not already reported. This method |
| 182 | * must be called from the transport thread. |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 183 | * |
| 184 | * @param newStatus the new status to set |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 185 | * @param stopDelivery if {@code true}, interrupts any further delivery of inbound messages that |
| 186 | * may already be queued up in the deframer. If {@code false}, the listener will be |
| 187 | * notified immediately after all currently completed messages in the deframer have been |
| 188 | * delivered to the application. |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 189 | */ |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 190 | public void transportReportStatus(final Status newStatus, boolean stopDelivery, |
| 191 | final Metadata.Trailers trailers) { |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 192 | Preconditions.checkNotNull(newStatus, "newStatus"); |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 193 | |
| 194 | boolean closingLater = closeListenerTask != null && !stopDelivery; |
| 195 | if (listenerClosed || closingLater) { |
| 196 | // We already closed (or are about to close) the listener. |
| 197 | return; |
| 198 | } |
| 199 | |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 200 | inboundPhase(Phase.STATUS); |
| 201 | status = newStatus; |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 202 | closeListenerTask = null; |
| 203 | |
| 204 | // Determine if the deframer is stalled (i.e. currently has no complete messages to deliver). |
nmittler | de3a131 | 2015-01-16 11:54:24 -0800 | [diff] [blame^] | 205 | boolean deliveryStalled = deframer.isStalled(); |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 206 | |
| 207 | if (stopDelivery || deliveryStalled) { |
| 208 | // Close the listener immediately. |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 209 | listenerClosed = true; |
| 210 | listener.closed(newStatus, trailers); |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 211 | } else { |
| 212 | // Delay close until inboundDeliveryStalled() |
| 213 | closeListenerTask = newCloseListenerTask(newStatus, trailers); |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 214 | } |
nathanmittler | 11c363a | 2015-01-09 11:22:19 -0800 | [diff] [blame] | 215 | } |
| 216 | |
| 217 | /** |
| 218 | * Creates a new {@link Runnable} to close the listener with the given status/trailers. |
| 219 | */ |
| 220 | private Runnable newCloseListenerTask(final Status status, final Metadata.Trailers trailers) { |
| 221 | return new Runnable() { |
| 222 | @Override |
| 223 | public void run() { |
| 224 | if (!listenerClosed) { |
| 225 | // Status has not been reported to the application layer |
| 226 | listenerClosed = true; |
| 227 | listener.closed(status, trailers); |
| 228 | } |
| 229 | } |
| 230 | }; |
| 231 | } |
| 232 | |
| 233 | /** |
| 234 | * Executes the pending listener close task, if one exists. |
| 235 | */ |
| 236 | private void runCloseListenerTask() { |
| 237 | if (closeListenerTask != null) { |
| 238 | closeListenerTask.run(); |
| 239 | closeListenerTask = null; |
| 240 | } |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 241 | } |
| 242 | |
| 243 | @Override |
| 244 | public final void halfClose() { |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 245 | if (outboundPhase(Phase.STATUS) != Phase.STATUS) { |
ejona | c0f4192 | 2014-12-23 12:16:17 -0800 | [diff] [blame] | 246 | closeFramer(); |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 247 | } |
| 248 | } |
| 249 | |
| 250 | /** |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 251 | * Cancel the stream. Called by the application layer, never called by the transport. |
| 252 | */ |
| 253 | @Override |
| 254 | public void cancel() { |
| 255 | outboundPhase(Phase.STATUS); |
| 256 | if (id() != null) { |
| 257 | // Only send a cancellation to remote side if we have actually been allocated |
| 258 | // a stream id and we are not already closed. i.e. the server side is aware of the stream. |
| 259 | sendCancel(); |
| 260 | } |
| 261 | dispose(); |
| 262 | } |
| 263 | |
| 264 | /** |
| 265 | * Send a stream cancellation message to the remote server. Can be called by either the |
| 266 | * application or transport layers. |
lryan | 669724a | 2014-11-10 10:21:45 -0800 | [diff] [blame] | 267 | */ |
| 268 | protected abstract void sendCancel(); |
lryan | c5e70c2 | 2014-11-24 16:41:02 -0800 | [diff] [blame] | 269 | |
| 270 | @Override |
| 271 | protected MoreObjects.ToStringHelper toStringHelper() { |
| 272 | MoreObjects.ToStringHelper toStringHelper = super.toStringHelper(); |
| 273 | if (status != null) { |
| 274 | toStringHelper.add("status", status); |
| 275 | } |
| 276 | return toStringHelper; |
| 277 | } |
| 278 | |
| 279 | @Override |
| 280 | public boolean isClosed() { |
| 281 | return super.isClosed() || listenerClosed; |
| 282 | } |
zhangkun | 048649e | 2014-08-28 15:52:03 -0700 | [diff] [blame] | 283 | } |