| OLD | NEW |
| 1 // Copyright (c) 2013, the Dart project authors. Please see the AUTHORS file | 1 // Copyright (c) 2013, the Dart project authors. Please see the AUTHORS file |
| 2 // for details. All rights reserved. Use of this source code is governed by a | 2 // for details. All rights reserved. Use of this source code is governed by a |
| 3 // BSD-style license that can be found in the LICENSE file. | 3 // BSD-style license that can be found in the LICENSE file. |
| 4 | 4 |
| 5 part of dart.io; | 5 part of dart.io; |
| 6 | 6 |
| 7 const String _webSocketGUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; | 7 const String _webSocketGUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; |
| 8 | 8 |
| 9 class _WebSocketMessageType { | 9 class _WebSocketMessageType { |
| 10 static const int NONE = 0; | 10 static const int NONE = 0; |
| (...skipping 554 matching lines...) Expand 10 before | Expand all | Expand 10 after Loading... |
| 565 } | 565 } |
| 566 } | 566 } |
| 567 | 567 |
| 568 | 568 |
| 569 class _WebSocketConsumer implements StreamConsumer { | 569 class _WebSocketConsumer implements StreamConsumer { |
| 570 final _WebSocketImpl webSocket; | 570 final _WebSocketImpl webSocket; |
| 571 final Socket socket; | 571 final Socket socket; |
| 572 StreamController _controller; | 572 StreamController _controller; |
| 573 StreamSubscription _subscription; | 573 StreamSubscription _subscription; |
| 574 bool _issuedPause = false; | 574 bool _issuedPause = false; |
| 575 bool _reportError = false; |
| 575 Completer _closeCompleter = new Completer(); | 576 Completer _closeCompleter = new Completer(); |
| 576 Completer _completer; | 577 Completer _completer; |
| 577 | 578 |
| 578 _WebSocketConsumer(_WebSocketImpl this.webSocket, Socket this.socket); | 579 _WebSocketConsumer(_WebSocketImpl this.webSocket, Socket this.socket); |
| 579 | 580 |
| 580 void _onListen() { | 581 void _onListen() { |
| 581 if (_subscription != null) { | 582 if (_subscription != null) { |
| 582 _subscription.cancel(); | 583 _subscription.cancel(); |
| 583 } | 584 } |
| 584 } | 585 } |
| (...skipping 21 matching lines...) Expand all Loading... |
| 606 onResume: _onResume, | 607 onResume: _onResume, |
| 607 onCancel: _onListen); | 608 onCancel: _onListen); |
| 608 var stream = _controller.stream.transform( | 609 var stream = _controller.stream.transform( |
| 609 new _WebSocketOutgoingTransformer(webSocket)); | 610 new _WebSocketOutgoingTransformer(webSocket)); |
| 610 socket.addStream(stream) | 611 socket.addStream(stream) |
| 611 .then((_) { | 612 .then((_) { |
| 612 _done(); | 613 _done(); |
| 613 _closeCompleter.complete(webSocket); | 614 _closeCompleter.complete(webSocket); |
| 614 }, | 615 }, |
| 615 onError: (error) { | 616 onError: (error) { |
| 616 if (!_done(error)) { | 617 if (_reportError) { |
| 617 _closeCompleter.completeError(error); | 618 if (!_done(error)) { |
| 619 _closeCompleter.completeError(error); |
| 620 } |
| 621 } else { |
| 622 _done(); |
| 623 _closeCompleter.complete(webSocket); |
| 618 } | 624 } |
| 619 }); | 625 }); |
| 620 } | 626 } |
| 621 | 627 |
| 622 bool _done([error]) { | 628 bool _done([error]) { |
| 623 if (_completer == null) return false; | 629 if (_completer == null) return false; |
| 624 if (error != null) { | 630 if (error != null) { |
| 625 _completer.completeError(error); | 631 _completer.completeError(error); |
| 626 } else { | 632 } else { |
| 627 _completer.complete(webSocket); | 633 _completer.complete(webSocket); |
| 628 } | 634 } |
| 629 _completer = null; | 635 _completer = null; |
| 630 return true; | 636 return true; |
| 631 } | 637 } |
| 632 | 638 |
| 633 Future addStream(var stream) { | 639 Future addStream(var stream) { |
| 634 _ensureController(); | 640 _ensureController(); |
| 635 _completer = new Completer(); | 641 _completer = new Completer(); |
| 636 _subscription = stream.listen( | 642 _subscription = stream.listen( |
| 637 (data) { | 643 (data) { |
| 644 _reportError = true; |
| 638 _controller.add(data); | 645 _controller.add(data); |
| 639 }, | 646 }, |
| 640 onDone: () { | 647 onDone: () { |
| 641 _done(); | 648 _done(); |
| 642 }, | 649 }, |
| 643 onError: (error) { | 650 onError: (error) { |
| 644 _done(error); | 651 _done(error); |
| 645 }, | 652 }, |
| 646 cancelOnError: true); | 653 cancelOnError: true); |
| 647 if (_issuedPause) { | 654 if (_issuedPause) { |
| 648 _subscription.pause(); | 655 _subscription.pause(); |
| 649 _issuedPause = false; | 656 _issuedPause = false; |
| 650 } | 657 } |
| 651 return _completer.future; | 658 return _completer.future; |
| 652 } | 659 } |
| 653 | 660 |
| 654 Future close() { | 661 Future close() { |
| 655 _ensureController(); | 662 _ensureController(); |
| 656 Future closeSocket() { | 663 Future closeSocket() { |
| 657 return socket.close().then((_) => webSocket); | 664 return socket.close().then((_) => webSocket); |
| 658 } | 665 } |
| 659 _controller.close(); | 666 _controller.close(); |
| 660 return _closeCompleter.future.then((_) => closeSocket()); | 667 return _closeCompleter.future.then((_) => closeSocket()); |
| 661 } | 668 } |
| 662 | 669 |
| 663 void add(data) { | 670 void add(data) { |
| 664 _ensureController(); | 671 _ensureController(); |
| 672 _reportError = false; |
| 665 _controller.add(data); | 673 _controller.add(data); |
| 666 } | 674 } |
| 667 } | 675 } |
| 668 | 676 |
| 669 | 677 |
| 670 class _WebSocketImpl extends Stream implements WebSocket { | 678 class _WebSocketImpl extends Stream implements WebSocket { |
| 671 final StreamController _controller = new StreamController(sync: true); | 679 StreamController _controller; |
| 672 StreamSink _sink; | 680 StreamSink _sink; |
| 673 | 681 |
| 674 final Socket _socket; | 682 final Socket _socket; |
| 675 final bool _serverSide; | 683 final bool _serverSide; |
| 676 int _readyState = WebSocket.CONNECTING; | 684 int _readyState = WebSocket.CONNECTING; |
| 677 bool _writeClosed = false; | 685 bool _writeClosed = false; |
| 678 int _closeCode; | 686 int _closeCode; |
| 679 String _closeReason; | 687 String _closeReason; |
| 680 Duration _pingInterval; | 688 Duration _pingInterval; |
| 681 Timer _pingTimer; | 689 Timer _pingTimer; |
| 682 _WebSocketConsumer _consumer; | 690 _WebSocketConsumer _consumer; |
| 691 var _subscription; |
| 683 | 692 |
| 684 int _outCloseCode; | 693 int _outCloseCode; |
| 685 String _outCloseReason; | 694 String _outCloseReason; |
| 686 | 695 |
| 687 static final HttpClient _httpClient = new HttpClient(); | 696 static final HttpClient _httpClient = new HttpClient(); |
| 688 | 697 |
| 689 static Future<WebSocket> connect(String url, [protocols]) { | 698 static Future<WebSocket> connect(String url, [protocols]) { |
| 690 Uri uri = Uri.parse(url); | 699 Uri uri = Uri.parse(url); |
| 691 if (uri.scheme != "ws" && uri.scheme != "wss") { | 700 if (uri.scheme != "ws" && uri.scheme != "wss") { |
| 692 throw new WebSocketException("Unsupported URL scheme '${uri.scheme}'"); | 701 throw new WebSocketException("Unsupported URL scheme '${uri.scheme}'"); |
| (...skipping 63 matching lines...) Expand 10 before | Expand all | Expand 10 after Loading... |
| 756 }); | 765 }); |
| 757 } | 766 } |
| 758 | 767 |
| 759 _WebSocketImpl._fromSocket(Socket this._socket, | 768 _WebSocketImpl._fromSocket(Socket this._socket, |
| 760 [bool this._serverSide = false]) { | 769 [bool this._serverSide = false]) { |
| 761 _consumer = new _WebSocketConsumer(this, _socket); | 770 _consumer = new _WebSocketConsumer(this, _socket); |
| 762 _sink = new _StreamSinkImpl(_consumer); | 771 _sink = new _StreamSinkImpl(_consumer); |
| 763 _readyState = WebSocket.OPEN; | 772 _readyState = WebSocket.OPEN; |
| 764 | 773 |
| 765 var transformer = new _WebSocketProtocolTransformer(_serverSide); | 774 var transformer = new _WebSocketProtocolTransformer(_serverSide); |
| 766 _socket.transform(transformer).listen( | 775 _subscription = _socket.transform(transformer).listen( |
| 767 (data) { | 776 (data) { |
| 777 // Simply set pingInterval, as it'll cancel any timers. |
| 778 pingInterval = _pingInterval; |
| 768 if (data is _WebSocketPing) { | 779 if (data is _WebSocketPing) { |
| 769 if (!_writeClosed) _consumer.add(new _WebSocketPong(data.payload)); | 780 if (!_writeClosed) _consumer.add(new _WebSocketPong(data.payload)); |
| 770 } else if (data is _WebSocketPong) { | 781 } else if (data is _WebSocketPong) { |
| 771 // Simply set pingInterval, as it'll cancel any timers. | 782 // Already handled. |
| 772 pingInterval = _pingInterval; | |
| 773 } else { | 783 } else { |
| 774 _controller.add(data); | 784 _controller.add(data); |
| 775 } | 785 } |
| 776 }, | 786 }, |
| 777 onError: (error) { | 787 onError: (error) { |
| 778 if (error is ArgumentError) { | 788 if (error is ArgumentError) { |
| 779 close(WebSocketStatus.INVALID_FRAME_PAYLOAD_DATA); | 789 close(WebSocketStatus.INVALID_FRAME_PAYLOAD_DATA); |
| 780 } else { | 790 } else { |
| 781 close(WebSocketStatus.PROTOCOL_ERROR); | 791 close(WebSocketStatus.PROTOCOL_ERROR); |
| 782 } | 792 } |
| 783 _controller.addError(error); | 793 _controller.addError(error); |
| 784 _controller.close(); | 794 _controller.close(); |
| 785 }, | 795 }, |
| 786 onDone: () { | 796 onDone: () { |
| 787 if (_readyState == WebSocket.OPEN) { | 797 if (_readyState == WebSocket.OPEN) { |
| 788 _readyState = WebSocket.CLOSING; | 798 _readyState = WebSocket.CLOSING; |
| 789 if (!_isReservedStatusCode(transformer.closeCode)) { | 799 if (!_isReservedStatusCode(transformer.closeCode)) { |
| 790 close(transformer.closeCode); | 800 close(transformer.closeCode); |
| 791 } else { | 801 } else { |
| 792 close(); | 802 close(); |
| 793 } | 803 } |
| 794 _readyState = WebSocket.CLOSED; | 804 _readyState = WebSocket.CLOSED; |
| 795 } | 805 } |
| 796 _closeCode = transformer.closeCode; | 806 _closeCode = transformer.closeCode; |
| 797 _closeReason = transformer.closeReason; | 807 _closeReason = transformer.closeReason; |
| 798 _controller.close(); | 808 _controller.close(); |
| 799 }, | 809 }, |
| 800 cancelOnError: true); | 810 cancelOnError: true); |
| 811 _subscription.pause(); |
| 812 _controller = new StreamController(sync: true, |
| 813 onListen: _subscription.resume, |
| 814 onPause: _subscription.pause, |
| 815 onResume: _subscription.resume); |
| 801 } | 816 } |
| 802 | 817 |
| 803 StreamSubscription listen(void onData(message), | 818 StreamSubscription listen(void onData(message), |
| 804 {void onError(error), | 819 {void onError(error), |
| 805 void onDone(), | 820 void onDone(), |
| 806 bool cancelOnError}) { | 821 bool cancelOnError}) { |
| 807 return _controller.stream.listen(onData, | 822 return _controller.stream.listen(onData, |
| 808 onError: onError, | 823 onError: onError, |
| 809 onDone: onDone, | 824 onDone: onDone, |
| 810 cancelOnError: cancelOnError); | 825 cancelOnError: cancelOnError); |
| 811 } | 826 } |
| 812 | 827 |
| 813 Duration get pingnterval => _pingInterval; | 828 Duration get pingnterval => _pingInterval; |
| 814 | 829 |
| 815 void set pingInterval(Duration interval) { | 830 void set pingInterval(Duration interval) { |
| 816 if (_writeClosed) return; | 831 if (_writeClosed) return; |
| 817 if (_pingTimer != null) _pingTimer.cancel(); | 832 if (_pingTimer != null) _pingTimer.cancel(); |
| 818 _pingInterval = interval; | 833 _pingInterval = interval; |
| 819 | 834 |
| 820 if (_pingInterval == null) return; | 835 if (_pingInterval == null) return; |
| 821 | 836 |
| 822 _pingTimer = new Timer(_pingInterval, () { | 837 _pingTimer = new Timer(_pingInterval, () { |
| 823 if (_writeClosed) return; | 838 if (_writeClosed) return; |
| 824 _consumer.add(new _WebSocketPing()); | 839 _consumer.add(new _WebSocketPing()); |
| 825 _pingTimer = new Timer(_pingInterval, () { | 840 _pingTimer = new Timer(_pingInterval, () { |
| 826 // No pong received. | 841 // No pong received. |
| 842 _closeCode = WebSocketStatus.GOING_AWAY; |
| 843 _subscription.cancel(); |
| 827 close(WebSocketStatus.GOING_AWAY); | 844 close(WebSocketStatus.GOING_AWAY); |
| 828 }); | 845 }); |
| 829 }); | 846 }); |
| 830 } | 847 } |
| 831 | 848 |
| 832 int get readyState => _readyState; | 849 int get readyState => _readyState; |
| 833 | 850 |
| 834 String get extensions => null; | 851 String get extensions => null; |
| 835 String get protocol => null; | 852 String get protocol => null; |
| 836 int get closeCode => _closeCode; | 853 int get closeCode => _closeCode; |
| (...skipping 21 matching lines...) Expand all Loading... |
| 858 (code < WebSocketStatus.NORMAL_CLOSURE || | 875 (code < WebSocketStatus.NORMAL_CLOSURE || |
| 859 code == WebSocketStatus.RESERVED_1004 || | 876 code == WebSocketStatus.RESERVED_1004 || |
| 860 code == WebSocketStatus.NO_STATUS_RECEIVED || | 877 code == WebSocketStatus.NO_STATUS_RECEIVED || |
| 861 code == WebSocketStatus.ABNORMAL_CLOSURE || | 878 code == WebSocketStatus.ABNORMAL_CLOSURE || |
| 862 (code > WebSocketStatus.INTERNAL_SERVER_ERROR && | 879 (code > WebSocketStatus.INTERNAL_SERVER_ERROR && |
| 863 code < WebSocketStatus.RESERVED_1015) || | 880 code < WebSocketStatus.RESERVED_1015) || |
| 864 (code >= WebSocketStatus.RESERVED_1015 && | 881 (code >= WebSocketStatus.RESERVED_1015 && |
| 865 code < 3000)); | 882 code < 3000)); |
| 866 } | 883 } |
| 867 } | 884 } |
| OLD | NEW |