| OLD | NEW |
| 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 Loading... |
| 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 Loading... |
| 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 } |
| OLD | NEW |