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