1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16 package io.netty.handler.codec.spdy;
17
18 import io.netty.channel.ChannelDuplexHandler;
19 import io.netty.channel.ChannelFuture;
20 import io.netty.channel.ChannelFutureListener;
21 import io.netty.channel.ChannelHandlerContext;
22 import io.netty.channel.ChannelPromise;
23 import io.netty.util.internal.ObjectUtil;
24
25 import java.util.concurrent.atomic.AtomicInteger;
26
27 import static io.netty.handler.codec.spdy.SpdyCodecUtil.SPDY_SESSION_STREAM_ID;
28 import static io.netty.handler.codec.spdy.SpdyCodecUtil.isServerId;
29 import static io.netty.util.internal.ObjectUtil.checkPositive;
30 import static io.netty.util.internal.ObjectUtil.checkPositiveOrZero;
31
32
33
34
35 public class SpdySessionHandler extends ChannelDuplexHandler {
36
37 private static final SpdyProtocolException PROTOCOL_EXCEPTION =
38 SpdyProtocolException.newStatic(null, SpdySessionHandler.class, "handleOutboundMessage(...)");
39 private static final SpdyProtocolException STREAM_CLOSED =
40 SpdyProtocolException.newStatic("Stream closed", SpdySessionHandler.class, "removeStream(...)");
41
42 private static final int DEFAULT_WINDOW_SIZE = 64 * 1024;
43 private int initialSendWindowSize = DEFAULT_WINDOW_SIZE;
44 private int initialReceiveWindowSize = DEFAULT_WINDOW_SIZE;
45 private volatile int initialSessionReceiveWindowSize = DEFAULT_WINDOW_SIZE;
46
47 private final SpdySession spdySession = new SpdySession(initialSendWindowSize, initialReceiveWindowSize);
48 private int lastGoodStreamId;
49
50
51 private static final int DEFAULT_MAX_CONCURRENT_STREAMS = 100;
52 private int remoteConcurrentStreams = Integer.MAX_VALUE;
53 private int localConcurrentStreams = DEFAULT_MAX_CONCURRENT_STREAMS;
54
55 private final AtomicInteger pings = new AtomicInteger();
56
57 private boolean sentGoAwayFrame;
58 private boolean receivedGoAwayFrame;
59 private boolean sentInitialSettingsFrame;
60
61 private ChannelFutureListener closeSessionFutureListener;
62
63 private final boolean server;
64 private final int minorVersion;
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80 public SpdySessionHandler(SpdyVersion version, boolean server) {
81 this(version, server, DEFAULT_MAX_CONCURRENT_STREAMS);
82 }
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98 public SpdySessionHandler(SpdyVersion version, boolean server, int maxLocalConcurrentStreams) {
99 this.minorVersion = ObjectUtil.checkNotNull(version, "version").minorVersion();
100 this.server = server;
101 this.localConcurrentStreams = checkPositive(maxLocalConcurrentStreams, "maxLocalConcurrentStreams");
102 }
103
104 public void setSessionReceiveWindowSize(int sessionReceiveWindowSize) {
105 checkPositiveOrZero(sessionReceiveWindowSize, "sessionReceiveWindowSize");
106
107
108
109
110
111
112 initialSessionReceiveWindowSize = sessionReceiveWindowSize;
113 }
114
115 @Override
116 public void handlerAdded(ChannelHandlerContext ctx) throws Exception {
117 if (ctx.channel().isActive()) {
118 sendInitialSettingsFrame(ctx);
119 }
120 }
121
122 @Override
123 public void channelActive(ChannelHandlerContext ctx) throws Exception {
124 sendInitialSettingsFrame(ctx);
125 super.channelActive(ctx);
126 }
127
128 private void sendInitialSettingsFrame(ChannelHandlerContext ctx) {
129 if (sentInitialSettingsFrame) {
130 return;
131 }
132 sentInitialSettingsFrame = true;
133
134
135 SpdySettingsFrame settingsFrame = new DefaultSpdySettingsFrame();
136 settingsFrame.setValue(SpdySettingsFrame.SETTINGS_MAX_CONCURRENT_STREAMS, localConcurrentStreams);
137 ctx.writeAndFlush(settingsFrame);
138 }
139
140 @Override
141 public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
142 if (msg instanceof SpdyDataFrame) {
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166 SpdyDataFrame spdyDataFrame = (SpdyDataFrame) msg;
167 int streamId = spdyDataFrame.streamId();
168
169 int deltaWindowSize = -1 * spdyDataFrame.content().readableBytes();
170 int newSessionWindowSize =
171 spdySession.updateReceiveWindowSize(SPDY_SESSION_STREAM_ID, deltaWindowSize);
172
173
174 if (newSessionWindowSize < 0) {
175 issueSessionError(ctx, SpdySessionStatus.PROTOCOL_ERROR);
176 return;
177 }
178
179
180 if (newSessionWindowSize <= initialSessionReceiveWindowSize / 2) {
181 int sessionDeltaWindowSize = initialSessionReceiveWindowSize - newSessionWindowSize;
182 spdySession.updateReceiveWindowSize(SPDY_SESSION_STREAM_ID, sessionDeltaWindowSize);
183 SpdyWindowUpdateFrame spdyWindowUpdateFrame =
184 new DefaultSpdyWindowUpdateFrame(SPDY_SESSION_STREAM_ID, sessionDeltaWindowSize);
185 ctx.writeAndFlush(spdyWindowUpdateFrame);
186 }
187
188
189
190 if (!spdySession.isActiveStream(streamId)) {
191 spdyDataFrame.release();
192 if (streamId <= lastGoodStreamId) {
193 issueStreamError(ctx, streamId, SpdyStreamStatus.PROTOCOL_ERROR);
194 } else if (!sentGoAwayFrame) {
195 issueStreamError(ctx, streamId, SpdyStreamStatus.INVALID_STREAM);
196 }
197 return;
198 }
199
200
201
202 if (spdySession.isRemoteSideClosed(streamId)) {
203 spdyDataFrame.release();
204 issueStreamError(ctx, streamId, SpdyStreamStatus.STREAM_ALREADY_CLOSED);
205 return;
206 }
207
208
209 if (!isRemoteInitiatedId(streamId) && !spdySession.hasReceivedReply(streamId)) {
210 spdyDataFrame.release();
211 issueStreamError(ctx, streamId, SpdyStreamStatus.PROTOCOL_ERROR);
212 return;
213 }
214
215
216
217
218
219
220
221
222 int newWindowSize = spdySession.updateReceiveWindowSize(streamId, deltaWindowSize);
223
224
225
226
227
228
229 if (newWindowSize < spdySession.getReceiveWindowSizeLowerBound(streamId)) {
230 spdyDataFrame.release();
231 issueStreamError(ctx, streamId, SpdyStreamStatus.FLOW_CONTROL_ERROR);
232 return;
233 }
234
235
236
237 if (newWindowSize < 0) {
238 while (spdyDataFrame.content().readableBytes() > initialReceiveWindowSize) {
239 SpdyDataFrame partialDataFrame = new DefaultSpdyDataFrame(
240 streamId, spdyDataFrame.content().readRetainedSlice(initialReceiveWindowSize));
241 ctx.writeAndFlush(partialDataFrame);
242 }
243 }
244
245
246 if (newWindowSize <= initialReceiveWindowSize / 2 && !spdyDataFrame.isLast()) {
247 int streamDeltaWindowSize = initialReceiveWindowSize - newWindowSize;
248 spdySession.updateReceiveWindowSize(streamId, streamDeltaWindowSize);
249 SpdyWindowUpdateFrame spdyWindowUpdateFrame =
250 new DefaultSpdyWindowUpdateFrame(streamId, streamDeltaWindowSize);
251 ctx.writeAndFlush(spdyWindowUpdateFrame);
252 }
253
254
255 if (spdyDataFrame.isLast()) {
256 halfCloseStream(streamId, true, ctx.newSucceededFuture());
257 }
258
259 } else if (msg instanceof SpdySynStreamFrame) {
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275 SpdySynStreamFrame spdySynStreamFrame = (SpdySynStreamFrame) msg;
276 int streamId = spdySynStreamFrame.streamId();
277
278
279 if (spdySynStreamFrame.isInvalid() ||
280 !isRemoteInitiatedId(streamId) ||
281 spdySession.isActiveStream(streamId)) {
282 issueStreamError(ctx, streamId, SpdyStreamStatus.PROTOCOL_ERROR);
283 return;
284 }
285
286
287 if (streamId <= lastGoodStreamId) {
288 issueSessionError(ctx, SpdySessionStatus.PROTOCOL_ERROR);
289 return;
290 }
291
292
293 byte priority = spdySynStreamFrame.priority();
294 boolean remoteSideClosed = spdySynStreamFrame.isLast();
295 boolean localSideClosed = spdySynStreamFrame.isUnidirectional();
296 if (!acceptStream(streamId, priority, remoteSideClosed, localSideClosed)) {
297 issueStreamError(ctx, streamId, SpdyStreamStatus.REFUSED_STREAM);
298 return;
299 }
300
301 } else if (msg instanceof SpdySynReplyFrame) {
302
303
304
305
306
307
308
309
310 SpdySynReplyFrame spdySynReplyFrame = (SpdySynReplyFrame) msg;
311 int streamId = spdySynReplyFrame.streamId();
312
313
314 if (spdySynReplyFrame.isInvalid() ||
315 isRemoteInitiatedId(streamId) ||
316 spdySession.isRemoteSideClosed(streamId)) {
317 issueStreamError(ctx, streamId, SpdyStreamStatus.INVALID_STREAM);
318 return;
319 }
320
321
322 if (spdySession.hasReceivedReply(streamId)) {
323 issueStreamError(ctx, streamId, SpdyStreamStatus.STREAM_IN_USE);
324 return;
325 }
326
327 spdySession.receivedReply(streamId);
328
329
330 if (spdySynReplyFrame.isLast()) {
331 halfCloseStream(streamId, true, ctx.newSucceededFuture());
332 }
333
334 } else if (msg instanceof SpdyRstStreamFrame) {
335
336
337
338
339
340
341
342
343
344
345 SpdyRstStreamFrame spdyRstStreamFrame = (SpdyRstStreamFrame) msg;
346 removeStream(spdyRstStreamFrame.streamId(), ctx.newSucceededFuture());
347
348 } else if (msg instanceof SpdySettingsFrame) {
349
350 SpdySettingsFrame spdySettingsFrame = (SpdySettingsFrame) msg;
351
352 int settingsMinorVersion = spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_MINOR_VERSION);
353 if (settingsMinorVersion >= 0 && settingsMinorVersion != minorVersion) {
354
355 issueSessionError(ctx, SpdySessionStatus.PROTOCOL_ERROR);
356 return;
357 }
358
359 int newConcurrentStreams =
360 spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_MAX_CONCURRENT_STREAMS);
361 if (newConcurrentStreams >= 0) {
362 remoteConcurrentStreams = newConcurrentStreams;
363 }
364
365
366
367
368 if (spdySettingsFrame.isPersisted(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE)) {
369 spdySettingsFrame.removeValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE);
370 }
371 spdySettingsFrame.setPersistValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE, false);
372
373 int newInitialWindowSize =
374 spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE);
375 if (newInitialWindowSize >= 0) {
376 updateInitialSendWindowSize(newInitialWindowSize);
377 }
378
379 } else if (msg instanceof SpdyPingFrame) {
380
381
382
383
384
385
386
387
388
389
390 SpdyPingFrame spdyPingFrame = (SpdyPingFrame) msg;
391
392 if (isRemoteInitiatedId(spdyPingFrame.id())) {
393 ctx.writeAndFlush(spdyPingFrame);
394 return;
395 }
396
397
398 if (pings.get() == 0) {
399 return;
400 }
401 pings.getAndDecrement();
402
403 } else if (msg instanceof SpdyGoAwayFrame) {
404
405 receivedGoAwayFrame = true;
406
407 } else if (msg instanceof SpdyHeadersFrame) {
408
409 SpdyHeadersFrame spdyHeadersFrame = (SpdyHeadersFrame) msg;
410 int streamId = spdyHeadersFrame.streamId();
411
412
413 if (spdyHeadersFrame.isInvalid()) {
414 issueStreamError(ctx, streamId, SpdyStreamStatus.PROTOCOL_ERROR);
415 return;
416 }
417
418 if (spdySession.isRemoteSideClosed(streamId)) {
419 issueStreamError(ctx, streamId, SpdyStreamStatus.INVALID_STREAM);
420 return;
421 }
422
423
424 if (spdyHeadersFrame.isLast()) {
425 halfCloseStream(streamId, true, ctx.newSucceededFuture());
426 }
427
428 } else if (msg instanceof SpdyWindowUpdateFrame) {
429
430
431
432
433
434
435
436
437
438
439
440 SpdyWindowUpdateFrame spdyWindowUpdateFrame = (SpdyWindowUpdateFrame) msg;
441 int streamId = spdyWindowUpdateFrame.streamId();
442 int deltaWindowSize = spdyWindowUpdateFrame.deltaWindowSize();
443
444
445 if (streamId != SPDY_SESSION_STREAM_ID && spdySession.isLocalSideClosed(streamId)) {
446 return;
447 }
448
449
450 if (spdySession.getSendWindowSize(streamId) > Integer.MAX_VALUE - deltaWindowSize) {
451 if (streamId == SPDY_SESSION_STREAM_ID) {
452 issueSessionError(ctx, SpdySessionStatus.PROTOCOL_ERROR);
453 } else {
454 issueStreamError(ctx, streamId, SpdyStreamStatus.FLOW_CONTROL_ERROR);
455 }
456 return;
457 }
458
459 updateSendWindowSize(ctx, streamId, deltaWindowSize);
460 }
461
462 ctx.fireChannelRead(msg);
463 }
464
465 @Override
466 public void channelInactive(ChannelHandlerContext ctx) throws Exception {
467 for (Integer streamId: spdySession.activeStreams().keySet()) {
468 removeStream(streamId, ctx.newSucceededFuture());
469 }
470 ctx.fireChannelInactive();
471 }
472
473 @Override
474 public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
475 if (cause instanceof SpdyProtocolException) {
476 issueSessionError(ctx, SpdySessionStatus.PROTOCOL_ERROR);
477 }
478
479 ctx.fireExceptionCaught(cause);
480 }
481
482 @Override
483 public void close(ChannelHandlerContext ctx, ChannelPromise promise) throws Exception {
484 sendGoAwayFrame(ctx, promise);
485 }
486
487 @Override
488 public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception {
489 if (msg instanceof SpdyDataFrame ||
490 msg instanceof SpdySynStreamFrame ||
491 msg instanceof SpdySynReplyFrame ||
492 msg instanceof SpdyRstStreamFrame ||
493 msg instanceof SpdySettingsFrame ||
494 msg instanceof SpdyPingFrame ||
495 msg instanceof SpdyGoAwayFrame ||
496 msg instanceof SpdyHeadersFrame ||
497 msg instanceof SpdyWindowUpdateFrame) {
498
499 handleOutboundMessage(ctx, msg, promise);
500 } else {
501 ctx.write(msg, promise);
502 }
503 }
504
505 private void handleOutboundMessage(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception {
506 if (msg instanceof SpdyDataFrame) {
507
508 SpdyDataFrame spdyDataFrame = (SpdyDataFrame) msg;
509 int streamId = spdyDataFrame.streamId();
510
511
512 if (spdySession.isLocalSideClosed(streamId)) {
513 spdyDataFrame.release();
514 promise.setFailure(PROTOCOL_EXCEPTION);
515 return;
516 }
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531 int dataLength = spdyDataFrame.content().readableBytes();
532 int sendWindowSize = spdySession.getSendWindowSize(streamId);
533 int sessionSendWindowSize = spdySession.getSendWindowSize(SPDY_SESSION_STREAM_ID);
534 sendWindowSize = Math.min(sendWindowSize, sessionSendWindowSize);
535
536 if (sendWindowSize <= 0) {
537
538 spdySession.putPendingWrite(streamId, new SpdySession.PendingWrite(spdyDataFrame, promise));
539 return;
540 } else if (sendWindowSize < dataLength) {
541
542 spdySession.updateSendWindowSize(streamId, -1 * sendWindowSize);
543 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * sendWindowSize);
544
545
546 SpdyDataFrame partialDataFrame = new DefaultSpdyDataFrame(
547 streamId, spdyDataFrame.content().readRetainedSlice(sendWindowSize));
548
549
550 spdySession.putPendingWrite(streamId, new SpdySession.PendingWrite(spdyDataFrame, promise));
551
552
553
554 final ChannelHandlerContext context = ctx;
555 ctx.write(partialDataFrame).addListener(future -> {
556 if (!future.isSuccess()) {
557 issueSessionError(context, SpdySessionStatus.INTERNAL_ERROR);
558 }
559 });
560 return;
561 } else {
562
563 spdySession.updateSendWindowSize(streamId, -1 * dataLength);
564 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * dataLength);
565
566
567
568 final ChannelHandlerContext context = ctx;
569 promise.addListener(future -> {
570 if (!future.isSuccess()) {
571 issueSessionError(context, SpdySessionStatus.INTERNAL_ERROR);
572 }
573 });
574 }
575
576
577 if (spdyDataFrame.isLast()) {
578 halfCloseStream(streamId, false, promise);
579 }
580
581 } else if (msg instanceof SpdySynStreamFrame) {
582
583 SpdySynStreamFrame spdySynStreamFrame = (SpdySynStreamFrame) msg;
584 int streamId = spdySynStreamFrame.streamId();
585
586 if (isRemoteInitiatedId(streamId)) {
587 promise.setFailure(PROTOCOL_EXCEPTION);
588 return;
589 }
590
591 byte priority = spdySynStreamFrame.priority();
592 boolean remoteSideClosed = spdySynStreamFrame.isUnidirectional();
593 boolean localSideClosed = spdySynStreamFrame.isLast();
594 if (!acceptStream(streamId, priority, remoteSideClosed, localSideClosed)) {
595 promise.setFailure(PROTOCOL_EXCEPTION);
596 return;
597 }
598
599 } else if (msg instanceof SpdySynReplyFrame) {
600
601 SpdySynReplyFrame spdySynReplyFrame = (SpdySynReplyFrame) msg;
602 int streamId = spdySynReplyFrame.streamId();
603
604
605 if (!isRemoteInitiatedId(streamId) || spdySession.isLocalSideClosed(streamId)) {
606 promise.setFailure(PROTOCOL_EXCEPTION);
607 return;
608 }
609
610
611 if (spdySynReplyFrame.isLast()) {
612 halfCloseStream(streamId, false, promise);
613 }
614
615 } else if (msg instanceof SpdyRstStreamFrame) {
616
617 SpdyRstStreamFrame spdyRstStreamFrame = (SpdyRstStreamFrame) msg;
618 removeStream(spdyRstStreamFrame.streamId(), promise);
619
620 } else if (msg instanceof SpdySettingsFrame) {
621
622 SpdySettingsFrame spdySettingsFrame = (SpdySettingsFrame) msg;
623
624 int settingsMinorVersion = spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_MINOR_VERSION);
625 if (settingsMinorVersion >= 0 && settingsMinorVersion != minorVersion) {
626
627 promise.setFailure(PROTOCOL_EXCEPTION);
628 return;
629 }
630
631 int newConcurrentStreams =
632 spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_MAX_CONCURRENT_STREAMS);
633 if (newConcurrentStreams >= 0) {
634 localConcurrentStreams = newConcurrentStreams;
635 }
636
637
638
639
640 if (spdySettingsFrame.isPersisted(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE)) {
641 spdySettingsFrame.removeValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE);
642 }
643 spdySettingsFrame.setPersistValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE, false);
644
645 int newInitialWindowSize =
646 spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE);
647 if (newInitialWindowSize >= 0) {
648 updateInitialReceiveWindowSize(newInitialWindowSize);
649 }
650
651 } else if (msg instanceof SpdyPingFrame) {
652
653 SpdyPingFrame spdyPingFrame = (SpdyPingFrame) msg;
654 if (isRemoteInitiatedId(spdyPingFrame.id())) {
655 ctx.fireExceptionCaught(new IllegalArgumentException(
656 "invalid PING ID: " + spdyPingFrame.id()));
657 return;
658 }
659 pings.getAndIncrement();
660
661 } else if (msg instanceof SpdyGoAwayFrame) {
662
663
664
665 promise.setFailure(PROTOCOL_EXCEPTION);
666 return;
667
668 } else if (msg instanceof SpdyHeadersFrame) {
669
670 SpdyHeadersFrame spdyHeadersFrame = (SpdyHeadersFrame) msg;
671 int streamId = spdyHeadersFrame.streamId();
672
673
674 if (spdySession.isLocalSideClosed(streamId)) {
675 promise.setFailure(PROTOCOL_EXCEPTION);
676 return;
677 }
678
679
680 if (spdyHeadersFrame.isLast()) {
681 halfCloseStream(streamId, false, promise);
682 }
683
684 } else if (msg instanceof SpdyWindowUpdateFrame) {
685
686
687 promise.setFailure(PROTOCOL_EXCEPTION);
688 return;
689 }
690
691 ctx.write(msg, promise);
692 }
693
694
695
696
697
698
699
700
701
702
703 private void issueSessionError(
704 ChannelHandlerContext ctx, SpdySessionStatus status) {
705
706 sendGoAwayFrame(ctx, status).addListener(new ClosingChannelFutureListener(ctx, ctx.newPromise()));
707 }
708
709
710
711
712
713
714
715
716
717
718
719
720 private void issueStreamError(ChannelHandlerContext ctx, int streamId, SpdyStreamStatus status) {
721 boolean fireChannelRead = !spdySession.isRemoteSideClosed(streamId);
722 ChannelPromise promise = ctx.newPromise();
723 removeStream(streamId, promise);
724
725 SpdyRstStreamFrame spdyRstStreamFrame = new DefaultSpdyRstStreamFrame(streamId, status);
726 ctx.writeAndFlush(spdyRstStreamFrame, promise);
727 if (fireChannelRead) {
728 ctx.fireChannelRead(spdyRstStreamFrame);
729 }
730 }
731
732
733
734
735
736 private boolean isRemoteInitiatedId(int id) {
737 boolean serverId = isServerId(id);
738 return server && !serverId || !server && serverId;
739 }
740
741
742 private void updateInitialSendWindowSize(int newInitialWindowSize) {
743 int deltaWindowSize = newInitialWindowSize - initialSendWindowSize;
744 initialSendWindowSize = newInitialWindowSize;
745 spdySession.updateAllSendWindowSizes(deltaWindowSize);
746 }
747
748
749 private void updateInitialReceiveWindowSize(int newInitialWindowSize) {
750 int deltaWindowSize = newInitialWindowSize - initialReceiveWindowSize;
751 initialReceiveWindowSize = newInitialWindowSize;
752 spdySession.updateAllReceiveWindowSizes(deltaWindowSize);
753 }
754
755
756 private boolean acceptStream(
757 int streamId, byte priority, boolean remoteSideClosed, boolean localSideClosed) {
758
759 if (receivedGoAwayFrame || sentGoAwayFrame) {
760 return false;
761 }
762
763 boolean remote = isRemoteInitiatedId(streamId);
764 int maxConcurrentStreams = remote ? localConcurrentStreams : remoteConcurrentStreams;
765 if (spdySession.numActiveStreams(remote) >= maxConcurrentStreams) {
766 return false;
767 }
768 spdySession.acceptStream(
769 streamId, priority, remoteSideClosed, localSideClosed,
770 initialSendWindowSize, initialReceiveWindowSize, remote);
771 if (remote) {
772 lastGoodStreamId = streamId;
773 }
774 return true;
775 }
776
777 private void halfCloseStream(int streamId, boolean remote, ChannelFuture future) {
778 if (remote) {
779 spdySession.closeRemoteSide(streamId, isRemoteInitiatedId(streamId));
780 } else {
781 spdySession.closeLocalSide(streamId, isRemoteInitiatedId(streamId));
782 }
783 if (closeSessionFutureListener != null && spdySession.noActiveStreams()) {
784 future.addListener(closeSessionFutureListener);
785 }
786 }
787
788 private void removeStream(int streamId, ChannelFuture future) {
789 spdySession.removeStream(streamId, STREAM_CLOSED, isRemoteInitiatedId(streamId));
790
791 if (closeSessionFutureListener != null && spdySession.noActiveStreams()) {
792 future.addListener(closeSessionFutureListener);
793 }
794 }
795
796 private void updateSendWindowSize(final ChannelHandlerContext ctx, int streamId, int deltaWindowSize) {
797 spdySession.updateSendWindowSize(streamId, deltaWindowSize);
798
799 while (true) {
800
801 SpdySession.PendingWrite pendingWrite = spdySession.getPendingWrite(streamId);
802 if (pendingWrite == null) {
803 return;
804 }
805
806 SpdyDataFrame spdyDataFrame = pendingWrite.spdyDataFrame;
807 int dataFrameSize = spdyDataFrame.content().readableBytes();
808 int writeStreamId = spdyDataFrame.streamId();
809 int sendWindowSize = spdySession.getSendWindowSize(writeStreamId);
810 int sessionSendWindowSize = spdySession.getSendWindowSize(SPDY_SESSION_STREAM_ID);
811 sendWindowSize = Math.min(sendWindowSize, sessionSendWindowSize);
812
813 if (sendWindowSize <= 0) {
814 return;
815 } else if (sendWindowSize < dataFrameSize) {
816
817 spdySession.updateSendWindowSize(writeStreamId, -1 * sendWindowSize);
818 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * sendWindowSize);
819
820
821 SpdyDataFrame partialDataFrame = new DefaultSpdyDataFrame(
822 writeStreamId, spdyDataFrame.content().readRetainedSlice(sendWindowSize));
823
824
825 ctx.writeAndFlush(partialDataFrame).addListener(future -> {
826 if (!future.isSuccess()) {
827 issueSessionError(ctx, SpdySessionStatus.INTERNAL_ERROR);
828 }
829 });
830 } else {
831
832 spdySession.removePendingWrite(writeStreamId);
833 spdySession.updateSendWindowSize(writeStreamId, -1 * dataFrameSize);
834 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * dataFrameSize);
835
836
837 if (spdyDataFrame.isLast()) {
838 halfCloseStream(writeStreamId, false, pendingWrite.promise);
839 }
840
841
842
843 ctx.writeAndFlush(spdyDataFrame, pendingWrite.promise).addListener(future -> {
844 if (!future.isSuccess()) {
845 issueSessionError(ctx, SpdySessionStatus.INTERNAL_ERROR);
846 }
847 });
848 }
849 }
850 }
851
852 private void sendGoAwayFrame(ChannelHandlerContext ctx, ChannelPromise future) {
853
854 if (!ctx.channel().isActive()) {
855 ctx.close(future);
856 return;
857 }
858
859 ChannelFuture f = sendGoAwayFrame(ctx, SpdySessionStatus.OK);
860 if (spdySession.noActiveStreams()) {
861 f.addListener(new ClosingChannelFutureListener(ctx, future));
862 } else {
863 closeSessionFutureListener = new ClosingChannelFutureListener(ctx, future);
864 }
865
866 }
867
868 private ChannelFuture sendGoAwayFrame(
869 ChannelHandlerContext ctx, SpdySessionStatus status) {
870 if (!sentGoAwayFrame) {
871 sentGoAwayFrame = true;
872 SpdyGoAwayFrame spdyGoAwayFrame = new DefaultSpdyGoAwayFrame(lastGoodStreamId, status);
873 return ctx.writeAndFlush(spdyGoAwayFrame);
874 } else {
875 return ctx.newSucceededFuture();
876 }
877 }
878
879 private static final class ClosingChannelFutureListener implements ChannelFutureListener {
880 private final ChannelHandlerContext ctx;
881 private final ChannelPromise promise;
882
883 ClosingChannelFutureListener(ChannelHandlerContext ctx, ChannelPromise promise) {
884 this.ctx = ctx;
885 this.promise = promise;
886 }
887
888 @Override
889 public void operationComplete(ChannelFuture sentGoAwayFuture) throws Exception {
890 ctx.close(promise);
891 }
892 }
893 }