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