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(new ChannelFutureListener() {
556 @Override
557 public void operationComplete(ChannelFuture future) throws Exception {
558 if (!future.isSuccess()) {
559 issueSessionError(context, SpdySessionStatus.INTERNAL_ERROR);
560 }
561 }
562 });
563 return;
564 } else {
565
566 spdySession.updateSendWindowSize(streamId, -1 * dataLength);
567 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * dataLength);
568
569
570
571 final ChannelHandlerContext context = ctx;
572 promise.addListener(new ChannelFutureListener() {
573 @Override
574 public void operationComplete(ChannelFuture future) throws Exception {
575 if (!future.isSuccess()) {
576 issueSessionError(context, SpdySessionStatus.INTERNAL_ERROR);
577 }
578 }
579 });
580 }
581
582
583 if (spdyDataFrame.isLast()) {
584 halfCloseStream(streamId, false, promise);
585 }
586
587 } else if (msg instanceof SpdySynStreamFrame) {
588
589 SpdySynStreamFrame spdySynStreamFrame = (SpdySynStreamFrame) msg;
590 int streamId = spdySynStreamFrame.streamId();
591
592 if (isRemoteInitiatedId(streamId)) {
593 promise.setFailure(PROTOCOL_EXCEPTION);
594 return;
595 }
596
597 byte priority = spdySynStreamFrame.priority();
598 boolean remoteSideClosed = spdySynStreamFrame.isUnidirectional();
599 boolean localSideClosed = spdySynStreamFrame.isLast();
600 if (!acceptStream(streamId, priority, remoteSideClosed, localSideClosed)) {
601 promise.setFailure(PROTOCOL_EXCEPTION);
602 return;
603 }
604
605 } else if (msg instanceof SpdySynReplyFrame) {
606
607 SpdySynReplyFrame spdySynReplyFrame = (SpdySynReplyFrame) msg;
608 int streamId = spdySynReplyFrame.streamId();
609
610
611 if (!isRemoteInitiatedId(streamId) || spdySession.isLocalSideClosed(streamId)) {
612 promise.setFailure(PROTOCOL_EXCEPTION);
613 return;
614 }
615
616
617 if (spdySynReplyFrame.isLast()) {
618 halfCloseStream(streamId, false, promise);
619 }
620
621 } else if (msg instanceof SpdyRstStreamFrame) {
622
623 SpdyRstStreamFrame spdyRstStreamFrame = (SpdyRstStreamFrame) msg;
624 removeStream(spdyRstStreamFrame.streamId(), promise);
625
626 } else if (msg instanceof SpdySettingsFrame) {
627
628 SpdySettingsFrame spdySettingsFrame = (SpdySettingsFrame) msg;
629
630 int settingsMinorVersion = spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_MINOR_VERSION);
631 if (settingsMinorVersion >= 0 && settingsMinorVersion != minorVersion) {
632
633 promise.setFailure(PROTOCOL_EXCEPTION);
634 return;
635 }
636
637 int newConcurrentStreams =
638 spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_MAX_CONCURRENT_STREAMS);
639 if (newConcurrentStreams >= 0) {
640 localConcurrentStreams = newConcurrentStreams;
641 }
642
643
644
645
646 if (spdySettingsFrame.isPersisted(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE)) {
647 spdySettingsFrame.removeValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE);
648 }
649 spdySettingsFrame.setPersistValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE, false);
650
651 int newInitialWindowSize =
652 spdySettingsFrame.getValue(SpdySettingsFrame.SETTINGS_INITIAL_WINDOW_SIZE);
653 if (newInitialWindowSize >= 0) {
654 updateInitialReceiveWindowSize(newInitialWindowSize);
655 }
656
657 } else if (msg instanceof SpdyPingFrame) {
658
659 SpdyPingFrame spdyPingFrame = (SpdyPingFrame) msg;
660 if (isRemoteInitiatedId(spdyPingFrame.id())) {
661 ctx.fireExceptionCaught(new IllegalArgumentException(
662 "invalid PING ID: " + spdyPingFrame.id()));
663 return;
664 }
665 pings.getAndIncrement();
666
667 } else if (msg instanceof SpdyGoAwayFrame) {
668
669
670
671 promise.setFailure(PROTOCOL_EXCEPTION);
672 return;
673
674 } else if (msg instanceof SpdyHeadersFrame) {
675
676 SpdyHeadersFrame spdyHeadersFrame = (SpdyHeadersFrame) msg;
677 int streamId = spdyHeadersFrame.streamId();
678
679
680 if (spdySession.isLocalSideClosed(streamId)) {
681 promise.setFailure(PROTOCOL_EXCEPTION);
682 return;
683 }
684
685
686 if (spdyHeadersFrame.isLast()) {
687 halfCloseStream(streamId, false, promise);
688 }
689
690 } else if (msg instanceof SpdyWindowUpdateFrame) {
691
692
693 promise.setFailure(PROTOCOL_EXCEPTION);
694 return;
695 }
696
697 ctx.write(msg, promise);
698 }
699
700
701
702
703
704
705
706
707
708
709 private void issueSessionError(
710 ChannelHandlerContext ctx, SpdySessionStatus status) {
711
712 sendGoAwayFrame(ctx, status).addListener(new ClosingChannelFutureListener(ctx, ctx.newPromise()));
713 }
714
715
716
717
718
719
720
721
722
723
724
725
726 private void issueStreamError(ChannelHandlerContext ctx, int streamId, SpdyStreamStatus status) {
727 boolean fireChannelRead = !spdySession.isRemoteSideClosed(streamId);
728 ChannelPromise promise = ctx.newPromise();
729 removeStream(streamId, promise);
730
731 SpdyRstStreamFrame spdyRstStreamFrame = new DefaultSpdyRstStreamFrame(streamId, status);
732 ctx.writeAndFlush(spdyRstStreamFrame, promise);
733 if (fireChannelRead) {
734 ctx.fireChannelRead(spdyRstStreamFrame);
735 }
736 }
737
738
739
740
741
742 private boolean isRemoteInitiatedId(int id) {
743 boolean serverId = isServerId(id);
744 return server && !serverId || !server && serverId;
745 }
746
747
748 private void updateInitialSendWindowSize(int newInitialWindowSize) {
749 int deltaWindowSize = newInitialWindowSize - initialSendWindowSize;
750 initialSendWindowSize = newInitialWindowSize;
751 spdySession.updateAllSendWindowSizes(deltaWindowSize);
752 }
753
754
755 private void updateInitialReceiveWindowSize(int newInitialWindowSize) {
756 int deltaWindowSize = newInitialWindowSize - initialReceiveWindowSize;
757 initialReceiveWindowSize = newInitialWindowSize;
758 spdySession.updateAllReceiveWindowSizes(deltaWindowSize);
759 }
760
761
762 private boolean acceptStream(
763 int streamId, byte priority, boolean remoteSideClosed, boolean localSideClosed) {
764
765 if (receivedGoAwayFrame || sentGoAwayFrame) {
766 return false;
767 }
768
769 boolean remote = isRemoteInitiatedId(streamId);
770 int maxConcurrentStreams = remote ? localConcurrentStreams : remoteConcurrentStreams;
771 if (spdySession.numActiveStreams(remote) >= maxConcurrentStreams) {
772 return false;
773 }
774 spdySession.acceptStream(
775 streamId, priority, remoteSideClosed, localSideClosed,
776 initialSendWindowSize, initialReceiveWindowSize, remote);
777 if (remote) {
778 lastGoodStreamId = streamId;
779 }
780 return true;
781 }
782
783 private void halfCloseStream(int streamId, boolean remote, ChannelFuture future) {
784 if (remote) {
785 spdySession.closeRemoteSide(streamId, isRemoteInitiatedId(streamId));
786 } else {
787 spdySession.closeLocalSide(streamId, isRemoteInitiatedId(streamId));
788 }
789 if (closeSessionFutureListener != null && spdySession.noActiveStreams()) {
790 future.addListener(closeSessionFutureListener);
791 }
792 }
793
794 private void removeStream(int streamId, ChannelFuture future) {
795 spdySession.removeStream(streamId, STREAM_CLOSED, isRemoteInitiatedId(streamId));
796
797 if (closeSessionFutureListener != null && spdySession.noActiveStreams()) {
798 future.addListener(closeSessionFutureListener);
799 }
800 }
801
802 private void updateSendWindowSize(final ChannelHandlerContext ctx, int streamId, int deltaWindowSize) {
803 spdySession.updateSendWindowSize(streamId, deltaWindowSize);
804
805 while (true) {
806
807 SpdySession.PendingWrite pendingWrite = spdySession.getPendingWrite(streamId);
808 if (pendingWrite == null) {
809 return;
810 }
811
812 SpdyDataFrame spdyDataFrame = pendingWrite.spdyDataFrame;
813 int dataFrameSize = spdyDataFrame.content().readableBytes();
814 int writeStreamId = spdyDataFrame.streamId();
815 int sendWindowSize = spdySession.getSendWindowSize(writeStreamId);
816 int sessionSendWindowSize = spdySession.getSendWindowSize(SPDY_SESSION_STREAM_ID);
817 sendWindowSize = Math.min(sendWindowSize, sessionSendWindowSize);
818
819 if (sendWindowSize <= 0) {
820 return;
821 } else if (sendWindowSize < dataFrameSize) {
822
823 spdySession.updateSendWindowSize(writeStreamId, -1 * sendWindowSize);
824 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * sendWindowSize);
825
826
827 SpdyDataFrame partialDataFrame = new DefaultSpdyDataFrame(
828 writeStreamId, spdyDataFrame.content().readRetainedSlice(sendWindowSize));
829
830
831
832 ctx.writeAndFlush(partialDataFrame).addListener(new ChannelFutureListener() {
833 @Override
834 public void operationComplete(ChannelFuture future) throws Exception {
835 if (!future.isSuccess()) {
836 issueSessionError(ctx, SpdySessionStatus.INTERNAL_ERROR);
837 }
838 }
839 });
840 } else {
841
842 spdySession.removePendingWrite(writeStreamId);
843 spdySession.updateSendWindowSize(writeStreamId, -1 * dataFrameSize);
844 spdySession.updateSendWindowSize(SPDY_SESSION_STREAM_ID, -1 * dataFrameSize);
845
846
847 if (spdyDataFrame.isLast()) {
848 halfCloseStream(writeStreamId, false, pendingWrite.promise);
849 }
850
851
852
853 ctx.writeAndFlush(spdyDataFrame, pendingWrite.promise).addListener(new ChannelFutureListener() {
854 @Override
855 public void operationComplete(ChannelFuture future) throws Exception {
856 if (!future.isSuccess()) {
857 issueSessionError(ctx, SpdySessionStatus.INTERNAL_ERROR);
858 }
859 }
860 });
861 }
862 }
863 }
864
865 private void sendGoAwayFrame(ChannelHandlerContext ctx, ChannelPromise future) {
866
867 if (!ctx.channel().isActive()) {
868 ctx.close(future);
869 return;
870 }
871
872 ChannelFuture f = sendGoAwayFrame(ctx, SpdySessionStatus.OK);
873 if (spdySession.noActiveStreams()) {
874 f.addListener(new ClosingChannelFutureListener(ctx, future));
875 } else {
876 closeSessionFutureListener = new ClosingChannelFutureListener(ctx, future);
877 }
878
879 }
880
881 private ChannelFuture sendGoAwayFrame(
882 ChannelHandlerContext ctx, SpdySessionStatus status) {
883 if (!sentGoAwayFrame) {
884 sentGoAwayFrame = true;
885 SpdyGoAwayFrame spdyGoAwayFrame = new DefaultSpdyGoAwayFrame(lastGoodStreamId, status);
886 return ctx.writeAndFlush(spdyGoAwayFrame);
887 } else {
888 return ctx.newSucceededFuture();
889 }
890 }
891
892 private static final class ClosingChannelFutureListener implements ChannelFutureListener {
893 private final ChannelHandlerContext ctx;
894 private final ChannelPromise promise;
895
896 ClosingChannelFutureListener(ChannelHandlerContext ctx, ChannelPromise promise) {
897 this.ctx = ctx;
898 this.promise = promise;
899 }
900
901 @Override
902 public void operationComplete(ChannelFuture sentGoAwayFuture) throws Exception {
903 ctx.close(promise);
904 }
905 }
906 }