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

Side by Side Diff: runtime/bin/socket_stream.dart

Issue 9029001: Add close to input stream and cleanup socket and streams (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Created 9 years 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
OLDNEW
1 // Copyright (c) 2011, the Dart project authors. Please see the AUTHORS file 1 // Copyright (c) 2011, 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 class SocketInputStream implements InputStream { 5 class SocketInputStream implements InputStream {
6 SocketInputStream(Socket socket) : _socket = socket; 6 SocketInputStream(Socket socket) : _socket = socket {
7 _socket.closeHandler = _closeHandler;
8 }
7 9
8 List<int> read([int len]) { 10 List<int> read([int len]) {
9 int bytesToRead = available(); 11 int bytesToRead = available();
10 if (bytesToRead == 0) return null; 12 if (bytesToRead == 0) return null;
11 if (len !== null) { 13 if (len !== null) {
12 if (len <= 0) { 14 if (len <= 0) {
13 throw new StreamException("Illegal length $len"); 15 throw new StreamException("Illegal length $len");
14 } else if (bytesToRead > len) { 16 } else if (bytesToRead > len) {
15 bytesToRead = len; 17 bytesToRead = len;
16 } 18 }
(...skipping 16 matching lines...) Expand all
33 if (len < 0) throw new StreamException("Illegal length $len"); 35 if (len < 0) throw new StreamException("Illegal length $len");
34 return _socket.readList(buffer, offset, len); 36 return _socket.readList(buffer, offset, len);
35 } 37 }
36 38
37 int available() => _socket.available(); 39 int available() => _socket.available();
38 40
39 void pipe(OutputStream output, [bool close = true]) { 41 void pipe(OutputStream output, [bool close = true]) {
40 _pipe(this, output, close: close); 42 _pipe(this, output, close: close);
41 } 43 }
42 44
45 void close() {
46 if (!_closed) {
47 _socket.close();
48 if (_clientCloseHandler !== null) _clientCloseHandler();
49 }
50 }
51
52 bool get closed() => _closed;
53
43 void set dataHandler(void callback()) { 54 void set dataHandler(void callback()) {
44 _socket.dataHandler = callback; 55 _socket._dataHandler = callback;
45 } 56 }
46 57
47 void set closeHandler(void callback()) { 58 void set closeHandler(void callback()) {
48 _clientCloseHandler = callback; 59 _clientCloseHandler = callback;
49 _socket.closeHandler = callback; 60 _socket._closeHandler = _closeHandler;
50 } 61 }
51 62
52 void set errorHandler(void callback()) { 63 void set errorHandler(void callback()) {
53 _socket.errorHandler = callback; 64 _socket.errorHandler = callback;
54 } 65 }
55 66
67 void _closeHandler() {
68 _closed = true;
69 if (_clientCloseHandler !== null) _clientCloseHandler();
70 }
71
56 Socket _socket; 72 Socket _socket;
57 Function _clientCloseHandler; 73 Function _clientCloseHandler;
74 bool _closed = false;
58 } 75 }
59 76
60 77
61 class SocketOutputStream implements OutputStream { 78 class SocketOutputStream implements OutputStream {
62 SocketOutputStream(Socket socket) 79 SocketOutputStream(Socket socket)
63 : _socket = socket, _pendingWrites = new _BufferList(); 80 : _socket = socket, _pendingWrites = new _BufferList();
64 81
65 bool write(List<int> buffer, [bool copyBuffer = true]) { 82 bool write(List<int> buffer, [bool copyBuffer = true]) {
66 return _write(buffer, 0, buffer.length, copyBuffer); 83 return _write(buffer, 0, buffer.length, copyBuffer);
67 } 84 }
68 85
69 bool writeFrom(List<int> buffer, [int offset = 0, int len]) { 86 bool writeFrom(List<int> buffer, [int offset = 0, int len]) {
70 return _write( 87 return _write(
71 buffer, offset, (len == null) ? buffer.length - offset : len, true); 88 buffer, offset, (len == null) ? buffer.length - offset : len, true);
72 } 89 }
73 90
74 void close() { 91 void close() {
75 if (!_pendingWrites.isEmpty()) { 92 if (!_pendingWrites.isEmpty()) {
76 // Mark the socket for close when all data is written. 93 // Mark the socket for close when all data is written.
77 _closing = true; 94 _closing = true;
78 _socket.writeHandler = _writeHandler; 95 _socket._writeHandler = _writeHandler;
79 } else { 96 } else {
80 // Close the socket for writing. 97 // Close the socket for writing.
81 _socket._closeWrite(); 98 _socket._closeWrite();
82 _closed = true; 99 _closed = true;
83 } 100 }
84 } 101 }
85 102
86 void destroy() { 103 void destroy() {
87 _socket.writeHandler = null; 104 _socket.writeHandler = null;
88 _pendingWrites.clear(); 105 _pendingWrites.clear();
89 _socket.close(); 106 _socket.close();
90 _closed = true; 107 _closed = true;
91 } 108 }
92 109
93 void set noPendingWriteHandler(void callback()) { 110 void set noPendingWriteHandler(void callback()) {
94 _noPendingWriteHandler = callback; 111 _noPendingWriteHandler = callback;
95 if (_noPendingWriteHandler != null) { 112 if (_noPendingWriteHandler != null) {
96 _socket.writeHandler = _writeHandler; 113 _socket._writeHandler = _writeHandler;
97 } 114 }
98 } 115 }
99 116
100 void set closeHandler(void callback()) { 117 void set closeHandler(void callback()) {
101 _socket.closeHandler = callback; 118 _socket.closeHandler = callback;
102 } 119 }
103 120
104 void set errorHandler(void callback()) { 121 void set errorHandler(void callback()) {
105 _streamErrorHandler = callback; 122 _streamErrorHandler = callback;
106 if (_streamErrorHandler != null) { 123 if (_streamErrorHandler != null) {
(...skipping 16 matching lines...) Expand all
123 // Place remaining data on the pending writes queue. 140 // Place remaining data on the pending writes queue.
124 int notWrittenOffset = offset + bytesWritten; 141 int notWrittenOffset = offset + bytesWritten;
125 if (copyBuffer) { 142 if (copyBuffer) {
126 List<int> newBuffer = 143 List<int> newBuffer =
127 buffer.getRange(notWrittenOffset, len - bytesWritten); 144 buffer.getRange(notWrittenOffset, len - bytesWritten);
128 _pendingWrites.add(newBuffer); 145 _pendingWrites.add(newBuffer);
129 } else { 146 } else {
130 assert(offset + len == buffer.length); 147 assert(offset + len == buffer.length);
131 _pendingWrites.add(buffer, notWrittenOffset); 148 _pendingWrites.add(buffer, notWrittenOffset);
132 } 149 }
133 _socket.writeHandler = _writeHandler; 150 _socket._writeHandler = _writeHandler;
134 return false; 151 return false;
135 } 152 }
136 153
137 void _writeHandler() { 154 void _writeHandler() {
138 // Write as much buffered data to the socket as possible. 155 // Write as much buffered data to the socket as possible.
139 while (!_pendingWrites.isEmpty()) { 156 while (!_pendingWrites.isEmpty()) {
140 List<int> buffer = _pendingWrites.first; 157 List<int> buffer = _pendingWrites.first;
141 int offset = _pendingWrites.index; 158 int offset = _pendingWrites.index;
142 int bytesToWrite = buffer.length - offset; 159 int bytesToWrite = buffer.length - offset;
143 int bytesWritten = _socket.writeList(buffer, offset, bytesToWrite); 160 int bytesWritten = _socket.writeList(buffer, offset, bytesToWrite);
144 _pendingWrites.removeBytes(bytesWritten); 161 _pendingWrites.removeBytes(bytesWritten);
145 if (bytesWritten < bytesToWrite) { 162 if (bytesWritten < bytesToWrite) {
146 _socket.writeHandler = _writeHandler; 163 _socket._writeHandler = _writeHandler;
147 return; 164 return;
148 } 165 }
149 } 166 }
150 167
151 // All buffered data was written. 168 // All buffered data was written.
152 if (_closing) { 169 if (_closing) {
153 _socket._closeWrite(); 170 _socket._closeWrite();
154 _closed = true; 171 _closed = true;
155 } else { 172 } else {
156 if (_noPendingWriteHandler != null) _noPendingWriteHandler(); 173 if (_noPendingWriteHandler != null) _noPendingWriteHandler();
157 } 174 }
158 if (_noPendingWriteHandler == null) _socket.writeHandler = null; 175 if (_noPendingWriteHandler == null) _socket._writeHandler = null;
159 } 176 }
160 177
161 void _errorHandler() { 178 void _errorHandler() {
162 close(); 179 close();
163 if (_streamErrorHandler != null) _streamErrorHandler(); 180 if (_streamErrorHandler != null) _streamErrorHandler();
164 } 181 }
165 182
166 Socket _socket; 183 Socket _socket;
167 _BufferList _pendingWrites; 184 _BufferList _pendingWrites;
168 var _noPendingWriteHandler; 185 var _noPendingWriteHandler;
169 var _streamErrorHandler; 186 var _streamErrorHandler;
170 bool _closing = false; 187 bool _closing = false;
171 bool _closed = false; 188 bool _closed = false;
172 } 189 }
OLDNEW

Powered by Google App Engine
This is Rietveld 408576698