Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(217)

Side by Side Diff: sdk/lib/io/websocket_impl.dart

Issue 19563006: Revert 25425 and disable web_socket_ping_test on simulators for now. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Created 7 years, 5 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View unified diff | Download patch | Annotate | Revision Log
« no previous file with comments | « sdk/lib/io/websocket.dart ('k') | tests/standalone/io/web_socket_ping_test.dart » ('j') | no next file with comments »
Toggle Intra-line Diffs ('i') | Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
OLDNEW
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
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
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
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
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 }
OLDNEW
« no previous file with comments | « sdk/lib/io/websocket.dart ('k') | tests/standalone/io/web_socket_ping_test.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698